Skip to main content

bgpkit_parser/parser/mrt/
mrt_elem.rs

1#![allow(unused)]
2//! This module handles converting MRT records into individual per-prefix BGP elements.
3//!
4//! Each MRT record may contain reachability information for multiple prefixes. This module breaks
5//! down MRT records into corresponding BGP elements, and thus allowing users to more conveniently
6//! process BGP information on a per-prefix basis.
7use crate::models::*;
8use crate::parser::bgp::messages::parse_bgp_update_message;
9use crate::ParserError;
10use crate::ParserError::ParseError;
11use bytes::Bytes;
12use itertools::Itertools;
13use log::{error, warn};
14use std::collections::HashMap;
15use std::fmt::{Display, Formatter};
16use std::net::{IpAddr, Ipv4Addr};
17
18#[derive(Default, Debug, Clone)]
19pub struct Elementor {
20    pub peer_table: Option<PeerIndexTable>,
21}
22
23/// Error returned by [`Elementor::record_to_elems_iter`].
24#[derive(Debug)]
25pub enum ElemError {
26    /// The record contains a [`PeerIndexTable`]. The contained table can be
27    /// passed to [`Elementor::with_peer_table`] to create an initialized elementor.
28    UnexpectedPeerIndexTable(Box<PeerIndexTable>),
29    /// A peer table is required for processing TableDumpV2 RIB entries,
30    /// but none has been set on this elementor.
31    MissingPeerTable,
32    /// The record contains a [`RibGenericEntries`] which is not yet supported.
33    UnsupportedRibGeneric,
34}
35
36impl Display for ElemError {
37    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
38        match self {
39            ElemError::UnexpectedPeerIndexTable(_) => {
40                write!(f, "unexpected PeerIndexTable record")
41            }
42            ElemError::MissingPeerTable => {
43                write!(
44                    f,
45                    "peer table not set; call set_peer_table or use with_peer_table first"
46                )
47            }
48            ElemError::UnsupportedRibGeneric => {
49                write!(f, "RibGenericEntries not yet supported")
50            }
51        }
52    }
53}
54
55impl std::error::Error for ElemError {}
56
57// use macro_rules! <name of macro>{<Body>}
58macro_rules! get_attr_value {
59    ($a:tt, $b:expr) => {
60        if let Attribute::$a(x) = $b {
61            Some(x)
62        } else {
63            None
64        }
65    };
66}
67
68#[allow(clippy::type_complexity)]
69fn get_relevant_attributes(
70    attributes: Attributes,
71) -> (
72    Option<AsPath>,
73    Option<AsPath>,
74    Option<Origin>,
75    Option<IpAddr>,
76    Option<u32>,
77    Option<u32>,
78    Option<Vec<MetaCommunity>>,
79    bool,
80    Option<(Asn, BgpIdentifier)>,
81    Option<Nlri>,
82    Option<Nlri>,
83    Option<Asn>,
84    Option<Vec<AttrRaw>>,
85    Option<Vec<AttrRaw>>,
86) {
87    let mut as_path = None;
88    let mut as4_path = None;
89    let mut origin = None;
90    let mut next_hop = None;
91    let mut local_pref = Some(0);
92    let mut med = Some(0);
93    let mut atomic = false;
94    let mut aggregator = None;
95    let mut announced = None;
96    let mut withdrawn = None;
97    let mut otc = None;
98    let mut unknown = vec![];
99    let mut deprecated = vec![];
100
101    let mut communities_vec: Vec<MetaCommunity> = vec![];
102
103    for attr in attributes {
104        match attr {
105            AttributeValue::Origin(v) => origin = Some(v),
106            AttributeValue::AsPath(path) => as_path = Some(path),
107            AttributeValue::As4Path(path) => as4_path = Some(path),
108            AttributeValue::NextHop(v) => next_hop = Some(v),
109            AttributeValue::MultiExitDiscriminator(v) => med = Some(v),
110            AttributeValue::LocalPreference(v) => local_pref = Some(v),
111            AttributeValue::AtomicAggregate => atomic = true,
112            AttributeValue::Communities(v) => communities_vec.extend(
113                v.into_iter()
114                    .map(MetaCommunity::Plain)
115                    .collect::<Vec<MetaCommunity>>(),
116            ),
117            AttributeValue::ExtendedCommunities(v) => communities_vec.extend(
118                v.into_iter()
119                    .map(MetaCommunity::Extended)
120                    .collect::<Vec<MetaCommunity>>(),
121            ),
122            AttributeValue::Ipv6AddressSpecificExtendedCommunities(v) => communities_vec.extend(
123                v.into_iter()
124                    .map(MetaCommunity::Ipv6Extended)
125                    .collect::<Vec<MetaCommunity>>(),
126            ),
127            AttributeValue::LargeCommunities(v) => communities_vec.extend(
128                v.into_iter()
129                    .map(MetaCommunity::Large)
130                    .collect::<Vec<MetaCommunity>>(),
131            ),
132            AttributeValue::Aggregator { asn, id } | AttributeValue::As4Aggregator { asn, id } => {
133                aggregator = Some((asn, id))
134            }
135            AttributeValue::MpReachNlri(nlri) => announced = Some(nlri),
136            AttributeValue::MpUnreachNlri(nlri) => withdrawn = Some(nlri),
137            AttributeValue::OnlyToCustomer(o) => otc = Some(o),
138
139            AttributeValue::Unknown(t) | AttributeValue::Raw(t) => {
140                unknown.push(t);
141            }
142            AttributeValue::Deprecated(t) => {
143                deprecated.push(t);
144            }
145
146            AttributeValue::OriginatorId(_)
147            | AttributeValue::Clusters(_)
148            | AttributeValue::Development(_)
149            | AttributeValue::LinkState(_)
150            | AttributeValue::TunnelEncapsulation(_)
151            | AttributeValue::TrafficEngineering(_)
152            | AttributeValue::Aigp(_)
153            | AttributeValue::BfdDiscriminator(_)
154            | AttributeValue::BgpPrefixSid(_)
155            | AttributeValue::Bier(_)
156            | AttributeValue::Sfp(_)
157            | AttributeValue::AttrSet(_) => {}
158        };
159    }
160
161    let communities = match !communities_vec.is_empty() {
162        true => Some(communities_vec),
163        false => None,
164    };
165
166    // If the next_hop is not set, we try to get it from the announced NLRI.
167    let next_hop = next_hop.or_else(|| {
168        announced
169            .as_ref()
170            .and_then(|v| v.next_hop.as_ref().map(NextHopAddress::global_addr))
171    });
172
173    (
174        as_path,
175        as4_path,
176        origin,
177        next_hop,
178        local_pref,
179        med,
180        communities,
181        atomic,
182        aggregator,
183        announced,
184        withdrawn,
185        otc,
186        if unknown.is_empty() {
187            None
188        } else {
189            Some(unknown)
190        },
191        if deprecated.is_empty() {
192            None
193        } else {
194            Some(deprecated)
195        },
196    )
197}
198
199fn rib_entry_to_elem(prefix: NetworkPrefix, peer: &Peer, entry: RibEntry) -> BgpElem {
200    let (
201        as_path,
202        as4_path,
203        origin,
204        next_hop,
205        local_pref,
206        med,
207        communities,
208        atomic,
209        aggregator,
210        announced,
211        _withdrawn,
212        only_to_customer,
213        unknown,
214        deprecated,
215    ) = get_relevant_attributes(entry.attributes);
216
217    let path = match (as_path, as4_path) {
218        (None, None) => None,
219        (Some(v), None) => Some(v),
220        (None, Some(v)) => Some(v),
221        (Some(v1), Some(v2)) => Some(AsPath::merge_aspath_as4path(&v1, &v2)),
222    };
223
224    let next_hop = match next_hop {
225        Some(v) => Some(v),
226        None => announced.and_then(|v| v.next_hop.as_ref().map(NextHopAddress::global_addr)),
227    };
228
229    let origin_asns = path
230        .as_ref()
231        .map(|as_path| as_path.iter_origins().collect());
232
233    BgpElem {
234        timestamp: entry.originated_time as f64,
235        elem_type: ElemType::ANNOUNCE,
236        peer_ip: peer.peer_ip,
237        peer_asn: peer.peer_asn,
238        peer_bgp_id: Some(peer.peer_bgp_id),
239        prefix,
240        next_hop,
241        as_path: path,
242        origin,
243        origin_asns,
244        local_pref,
245        med,
246        communities,
247        atomic,
248        aggr_asn: aggregator.map(|v| v.0),
249        aggr_ip: aggregator.map(|v| v.1),
250        only_to_customer,
251        unknown,
252        deprecated,
253    }
254}
255
256/// Iterator over [`BgpElem`]s produced from a single [`MrtRecord`],
257/// without requiring a mutable reference to the [`Elementor`].
258///
259/// This avoids allocating a `Vec` for the common RIB table dump case
260/// by lazily converting each [`RibEntry`] into a [`BgpElem`] on demand.
261pub enum RecordElemIter<'a> {
262    #[doc(hidden)]
263    Empty,
264    #[doc(hidden)]
265    TableDump(Option<BgpElem>),
266    #[doc(hidden)]
267    TableDumpBatch(std::vec::IntoIter<TableDumpMessage>),
268    #[doc(hidden)]
269    RibAfi {
270        peer_table: &'a PeerIndexTable,
271        prefix: NetworkPrefix,
272        entries: std::vec::IntoIter<RibEntry>,
273    },
274    #[doc(hidden)]
275    Bgp4Mp(BgpUpdateElemIter),
276}
277
278impl Iterator for RecordElemIter<'_> {
279    type Item = BgpElem;
280
281    fn next(&mut self) -> Option<BgpElem> {
282        match self {
283            RecordElemIter::Empty => None,
284            RecordElemIter::TableDump(elem) => elem.take(),
285            RecordElemIter::TableDumpBatch(entries) => entries.next().map(table_dump_to_elem),
286            RecordElemIter::Bgp4Mp(iter) => iter.next(),
287            RecordElemIter::RibAfi {
288                peer_table,
289                prefix,
290                entries,
291            } => {
292                let entry = entries.next()?;
293                let pid = entry.peer_index;
294                match peer_table.get_peer_by_id(&pid) {
295                    Some(peer) => Some(rib_entry_to_elem(*prefix, peer, entry)),
296                    None => {
297                        error!("peer ID {} not found in peer_index table", pid);
298                        *self = RecordElemIter::Empty;
299                        None
300                    }
301                }
302            }
303        }
304    }
305
306    fn size_hint(&self) -> (usize, Option<usize>) {
307        match self {
308            RecordElemIter::Empty => (0, Some(0)),
309            RecordElemIter::TableDump(elem) => {
310                let n = elem.is_some() as usize;
311                (n, Some(n))
312            }
313            RecordElemIter::TableDumpBatch(entries) => {
314                let len = entries.len();
315                (len, Some(len))
316            }
317            RecordElemIter::Bgp4Mp(iter) => iter.size_hint(),
318            RecordElemIter::RibAfi { entries, .. } => {
319                let len = entries.len();
320                (len, Some(len))
321            }
322        }
323    }
324}
325
326/// Iterator over [`BgpElem`]s produced from a [`BgpUpdateMessage`],
327/// avoiding allocation by lazily yielding elements from announced and
328/// withdrawn prefixes in two phases.
329pub struct BgpUpdateElemIter {
330    timestamp: f64,
331    peer_ip: IpAddr,
332    peer_asn: Asn,
333    peer_bgp_id: Option<BgpIdentifier>,
334    only_to_customer: Option<Asn>,
335    // Announce-specific shared attributes
336    path: Option<AsPath>,
337    origin_asns: Option<Vec<Asn>>,
338    origin: Option<Origin>,
339    next_hop: Option<IpAddr>,
340    local_pref: Option<u32>,
341    med: Option<u32>,
342    communities: Option<Vec<MetaCommunity>>,
343    atomic: bool,
344    aggr_asn: Option<Asn>,
345    aggr_ip: Option<BgpIdentifier>,
346    unknown: Option<Vec<AttrRaw>>,
347    deprecated: Option<Vec<AttrRaw>>,
348    // Prefix iterators (two chained sources each)
349    announced:
350        std::iter::Chain<std::vec::IntoIter<NetworkPrefix>, std::vec::IntoIter<NetworkPrefix>>,
351    withdrawn:
352        std::iter::Chain<std::vec::IntoIter<NetworkPrefix>, std::vec::IntoIter<NetworkPrefix>>,
353    in_withdrawn_phase: bool,
354}
355
356impl Iterator for BgpUpdateElemIter {
357    type Item = BgpElem;
358
359    fn next(&mut self) -> Option<BgpElem> {
360        if !self.in_withdrawn_phase {
361            if let Some(prefix) = self.announced.next() {
362                return Some(BgpElem {
363                    timestamp: self.timestamp,
364                    elem_type: ElemType::ANNOUNCE,
365                    peer_ip: self.peer_ip,
366                    peer_asn: self.peer_asn,
367                    peer_bgp_id: self.peer_bgp_id,
368                    prefix,
369                    next_hop: self.next_hop,
370                    as_path: self.path.clone(),
371                    origin: self.origin,
372                    origin_asns: self.origin_asns.clone(),
373                    local_pref: self.local_pref,
374                    med: self.med,
375                    communities: self.communities.clone(),
376                    atomic: self.atomic,
377                    aggr_asn: self.aggr_asn,
378                    aggr_ip: self.aggr_ip,
379                    only_to_customer: self.only_to_customer,
380                    unknown: self.unknown.clone(),
381                    deprecated: self.deprecated.clone(),
382                });
383            }
384            self.in_withdrawn_phase = true;
385        }
386
387        self.withdrawn.next().map(|prefix| BgpElem {
388            timestamp: self.timestamp,
389            elem_type: ElemType::WITHDRAW,
390            peer_ip: self.peer_ip,
391            peer_asn: self.peer_asn,
392            peer_bgp_id: self.peer_bgp_id,
393            prefix,
394            next_hop: None,
395            as_path: None,
396            origin: None,
397            origin_asns: None,
398            local_pref: None,
399            med: None,
400            communities: None,
401            atomic: false,
402            aggr_asn: None,
403            aggr_ip: None,
404            only_to_customer: None,
405            unknown: None,
406            deprecated: None,
407        })
408    }
409
410    fn size_hint(&self) -> (usize, Option<usize>) {
411        let (ann_lo, ann_hi) = if self.in_withdrawn_phase {
412            (0, Some(0))
413        } else {
414            self.announced.size_hint()
415        };
416        let (wd_lo, wd_hi) = self.withdrawn.size_hint();
417        (ann_lo + wd_lo, ann_hi.and_then(|a| wd_hi.map(|w| a + w)))
418    }
419}
420
421impl Elementor {
422    pub fn new() -> Elementor {
423        Self::default()
424    }
425
426    /// Sets the peer index table for the elementor.
427    ///
428    /// This method takes an MRT record and extracts the peer index table from it if the record contains one.
429    /// The peer index table is required for processing TableDumpV2 records, as it contains the mapping between
430    /// peer indices and their corresponding IP addresses and ASNs.
431    ///
432    /// # Arguments
433    ///
434    /// * `record` - An MRT record that should contain a peer index table
435    ///
436    /// # Returns
437    ///
438    /// * `Ok(())` - If the peer table was successfully extracted and set
439    /// * `Err(ParserError)` - If the record does not contain a peer index table
440    ///
441    /// # Example
442    ///
443    /// ```no_run
444    /// use bgpkit_parser::{BgpkitParser, Elementor};
445    ///
446    /// let mut parser = BgpkitParser::new("rib.dump.bz2").unwrap();
447    /// let mut elementor = Elementor::new();
448    ///
449    /// // Get the first record which should be the peer index table
450    /// if let Ok(record) = parser.next_record() {
451    ///     elementor.set_peer_table(record).unwrap();
452    /// }
453    /// ```
454    pub fn set_peer_table(&mut self, record: MrtRecord) -> Result<(), ParserError> {
455        if let MrtMessage::TableDumpV2Message(TableDumpV2Message::PeerIndexTable(p)) =
456            record.message
457        {
458            self.peer_table = Some(p);
459            Ok(())
460        } else {
461            Err(ParseError("peer_table is not a PeerIndexTable".to_string()))
462        }
463    }
464
465    /// Creates an [`Elementor`] with the given [`PeerIndexTable`] already set.
466    pub fn with_peer_table(peer_table: PeerIndexTable) -> Elementor {
467        Elementor {
468            peer_table: Some(peer_table),
469        }
470    }
471
472    /// Convert a [`MrtRecord`] into an iterator of [`BgpElem`]s without
473    /// requiring `&mut self`.
474    ///
475    /// Unlike [`record_to_elems`](Elementor::record_to_elems), this method:
476    /// - Takes `&self` instead of `&mut self`, since the peer table must
477    ///   already be set via [`set_peer_table`](Elementor::set_peer_table) or
478    ///   [`with_peer_table`](Elementor::with_peer_table).
479    /// - Returns an error if the record contains a [`PeerIndexTable`] (which
480    ///   would require mutation).
481    /// - Returns a lazy [`RecordElemIter`] instead of collecting into a `Vec`,
482    ///   avoiding allocation for the common RIB table dump case.
483    ///
484    /// # Errors
485    ///
486    /// - [`ElemError::UnexpectedPeerIndexTable`] if the record is a PeerIndexTable message.
487    /// - [`ElemError::MissingPeerTable`] if the record requires a peer table but none is set.
488    pub fn record_to_elems_iter(&self, record: MrtRecord) -> Result<RecordElemIter<'_>, ElemError> {
489        let timestamp = {
490            let t = record.common_header.timestamp;
491            if let Some(micro) = &record.common_header.microsecond_timestamp {
492                let m = (*micro as f64) / 1000000.0;
493                t as f64 + m
494            } else {
495                f64::from(t)
496            }
497        };
498
499        match record.message {
500            MrtMessage::TableDumpMessage(msg) => {
501                Ok(RecordElemIter::TableDump(Some(table_dump_to_elem(msg))))
502            }
503            MrtMessage::TableDumpMessageBatch(messages) => {
504                Ok(RecordElemIter::TableDumpBatch(messages.into_iter()))
505            }
506
507            MrtMessage::TableDumpV2Message(msg) => match msg {
508                TableDumpV2Message::PeerIndexTable(p) => {
509                    Err(ElemError::UnexpectedPeerIndexTable(Box::new(p)))
510                }
511                TableDumpV2Message::RibAfi(t) => {
512                    let peer_table = self
513                        .peer_table
514                        .as_ref()
515                        .ok_or(ElemError::MissingPeerTable)?;
516                    Ok(RecordElemIter::RibAfi {
517                        peer_table,
518                        prefix: t.prefix,
519                        entries: t.rib_entries.into_iter(),
520                    })
521                }
522                TableDumpV2Message::RibGeneric(_) => Err(ElemError::UnsupportedRibGeneric),
523                TableDumpV2Message::GeoPeerTable(_) => Ok(RecordElemIter::Empty),
524            },
525
526            MrtMessage::Bgp4Mp(msg) => match msg {
527                Bgp4MpEnum::StateChange(_) => Ok(RecordElemIter::Empty),
528                Bgp4MpEnum::Message(v) => {
529                    match Elementor::bgp_to_elems_iter(
530                        v.bgp_message,
531                        timestamp,
532                        &v.peer_ip,
533                        &v.peer_asn,
534                    ) {
535                        Some(iter) => Ok(RecordElemIter::Bgp4Mp(iter)),
536                        None => Ok(RecordElemIter::Empty),
537                    }
538                }
539            },
540            MrtMessage::LegacyBgp(msg) => match msg {
541                LegacyBgp::StateChange(_) => Ok(RecordElemIter::Empty),
542                LegacyBgp::Message(message) => match Elementor::bgp_to_elems_iter(
543                    message.bgp_message,
544                    timestamp,
545                    &message.peer_ip,
546                    &message.peer_asn,
547                ) {
548                    Some(iter) => Ok(RecordElemIter::Bgp4Mp(iter)),
549                    None => Ok(RecordElemIter::Empty),
550                },
551            },
552        }
553    }
554
555    /// Convert a [BgpMessage] to a vector of [BgpElem]s.
556    ///
557    /// A [BgpMessage] may include `Update`, `Open`, `Notification` or `KeepAlive` messages,
558    /// and only `Update` message contains [BgpElem]s.
559    pub fn bgp_to_elems(
560        msg: BgpMessage,
561        timestamp: f64,
562        peer_ip: &IpAddr,
563        peer_asn: &Asn,
564    ) -> Vec<BgpElem> {
565        Elementor::bgp_to_elems_iter(msg, timestamp, peer_ip, peer_asn)
566            .map(|iter| iter.collect())
567            .unwrap_or_default()
568    }
569
570    /// Convert a [BgpMessage] into an iterator of [BgpElem]s.
571    ///
572    /// Returns `None` for non-Update messages (Open, Notification, KeepAlive, RouteRefresh).
573    pub fn bgp_to_elems_iter(
574        msg: BgpMessage,
575        timestamp: f64,
576        peer_ip: &IpAddr,
577        peer_asn: &Asn,
578    ) -> Option<BgpUpdateElemIter> {
579        match msg {
580            BgpMessage::Update(msg) => Some(Elementor::bgp_update_to_elems_iter(
581                msg, timestamp, peer_ip, peer_asn,
582            )),
583            BgpMessage::Open(_)
584            | BgpMessage::Notification(_)
585            | BgpMessage::KeepAlive
586            | BgpMessage::RouteRefresh(_) => None,
587        }
588    }
589
590    /// Convert a [BgpUpdateMessage] to a vector of [BgpElem]s.
591    pub fn bgp_update_to_elems(
592        msg: BgpUpdateMessage,
593        timestamp: f64,
594        peer_ip: &IpAddr,
595        peer_asn: &Asn,
596    ) -> Vec<BgpElem> {
597        Elementor::bgp_update_to_elems_iter(msg, timestamp, peer_ip, peer_asn).collect()
598    }
599
600    /// Convert a [BgpUpdateMessage] into a [`BgpUpdateElemIter`] that lazily
601    /// yields [BgpElem]s without allocating a `Vec`.
602    pub fn bgp_update_to_elems_iter(
603        msg: BgpUpdateMessage,
604        timestamp: f64,
605        peer_ip: &IpAddr,
606        peer_asn: &Asn,
607    ) -> BgpUpdateElemIter {
608        let (
609            as_path,
610            as4_path,
611            origin,
612            next_hop,
613            local_pref,
614            med,
615            communities,
616            atomic,
617            aggregator,
618            announced,
619            withdrawn,
620            only_to_customer,
621            unknown,
622            deprecated,
623        ) = get_relevant_attributes(msg.attributes);
624
625        let path = match (as_path, as4_path) {
626            (None, None) => None,
627            (Some(v), None) => Some(v),
628            (None, Some(v)) => Some(v),
629            (Some(v1), Some(v2)) => Some(AsPath::merge_aspath_as4path(&v1, &v2)),
630        };
631
632        let origin_asns = path
633            .as_ref()
634            .map(|as_path| as_path.iter_origins().collect());
635
636        let nlri_announced = announced.map(|n| n.prefixes).unwrap_or_default();
637        let nlri_withdrawn = withdrawn.map(|n| n.prefixes).unwrap_or_default();
638
639        BgpUpdateElemIter {
640            timestamp,
641            peer_ip: *peer_ip,
642            peer_asn: *peer_asn,
643            peer_bgp_id: None,
644            only_to_customer,
645            path,
646            origin_asns,
647            origin,
648            next_hop,
649            local_pref,
650            med,
651            communities,
652            atomic,
653            aggr_asn: aggregator.as_ref().map(|v| v.0),
654            aggr_ip: aggregator.as_ref().map(|v| v.1),
655            unknown,
656            deprecated,
657            announced: msg.announced_prefixes.into_iter().chain(nlri_announced),
658            withdrawn: msg.withdrawn_prefixes.into_iter().chain(nlri_withdrawn),
659            in_withdrawn_phase: false,
660        }
661    }
662
663    /// Convert a [MrtRecord] to a vector of [BgpElem]s.
664    ///
665    /// If the record is a [`PeerIndexTable`], it is consumed to set the internal
666    /// peer table. Errors are logged.
667    ///
668    /// For a non-mutating, lazy alternative, see
669    /// [`record_to_elems_iter`](Elementor::record_to_elems_iter).
670    pub fn record_to_elems(&mut self, record: MrtRecord) -> Vec<BgpElem> {
671        match record.message {
672            MrtMessage::TableDumpV2Message(TableDumpV2Message::PeerIndexTable(_)) => {
673                self.set_peer_table(record);
674                vec![]
675            }
676            _ => match self.record_to_elems_iter(record) {
677                Ok(iter) => iter.collect(),
678                Err(e) => {
679                    error!("{}", e);
680                    vec![]
681                }
682            },
683        }
684    }
685}
686
687fn table_dump_to_elem(msg: TableDumpMessage) -> BgpElem {
688    let (
689        as_path,
690        _as4_path,
691        origin,
692        next_hop,
693        local_pref,
694        med,
695        communities,
696        atomic,
697        aggregator,
698        _announced,
699        _withdrawn,
700        only_to_customer,
701        unknown,
702        deprecated,
703    ) = get_relevant_attributes(msg.attributes);
704
705    let origin_asns = as_path
706        .as_ref()
707        .map(|as_path| as_path.iter_origins().collect());
708
709    BgpElem {
710        timestamp: msg.originated_time as f64,
711        elem_type: ElemType::ANNOUNCE,
712        peer_ip: msg.peer_ip,
713        peer_asn: msg.peer_asn,
714        peer_bgp_id: None,
715        prefix: msg.prefix,
716        next_hop,
717        as_path,
718        origin,
719        origin_asns,
720        local_pref,
721        med,
722        communities,
723        atomic,
724        aggr_asn: aggregator.map(|v| v.0),
725        aggr_ip: aggregator.map(|v| v.1),
726        only_to_customer,
727        unknown,
728        deprecated,
729    }
730}
731
732#[inline(always)]
733pub fn option_to_string<T>(o: &Option<T>) -> String
734where
735    T: Display,
736{
737    if let Some(v) = o {
738        v.to_string()
739    } else {
740        String::new()
741    }
742}
743
744impl From<&BgpElem> for Attributes {
745    fn from(value: &BgpElem) -> Self {
746        let mut values = Vec::<AttributeValue>::new();
747        let mut attributes = Attributes::default();
748        let prefix = value.prefix;
749
750        if value.elem_type == ElemType::WITHDRAW {
751            values.push(AttributeValue::MpUnreachNlri(Nlri::new_unreachable(prefix)));
752            attributes.extend(values);
753            return attributes;
754        }
755
756        values.push(AttributeValue::MpReachNlri(Nlri::new_reachable(
757            prefix,
758            value.next_hop,
759        )));
760
761        if let Some(v) = value.next_hop {
762            values.push(AttributeValue::NextHop(v));
763        }
764
765        if let Some(v) = value.as_path.as_ref() {
766            // The elem path is the RFC 6793 merged (effective) path, so it
767            // always maps to the plain AS_PATH attribute; the segment width is
768            // decided by the session's `asn_len` at encode time.
769            values.push(AttributeValue::AsPath(v.clone()));
770        }
771
772        if let Some(v) = value.origin {
773            values.push(AttributeValue::Origin(v));
774        }
775
776        if let Some(v) = value.local_pref {
777            values.push(AttributeValue::LocalPreference(v));
778        }
779
780        if let Some(v) = value.med {
781            values.push(AttributeValue::MultiExitDiscriminator(v));
782        }
783
784        if let Some(v) = value.communities.as_ref() {
785            let mut communites = vec![];
786            let mut extended_communities = vec![];
787            let mut ipv6_extended_communities = vec![];
788            let mut large_communities = vec![];
789            for c in v {
790                match c {
791                    MetaCommunity::Plain(v) => communites.push(*v),
792                    MetaCommunity::Extended(v) => extended_communities.push(*v),
793                    MetaCommunity::Large(v) => large_communities.push(*v),
794                    MetaCommunity::Ipv6Extended(v) => ipv6_extended_communities.push(*v),
795                }
796            }
797            if !communites.is_empty() {
798                values.push(AttributeValue::Communities(communites));
799            }
800            if !extended_communities.is_empty() {
801                values.push(AttributeValue::ExtendedCommunities(extended_communities));
802            }
803            if !large_communities.is_empty() {
804                values.push(AttributeValue::LargeCommunities(large_communities));
805            }
806            if !ipv6_extended_communities.is_empty() {
807                values.push(AttributeValue::Ipv6AddressSpecificExtendedCommunities(
808                    ipv6_extended_communities,
809                ));
810            }
811        }
812
813        if let Some(v) = value.aggr_asn {
814            let aggregator_id = match value.aggr_ip {
815                Some(v) => v,
816                None => Ipv4Addr::UNSPECIFIED,
817            };
818            values.push(AttributeValue::Aggregator {
819                asn: v,
820                id: aggregator_id,
821            });
822        }
823
824        if let Some(v) = value.only_to_customer {
825            values.push(AttributeValue::OnlyToCustomer(v));
826        }
827
828        if let Some(v) = value.unknown.as_ref() {
829            for t in v {
830                values.push(AttributeValue::Unknown(t.clone()));
831            }
832        }
833
834        if let Some(v) = value.deprecated.as_ref() {
835            for t in v {
836                values.push(AttributeValue::Deprecated(t.clone()));
837            }
838        }
839
840        attributes.extend(values);
841        attributes
842    }
843}
844
845#[cfg(test)]
846mod tests {
847    use super::*;
848    use crate::BgpkitParser;
849    use std::net::{Ipv4Addr, Ipv6Addr};
850    use std::str::FromStr;
851
852    #[test]
853    fn test_option_to_string() {
854        let o1 = Some(1);
855        let o2: Option<u32> = None;
856        assert_eq!(option_to_string(&o1), "1");
857        assert_eq!(option_to_string(&o2), "");
858    }
859
860    #[test]
861    fn test_record_to_elems() {
862        let url_table_dump_v1 = "https://data.ris.ripe.net/rrc00/2003.01/bview.20030101.0000.gz";
863        let url_table_dump_v2 = "https://data.ris.ripe.net/rrc00/2023.01/bview.20230101.0000.gz";
864        let url_bgp4mp = "https://data.ris.ripe.net/rrc00/2021.10/updates.20211001.0000.gz";
865
866        let mut elementor = Elementor::new();
867        let parser = BgpkitParser::new(url_table_dump_v1).unwrap();
868        let mut record_iter = parser.into_record_iter();
869        let record = record_iter.next().unwrap();
870        let elems = elementor.record_to_elems(record);
871        assert_eq!(elems.len(), 1);
872
873        let parser = BgpkitParser::new(url_table_dump_v2).unwrap();
874        let mut record_iter = parser.into_record_iter();
875        let peer_index_table = record_iter.next().unwrap();
876        let _elems = elementor.record_to_elems(peer_index_table);
877        let record = record_iter.next().unwrap();
878        let elems = elementor.record_to_elems(record);
879        assert!(!elems.is_empty());
880
881        let parser = BgpkitParser::new(url_bgp4mp).unwrap();
882        let mut record_iter = parser.into_record_iter();
883        let record = record_iter.next().unwrap();
884        let elems = elementor.record_to_elems(record);
885        assert!(!elems.is_empty());
886    }
887
888    #[test]
889    fn test_attributes_from_bgp_elem() {
890        let mut elem = BgpElem {
891            timestamp: 0.0,
892            elem_type: ElemType::ANNOUNCE,
893            peer_ip: IpAddr::from_str("10.0.0.1").unwrap(),
894            peer_asn: Asn::new_32bit(65000),
895            peer_bgp_id: None,
896            prefix: NetworkPrefix::from_str("10.0.1.0/24").unwrap(),
897            next_hop: Some(IpAddr::from_str("10.0.0.2").unwrap()),
898            as_path: Some(AsPath::from_sequence([65000, 65001, 65002])),
899            origin: Some(Origin::EGP),
900            origin_asns: Some(vec![Asn::new_32bit(65000)]),
901            local_pref: Some(100),
902            med: Some(200),
903            communities: Some(vec![
904                MetaCommunity::Plain(Community::NoAdvertise),
905                MetaCommunity::Extended(ExtendedCommunity::Raw([0, 0, 0, 0, 0, 0, 0, 0])),
906                MetaCommunity::Large(LargeCommunity {
907                    global_admin: 0,
908                    local_data: [0, 0],
909                }),
910                MetaCommunity::Ipv6Extended(Ipv6AddrExtCommunity {
911                    community_type: ExtendedCommunityType::TransitiveTwoOctetAs,
912                    subtype: 0,
913                    global_admin: Ipv6Addr::from_str("2001:db8::").unwrap(),
914                    local_admin: [0, 0],
915                }),
916            ]),
917            atomic: false,
918            aggr_asn: Some(Asn::new_32bit(65000)),
919            aggr_ip: Some(Ipv4Addr::from_str("10.2.0.0").unwrap()),
920            only_to_customer: Some(Asn::new_32bit(65000)),
921            unknown: Some(vec![AttrRaw {
922                code: AttrType::RESERVED.into(),
923                bytes: Bytes::new(),
924            }]),
925            deprecated: Some(vec![AttrRaw {
926                code: AttrType::RESERVED.into(),
927                bytes: Bytes::new(),
928            }]),
929        };
930
931        let _attributes = Attributes::from(&elem);
932        elem.elem_type = ElemType::WITHDRAW;
933        let _attributes = Attributes::from(&elem);
934    }
935
936    #[test]
937    fn test_get_relevant_attributes() {
938        let attributes = vec![
939            AttributeValue::Origin(Origin::IGP),
940            AttributeValue::As4Path(AsPath::from_sequence([65000, 65001, 65002])),
941            AttributeValue::NextHop(IpAddr::from_str("10.0.0.1").unwrap()),
942            AttributeValue::MultiExitDiscriminator(100),
943            AttributeValue::LocalPreference(200),
944            AttributeValue::AtomicAggregate,
945            AttributeValue::Aggregator {
946                asn: Asn::new_32bit(65000),
947                id: Ipv4Addr::from_str("10.0.0.1").unwrap(),
948            },
949            AttributeValue::Communities(vec![Community::NoExport]),
950            AttributeValue::ExtendedCommunities(vec![ExtendedCommunity::Raw([
951                0, 0, 0, 0, 0, 0, 0, 0,
952            ])]),
953            AttributeValue::LargeCommunities(vec![LargeCommunity {
954                global_admin: 0,
955                local_data: [0, 0],
956            }]),
957            AttributeValue::Ipv6AddressSpecificExtendedCommunities(vec![Ipv6AddrExtCommunity {
958                community_type: ExtendedCommunityType::TransitiveTwoOctetAs,
959                subtype: 0,
960                global_admin: Ipv6Addr::from_str("2001:db8::").unwrap(),
961                local_admin: [0, 0],
962            }]),
963            AttributeValue::MpReachNlri(Nlri::new_reachable(
964                NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
965                Some(IpAddr::from_str("10.0.0.1").unwrap()),
966            )),
967            AttributeValue::MpUnreachNlri(Nlri::new_unreachable(
968                NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
969            )),
970            AttributeValue::OnlyToCustomer(Asn::new_32bit(65000)),
971            AttributeValue::Unknown(AttrRaw {
972                code: AttrType::RESERVED.into(),
973                bytes: Bytes::new(),
974            }),
975            AttributeValue::Deprecated(AttrRaw {
976                code: AttrType::RESERVED.into(),
977                bytes: Bytes::new(),
978            }),
979        ]
980        .into_iter()
981        .map(Attribute::from)
982        .collect::<Vec<Attribute>>();
983
984        let attributes = Attributes::from(attributes);
985
986        let (
987            _as_path,
988            _as4_path, // Table dump v1 does not have 4-byte AS number
989            _origin,
990            _next_hop,
991            _local_pref,
992            _med,
993            _communities,
994            _atomic,
995            _aggregator,
996            _announced,
997            _withdrawn,
998            _only_to_customer,
999            _unknown,
1000            _deprecated,
1001        ) = get_relevant_attributes(attributes);
1002    }
1003
1004    #[test]
1005    fn test_next_hop_from_nlri() {
1006        let attributes = vec![AttributeValue::NextHop(
1007            IpAddr::from_str("10.0.0.1").unwrap(),
1008        )]
1009        .into_iter()
1010        .map(Attribute::from)
1011        .collect::<Vec<Attribute>>();
1012
1013        let attributes = Attributes::from(attributes);
1014
1015        let (
1016            _as_path,
1017            _as4_path, // Table dump v1 does not have 4-byte AS number
1018            _origin,
1019            next_hop,
1020            _local_pref,
1021            _med,
1022            _communities,
1023            _atomic,
1024            _aggregator,
1025            _announced,
1026            _withdrawn,
1027            _only_to_customer,
1028            _unknown,
1029            _deprecated,
1030        ) = get_relevant_attributes(attributes);
1031
1032        assert_eq!(next_hop, Some(IpAddr::from_str("10.0.0.1").unwrap()));
1033
1034        let attributes = vec![AttributeValue::MpReachNlri(Nlri::new_reachable(
1035            NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
1036            Some(IpAddr::from_str("10.0.0.2").unwrap()),
1037        ))]
1038        .into_iter()
1039        .map(Attribute::from)
1040        .collect::<Vec<Attribute>>();
1041
1042        let attributes = Attributes::from(attributes);
1043
1044        let (
1045            _as_path,
1046            _as4_path, // Table dump v1 does not have 4-byte AS number
1047            _origin,
1048            next_hop,
1049            _local_pref,
1050            _med,
1051            _communities,
1052            _atomic,
1053            _aggregator,
1054            _announced,
1055            _withdrawn,
1056            _only_to_customer,
1057            _unknown,
1058            _deprecated,
1059        ) = get_relevant_attributes(attributes);
1060
1061        assert_eq!(next_hop, Some(IpAddr::from_str("10.0.0.2").unwrap()));
1062    }
1063
1064    #[test]
1065    fn test_record_to_elems_iter_equivalence_tabledumpv2_small() {
1066        // rib-example-small.bz2 is a TableDumpV2 file (starts with PeerIndexTable)
1067        let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1068
1069        let mut elementor = Elementor::new();
1070        let parser = BgpkitParser::new(url).unwrap();
1071        let mut record_iter = parser.into_record_iter();
1072
1073        // Skip the PeerIndexTable
1074        let peer_index_table = record_iter.next().unwrap();
1075        let _ = elementor.record_to_elems(peer_index_table);
1076
1077        // Process the first RIB entry
1078        let record = record_iter.next().unwrap();
1079        let elems_vec = elementor.record_to_elems(record.clone());
1080        let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1081        assert_eq!(elems_vec, elems_iter);
1082        assert!(!elems_vec.is_empty());
1083    }
1084
1085    #[test]
1086    fn test_record_to_elems_iter_equivalence_bgp4mp() {
1087        let url = "https://spaces.bgpkit.org/parser/update-example.gz";
1088
1089        let mut elementor = Elementor::new();
1090        let parser = BgpkitParser::new(url).unwrap();
1091        let mut record_iter = parser.into_record_iter();
1092        let record = record_iter.next().unwrap();
1093
1094        let elems_vec = elementor.record_to_elems(record.clone());
1095        let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1096        assert_eq!(elems_vec, elems_iter);
1097        assert!(!elems_vec.is_empty());
1098    }
1099
1100    #[test]
1101    #[ignore = "requires large RIB file download"]
1102    fn test_record_to_elems_iter_equivalence_tabledumpv2() {
1103        let url = "https://data.ris.ripe.net/rrc00/2023.01/bview.20230101.0000.gz";
1104
1105        let mut elementor = Elementor::new();
1106        let parser = BgpkitParser::new(url).unwrap();
1107        let mut record_iter = parser.into_record_iter();
1108
1109        let peer_index_table = record_iter.next().unwrap();
1110        let _ = elementor.record_to_elems(peer_index_table);
1111
1112        let record = record_iter.next().unwrap();
1113        let elems_vec = elementor.record_to_elems(record.clone());
1114        let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1115        assert_eq!(elems_vec, elems_iter);
1116        assert!(!elems_vec.is_empty());
1117    }
1118
1119    #[test]
1120    fn test_record_to_elems_iter_tabledumpv2_with_peer_table() {
1121        let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1122
1123        let parser = BgpkitParser::new(url).unwrap();
1124        let mut record_iter = parser.into_record_iter();
1125
1126        let peer_index_table = record_iter.next().unwrap();
1127        let mut elementor = Elementor::with_peer_table(
1128            if let MrtMessage::TableDumpV2Message(TableDumpV2Message::PeerIndexTable(pit)) =
1129                peer_index_table.message
1130            {
1131                pit
1132            } else {
1133                panic!("Expected PeerIndexTable");
1134            },
1135        );
1136
1137        let record = record_iter.next().unwrap();
1138        let elems_vec = elementor.record_to_elems(record.clone());
1139        let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1140        assert_eq!(elems_vec, elems_iter);
1141        assert!(!elems_vec.is_empty());
1142    }
1143
1144    #[test]
1145    fn test_record_to_elems_iter_error_unexpected_peer_index_table() {
1146        let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1147
1148        let elementor = Elementor::new();
1149        let parser = BgpkitParser::new(url).unwrap();
1150        let mut record_iter = parser.into_record_iter();
1151        let record = record_iter.next().unwrap();
1152
1153        let result = elementor.record_to_elems_iter(record);
1154        assert!(matches!(
1155            result,
1156            Err(ElemError::UnexpectedPeerIndexTable(_))
1157        ));
1158    }
1159
1160    #[test]
1161    fn test_record_to_elems_iter_error_missing_peer_table() {
1162        // rib-example-small.bz2 is a TableDumpV2 file (starts with PeerIndexTable)
1163        let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1164
1165        let elementor = Elementor::new();
1166        let parser = BgpkitParser::new(url).unwrap();
1167        let mut record_iter = parser.into_record_iter();
1168
1169        // Skip the PeerIndexTable without consuming it via record_to_elems
1170        // which would set the peer table in the elementor
1171        let _peer_index_table = record_iter.next().unwrap();
1172
1173        // Now try to process a RIB entry without having set the peer table
1174        let record = record_iter.next().unwrap();
1175        let result = elementor.record_to_elems_iter(record);
1176        assert!(matches!(result, Err(ElemError::MissingPeerTable)));
1177    }
1178
1179    #[test]
1180    fn test_bgp_to_elems_iter_equivalence() {
1181        let timestamp = 0.0;
1182        let peer_ip = IpAddr::from_str("10.0.0.1").unwrap();
1183        let peer_asn = Asn::new_32bit(65000);
1184
1185        let attributes = vec![
1186            AttributeValue::Origin(Origin::IGP),
1187            AttributeValue::AsPath(AsPath::from_sequence([65000, 65001, 65002])),
1188            AttributeValue::NextHop(peer_ip),
1189        ]
1190        .into_iter()
1191        .map(Attribute::from)
1192        .collect::<Vec<Attribute>>();
1193        let attributes = Attributes::from(attributes);
1194
1195        let announced_prefixes = vec![NetworkPrefix::from_str("10.0.0.0/24").unwrap()];
1196
1197        let bgp_message = BgpMessage::Update(BgpUpdateMessage {
1198            attributes,
1199            announced_prefixes,
1200            withdrawn_prefixes: vec![],
1201        });
1202
1203        let elems_vec =
1204            Elementor::bgp_to_elems(bgp_message.clone(), timestamp, &peer_ip, &peer_asn);
1205        let elems_iter: Vec<BgpElem> =
1206            Elementor::bgp_to_elems_iter(bgp_message, timestamp, &peer_ip, &peer_asn)
1207                .unwrap()
1208                .collect();
1209        assert_eq!(elems_vec, elems_iter);
1210        assert_eq!(elems_vec.len(), 1);
1211    }
1212
1213    #[test]
1214    fn test_bgp_to_elems_iter_non_update_messages() {
1215        use std::net::Ipv4Addr;
1216
1217        let timestamp = 0.0;
1218        let peer_ip = IpAddr::from_str("10.0.0.1").unwrap();
1219        let peer_asn = Asn::new_32bit(65000);
1220
1221        let open_msg = BgpOpenMessage {
1222            version: 4,
1223            asn: Asn::new_32bit(1),
1224            hold_time: 180,
1225            bgp_identifier: Ipv4Addr::new(192, 0, 2, 1),
1226            extended_length: false,
1227            opt_params: vec![],
1228        };
1229        assert!(Elementor::bgp_to_elems_iter(
1230            BgpMessage::Open(open_msg),
1231            timestamp,
1232            &peer_ip,
1233            &peer_asn
1234        )
1235        .is_none());
1236
1237        let notification_msg = BgpNotificationMessage {
1238            error: BgpError::Unknown(0, 0),
1239            data: vec![],
1240        };
1241        assert!(Elementor::bgp_to_elems_iter(
1242            BgpMessage::Notification(notification_msg),
1243            timestamp,
1244            &peer_ip,
1245            &peer_asn
1246        )
1247        .is_none());
1248
1249        assert!(Elementor::bgp_to_elems_iter(
1250            BgpMessage::KeepAlive,
1251            timestamp,
1252            &peer_ip,
1253            &peer_asn
1254        )
1255        .is_none());
1256    }
1257
1258    #[test]
1259    fn test_bgp_update_to_elems_iter_equivalence() {
1260        let timestamp = 0.0;
1261        let peer_ip = IpAddr::from_str("10.0.0.1").unwrap();
1262        let peer_asn = Asn::new_32bit(65000);
1263
1264        let attributes = vec![
1265            AttributeValue::Origin(Origin::IGP),
1266            AttributeValue::AsPath(AsPath::from_sequence([65000, 65001, 65002])),
1267            AttributeValue::NextHop(peer_ip),
1268        ]
1269        .into_iter()
1270        .map(Attribute::from)
1271        .collect::<Vec<Attribute>>();
1272        let attributes = Attributes::from(attributes);
1273
1274        let announced_prefixes = vec![NetworkPrefix::from_str("10.0.0.0/24").unwrap()];
1275        let withdrawn_prefixes = vec![NetworkPrefix::from_str("10.0.1.0/24").unwrap()];
1276
1277        let update = BgpUpdateMessage {
1278            attributes,
1279            announced_prefixes,
1280            withdrawn_prefixes,
1281        };
1282
1283        let elems_vec =
1284            Elementor::bgp_update_to_elems(update.clone(), timestamp, &peer_ip, &peer_asn);
1285        let elems_iter: Vec<BgpElem> =
1286            Elementor::bgp_update_to_elems_iter(update, timestamp, &peer_ip, &peer_asn).collect();
1287        assert_eq!(elems_vec, elems_iter);
1288        assert_eq!(elems_vec.len(), 2);
1289    }
1290
1291    #[test]
1292    fn test_record_elem_iter_size_hint() {
1293        use std::collections::HashMap;
1294
1295        let peer_table = PeerIndexTable {
1296            collector_bgp_id: BgpIdentifier::from_str("10.0.0.1").unwrap(),
1297            view_name: "".to_string(),
1298            id_peer_map: HashMap::new(),
1299            peer_ip_id_map: HashMap::new(),
1300        };
1301
1302        let entries: Vec<RibEntry> = vec![];
1303        let iter = RecordElemIter::RibAfi {
1304            peer_table: &peer_table,
1305            prefix: NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
1306            entries: entries.into_iter(),
1307        };
1308        assert_eq!(iter.size_hint(), (0, Some(0)));
1309
1310        let entries: Vec<RibEntry> = (0..5)
1311            .map(|i| RibEntry {
1312                peer_index: i as u16,
1313                originated_time: 0,
1314                path_id: None,
1315                attributes: Attributes::default(),
1316            })
1317            .collect();
1318        let iter = RecordElemIter::RibAfi {
1319            peer_table: &peer_table,
1320            prefix: NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
1321            entries: entries.into_iter(),
1322        };
1323        assert_eq!(iter.size_hint(), (5, Some(5)));
1324    }
1325}