1#![allow(unused)]
2use 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#[derive(Debug)]
25pub enum ElemError {
26 UnexpectedPeerIndexTable(Box<PeerIndexTable>),
29 MissingPeerTable,
32 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
57macro_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 let next_hop = next_hop.or_else(|| {
168 announced.as_ref().and_then(|v| {
169 v.next_hop.as_ref().map(|h| match h {
170 NextHopAddress::Ipv4(v) => IpAddr::from(*v),
171 NextHopAddress::Ipv6(v) => IpAddr::from(*v),
172 NextHopAddress::Ipv6LinkLocal(v, _) => IpAddr::from(*v),
173 NextHopAddress::VpnIpv6(_, v) => IpAddr::from(*v),
175 NextHopAddress::VpnIpv6LinkLocal(_, v, _, _) => IpAddr::from(*v),
176 })
177 })
178 });
179
180 (
181 as_path,
182 as4_path,
183 origin,
184 next_hop,
185 local_pref,
186 med,
187 communities,
188 atomic,
189 aggregator,
190 announced,
191 withdrawn,
192 otc,
193 if unknown.is_empty() {
194 None
195 } else {
196 Some(unknown)
197 },
198 if deprecated.is_empty() {
199 None
200 } else {
201 Some(deprecated)
202 },
203 )
204}
205
206fn rib_entry_to_elem(prefix: NetworkPrefix, peer: &Peer, entry: RibEntry) -> BgpElem {
207 let (
208 as_path,
209 as4_path,
210 origin,
211 next_hop,
212 local_pref,
213 med,
214 communities,
215 atomic,
216 aggregator,
217 announced,
218 _withdrawn,
219 only_to_customer,
220 unknown,
221 deprecated,
222 ) = get_relevant_attributes(entry.attributes);
223
224 let path = match (as_path, as4_path) {
225 (None, None) => None,
226 (Some(v), None) => Some(v),
227 (None, Some(v)) => Some(v),
228 (Some(v1), Some(v2)) => Some(AsPath::merge_aspath_as4path(&v1, &v2)),
229 };
230
231 let next_hop = match next_hop {
232 Some(v) => Some(v),
233 None => announced.and_then(|v| {
234 v.next_hop.map(|h| match h {
235 NextHopAddress::Ipv4(v) => IpAddr::from(v),
236 NextHopAddress::Ipv6(v) => IpAddr::from(v),
237 NextHopAddress::Ipv6LinkLocal(v, _) => IpAddr::from(v),
238 NextHopAddress::VpnIpv6(_, v) => IpAddr::from(v),
239 NextHopAddress::VpnIpv6LinkLocal(_, v, _, _) => IpAddr::from(v),
240 })
241 }),
242 };
243
244 let origin_asns = path
245 .as_ref()
246 .map(|as_path| as_path.iter_origins().collect());
247
248 BgpElem {
249 timestamp: entry.originated_time as f64,
250 elem_type: ElemType::ANNOUNCE,
251 peer_ip: peer.peer_ip,
252 peer_asn: peer.peer_asn,
253 peer_bgp_id: Some(peer.peer_bgp_id),
254 prefix,
255 next_hop,
256 as_path: path,
257 origin,
258 origin_asns,
259 local_pref,
260 med,
261 communities,
262 atomic,
263 aggr_asn: aggregator.map(|v| v.0),
264 aggr_ip: aggregator.map(|v| v.1),
265 only_to_customer,
266 unknown,
267 deprecated,
268 }
269}
270
271pub enum RecordElemIter<'a> {
277 #[doc(hidden)]
278 Empty,
279 #[doc(hidden)]
280 TableDump(Option<BgpElem>),
281 #[doc(hidden)]
282 TableDumpBatch(std::vec::IntoIter<TableDumpMessage>),
283 #[doc(hidden)]
284 RibAfi {
285 peer_table: &'a PeerIndexTable,
286 prefix: NetworkPrefix,
287 entries: std::vec::IntoIter<RibEntry>,
288 },
289 #[doc(hidden)]
290 Bgp4Mp(BgpUpdateElemIter),
291}
292
293impl Iterator for RecordElemIter<'_> {
294 type Item = BgpElem;
295
296 fn next(&mut self) -> Option<BgpElem> {
297 match self {
298 RecordElemIter::Empty => None,
299 RecordElemIter::TableDump(elem) => elem.take(),
300 RecordElemIter::TableDumpBatch(entries) => entries.next().map(table_dump_to_elem),
301 RecordElemIter::Bgp4Mp(iter) => iter.next(),
302 RecordElemIter::RibAfi {
303 peer_table,
304 prefix,
305 entries,
306 } => {
307 let entry = entries.next()?;
308 let pid = entry.peer_index;
309 match peer_table.get_peer_by_id(&pid) {
310 Some(peer) => Some(rib_entry_to_elem(*prefix, peer, entry)),
311 None => {
312 error!("peer ID {} not found in peer_index table", pid);
313 *self = RecordElemIter::Empty;
314 None
315 }
316 }
317 }
318 }
319 }
320
321 fn size_hint(&self) -> (usize, Option<usize>) {
322 match self {
323 RecordElemIter::Empty => (0, Some(0)),
324 RecordElemIter::TableDump(elem) => {
325 let n = elem.is_some() as usize;
326 (n, Some(n))
327 }
328 RecordElemIter::TableDumpBatch(entries) => {
329 let len = entries.len();
330 (len, Some(len))
331 }
332 RecordElemIter::Bgp4Mp(iter) => iter.size_hint(),
333 RecordElemIter::RibAfi { entries, .. } => {
334 let len = entries.len();
335 (len, Some(len))
336 }
337 }
338 }
339}
340
341pub struct BgpUpdateElemIter {
345 timestamp: f64,
346 peer_ip: IpAddr,
347 peer_asn: Asn,
348 peer_bgp_id: Option<BgpIdentifier>,
349 only_to_customer: Option<Asn>,
350 path: Option<AsPath>,
352 origin_asns: Option<Vec<Asn>>,
353 origin: Option<Origin>,
354 next_hop: Option<IpAddr>,
355 local_pref: Option<u32>,
356 med: Option<u32>,
357 communities: Option<Vec<MetaCommunity>>,
358 atomic: bool,
359 aggr_asn: Option<Asn>,
360 aggr_ip: Option<BgpIdentifier>,
361 unknown: Option<Vec<AttrRaw>>,
362 deprecated: Option<Vec<AttrRaw>>,
363 announced:
365 std::iter::Chain<std::vec::IntoIter<NetworkPrefix>, std::vec::IntoIter<NetworkPrefix>>,
366 withdrawn:
367 std::iter::Chain<std::vec::IntoIter<NetworkPrefix>, std::vec::IntoIter<NetworkPrefix>>,
368 in_withdrawn_phase: bool,
369}
370
371impl Iterator for BgpUpdateElemIter {
372 type Item = BgpElem;
373
374 fn next(&mut self) -> Option<BgpElem> {
375 if !self.in_withdrawn_phase {
376 if let Some(prefix) = self.announced.next() {
377 return Some(BgpElem {
378 timestamp: self.timestamp,
379 elem_type: ElemType::ANNOUNCE,
380 peer_ip: self.peer_ip,
381 peer_asn: self.peer_asn,
382 peer_bgp_id: self.peer_bgp_id,
383 prefix,
384 next_hop: self.next_hop,
385 as_path: self.path.clone(),
386 origin: self.origin,
387 origin_asns: self.origin_asns.clone(),
388 local_pref: self.local_pref,
389 med: self.med,
390 communities: self.communities.clone(),
391 atomic: self.atomic,
392 aggr_asn: self.aggr_asn,
393 aggr_ip: self.aggr_ip,
394 only_to_customer: self.only_to_customer,
395 unknown: self.unknown.clone(),
396 deprecated: self.deprecated.clone(),
397 });
398 }
399 self.in_withdrawn_phase = true;
400 }
401
402 self.withdrawn.next().map(|prefix| BgpElem {
403 timestamp: self.timestamp,
404 elem_type: ElemType::WITHDRAW,
405 peer_ip: self.peer_ip,
406 peer_asn: self.peer_asn,
407 peer_bgp_id: self.peer_bgp_id,
408 prefix,
409 next_hop: None,
410 as_path: None,
411 origin: None,
412 origin_asns: None,
413 local_pref: None,
414 med: None,
415 communities: None,
416 atomic: false,
417 aggr_asn: None,
418 aggr_ip: None,
419 only_to_customer: None,
420 unknown: None,
421 deprecated: None,
422 })
423 }
424
425 fn size_hint(&self) -> (usize, Option<usize>) {
426 let (ann_lo, ann_hi) = if self.in_withdrawn_phase {
427 (0, Some(0))
428 } else {
429 self.announced.size_hint()
430 };
431 let (wd_lo, wd_hi) = self.withdrawn.size_hint();
432 (ann_lo + wd_lo, ann_hi.and_then(|a| wd_hi.map(|w| a + w)))
433 }
434}
435
436impl Elementor {
437 pub fn new() -> Elementor {
438 Self::default()
439 }
440
441 pub fn set_peer_table(&mut self, record: MrtRecord) -> Result<(), ParserError> {
470 if let MrtMessage::TableDumpV2Message(TableDumpV2Message::PeerIndexTable(p)) =
471 record.message
472 {
473 self.peer_table = Some(p);
474 Ok(())
475 } else {
476 Err(ParseError("peer_table is not a PeerIndexTable".to_string()))
477 }
478 }
479
480 pub fn with_peer_table(peer_table: PeerIndexTable) -> Elementor {
482 Elementor {
483 peer_table: Some(peer_table),
484 }
485 }
486
487 pub fn record_to_elems_iter(&self, record: MrtRecord) -> Result<RecordElemIter<'_>, ElemError> {
504 let timestamp = {
505 let t = record.common_header.timestamp;
506 if let Some(micro) = &record.common_header.microsecond_timestamp {
507 let m = (*micro as f64) / 1000000.0;
508 t as f64 + m
509 } else {
510 f64::from(t)
511 }
512 };
513
514 match record.message {
515 MrtMessage::TableDumpMessage(msg) => {
516 Ok(RecordElemIter::TableDump(Some(table_dump_to_elem(msg))))
517 }
518 MrtMessage::TableDumpMessageBatch(messages) => {
519 Ok(RecordElemIter::TableDumpBatch(messages.into_iter()))
520 }
521
522 MrtMessage::TableDumpV2Message(msg) => match msg {
523 TableDumpV2Message::PeerIndexTable(p) => {
524 Err(ElemError::UnexpectedPeerIndexTable(Box::new(p)))
525 }
526 TableDumpV2Message::RibAfi(t) => {
527 let peer_table = self
528 .peer_table
529 .as_ref()
530 .ok_or(ElemError::MissingPeerTable)?;
531 Ok(RecordElemIter::RibAfi {
532 peer_table,
533 prefix: t.prefix,
534 entries: t.rib_entries.into_iter(),
535 })
536 }
537 TableDumpV2Message::RibGeneric(_) => Err(ElemError::UnsupportedRibGeneric),
538 TableDumpV2Message::GeoPeerTable(_) => Ok(RecordElemIter::Empty),
539 },
540
541 MrtMessage::Bgp4Mp(msg) => match msg {
542 Bgp4MpEnum::StateChange(_) => Ok(RecordElemIter::Empty),
543 Bgp4MpEnum::Message(v) => {
544 match Elementor::bgp_to_elems_iter(
545 v.bgp_message,
546 timestamp,
547 &v.peer_ip,
548 &v.peer_asn,
549 ) {
550 Some(iter) => Ok(RecordElemIter::Bgp4Mp(iter)),
551 None => Ok(RecordElemIter::Empty),
552 }
553 }
554 },
555 MrtMessage::LegacyBgp(msg) => match msg {
556 LegacyBgp::StateChange(_) => Ok(RecordElemIter::Empty),
557 LegacyBgp::Message(message) => match Elementor::bgp_to_elems_iter(
558 message.bgp_message,
559 timestamp,
560 &message.peer_ip,
561 &message.peer_asn,
562 ) {
563 Some(iter) => Ok(RecordElemIter::Bgp4Mp(iter)),
564 None => Ok(RecordElemIter::Empty),
565 },
566 },
567 }
568 }
569
570 pub fn bgp_to_elems(
575 msg: BgpMessage,
576 timestamp: f64,
577 peer_ip: &IpAddr,
578 peer_asn: &Asn,
579 ) -> Vec<BgpElem> {
580 Elementor::bgp_to_elems_iter(msg, timestamp, peer_ip, peer_asn)
581 .map(|iter| iter.collect())
582 .unwrap_or_default()
583 }
584
585 pub fn bgp_to_elems_iter(
589 msg: BgpMessage,
590 timestamp: f64,
591 peer_ip: &IpAddr,
592 peer_asn: &Asn,
593 ) -> Option<BgpUpdateElemIter> {
594 match msg {
595 BgpMessage::Update(msg) => Some(Elementor::bgp_update_to_elems_iter(
596 msg, timestamp, peer_ip, peer_asn,
597 )),
598 BgpMessage::Open(_)
599 | BgpMessage::Notification(_)
600 | BgpMessage::KeepAlive
601 | BgpMessage::RouteRefresh(_) => None,
602 }
603 }
604
605 pub fn bgp_update_to_elems(
607 msg: BgpUpdateMessage,
608 timestamp: f64,
609 peer_ip: &IpAddr,
610 peer_asn: &Asn,
611 ) -> Vec<BgpElem> {
612 Elementor::bgp_update_to_elems_iter(msg, timestamp, peer_ip, peer_asn).collect()
613 }
614
615 pub fn bgp_update_to_elems_iter(
618 msg: BgpUpdateMessage,
619 timestamp: f64,
620 peer_ip: &IpAddr,
621 peer_asn: &Asn,
622 ) -> BgpUpdateElemIter {
623 let (
624 as_path,
625 as4_path,
626 origin,
627 next_hop,
628 local_pref,
629 med,
630 communities,
631 atomic,
632 aggregator,
633 announced,
634 withdrawn,
635 only_to_customer,
636 unknown,
637 deprecated,
638 ) = get_relevant_attributes(msg.attributes);
639
640 let path = match (as_path, as4_path) {
641 (None, None) => None,
642 (Some(v), None) => Some(v),
643 (None, Some(v)) => Some(v),
644 (Some(v1), Some(v2)) => Some(AsPath::merge_aspath_as4path(&v1, &v2)),
645 };
646
647 let origin_asns = path
648 .as_ref()
649 .map(|as_path| as_path.iter_origins().collect());
650
651 let nlri_announced = announced.map(|n| n.prefixes).unwrap_or_default();
652 let nlri_withdrawn = withdrawn.map(|n| n.prefixes).unwrap_or_default();
653
654 BgpUpdateElemIter {
655 timestamp,
656 peer_ip: *peer_ip,
657 peer_asn: *peer_asn,
658 peer_bgp_id: None,
659 only_to_customer,
660 path,
661 origin_asns,
662 origin,
663 next_hop,
664 local_pref,
665 med,
666 communities,
667 atomic,
668 aggr_asn: aggregator.as_ref().map(|v| v.0),
669 aggr_ip: aggregator.as_ref().map(|v| v.1),
670 unknown,
671 deprecated,
672 announced: msg.announced_prefixes.into_iter().chain(nlri_announced),
673 withdrawn: msg.withdrawn_prefixes.into_iter().chain(nlri_withdrawn),
674 in_withdrawn_phase: false,
675 }
676 }
677
678 pub fn record_to_elems(&mut self, record: MrtRecord) -> Vec<BgpElem> {
686 match record.message {
687 MrtMessage::TableDumpV2Message(TableDumpV2Message::PeerIndexTable(_)) => {
688 self.set_peer_table(record);
689 vec![]
690 }
691 _ => match self.record_to_elems_iter(record) {
692 Ok(iter) => iter.collect(),
693 Err(e) => {
694 error!("{}", e);
695 vec![]
696 }
697 },
698 }
699 }
700}
701
702fn table_dump_to_elem(msg: TableDumpMessage) -> BgpElem {
703 let (
704 as_path,
705 _as4_path,
706 origin,
707 next_hop,
708 local_pref,
709 med,
710 communities,
711 atomic,
712 aggregator,
713 _announced,
714 _withdrawn,
715 only_to_customer,
716 unknown,
717 deprecated,
718 ) = get_relevant_attributes(msg.attributes);
719
720 let origin_asns = as_path
721 .as_ref()
722 .map(|as_path| as_path.iter_origins().collect());
723
724 BgpElem {
725 timestamp: msg.originated_time as f64,
726 elem_type: ElemType::ANNOUNCE,
727 peer_ip: msg.peer_ip,
728 peer_asn: msg.peer_asn,
729 peer_bgp_id: None,
730 prefix: msg.prefix,
731 next_hop,
732 as_path,
733 origin,
734 origin_asns,
735 local_pref,
736 med,
737 communities,
738 atomic,
739 aggr_asn: aggregator.map(|v| v.0),
740 aggr_ip: aggregator.map(|v| v.1),
741 only_to_customer,
742 unknown,
743 deprecated,
744 }
745}
746
747#[inline(always)]
748pub fn option_to_string<T>(o: &Option<T>) -> String
749where
750 T: Display,
751{
752 if let Some(v) = o {
753 v.to_string()
754 } else {
755 String::new()
756 }
757}
758
759impl From<&BgpElem> for Attributes {
760 fn from(value: &BgpElem) -> Self {
761 let mut values = Vec::<AttributeValue>::new();
762 let mut attributes = Attributes::default();
763 let prefix = value.prefix;
764
765 if value.elem_type == ElemType::WITHDRAW {
766 values.push(AttributeValue::MpUnreachNlri(Nlri::new_unreachable(prefix)));
767 attributes.extend(values);
768 return attributes;
769 }
770
771 values.push(AttributeValue::MpReachNlri(Nlri::new_reachable(
772 prefix,
773 value.next_hop,
774 )));
775
776 if let Some(v) = value.next_hop {
777 values.push(AttributeValue::NextHop(v));
778 }
779
780 if let Some(v) = value.as_path.as_ref() {
781 values.push(AttributeValue::AsPath(v.clone()));
785 }
786
787 if let Some(v) = value.origin {
788 values.push(AttributeValue::Origin(v));
789 }
790
791 if let Some(v) = value.local_pref {
792 values.push(AttributeValue::LocalPreference(v));
793 }
794
795 if let Some(v) = value.med {
796 values.push(AttributeValue::MultiExitDiscriminator(v));
797 }
798
799 if let Some(v) = value.communities.as_ref() {
800 let mut communites = vec![];
801 let mut extended_communities = vec![];
802 let mut ipv6_extended_communities = vec![];
803 let mut large_communities = vec![];
804 for c in v {
805 match c {
806 MetaCommunity::Plain(v) => communites.push(*v),
807 MetaCommunity::Extended(v) => extended_communities.push(*v),
808 MetaCommunity::Large(v) => large_communities.push(*v),
809 MetaCommunity::Ipv6Extended(v) => ipv6_extended_communities.push(*v),
810 }
811 }
812 if !communites.is_empty() {
813 values.push(AttributeValue::Communities(communites));
814 }
815 if !extended_communities.is_empty() {
816 values.push(AttributeValue::ExtendedCommunities(extended_communities));
817 }
818 if !large_communities.is_empty() {
819 values.push(AttributeValue::LargeCommunities(large_communities));
820 }
821 if !ipv6_extended_communities.is_empty() {
822 values.push(AttributeValue::Ipv6AddressSpecificExtendedCommunities(
823 ipv6_extended_communities,
824 ));
825 }
826 }
827
828 if let Some(v) = value.aggr_asn {
829 let aggregator_id = match value.aggr_ip {
830 Some(v) => v,
831 None => Ipv4Addr::UNSPECIFIED,
832 };
833 values.push(AttributeValue::Aggregator {
834 asn: v,
835 id: aggregator_id,
836 });
837 }
838
839 if let Some(v) = value.only_to_customer {
840 values.push(AttributeValue::OnlyToCustomer(v));
841 }
842
843 if let Some(v) = value.unknown.as_ref() {
844 for t in v {
845 values.push(AttributeValue::Unknown(t.clone()));
846 }
847 }
848
849 if let Some(v) = value.deprecated.as_ref() {
850 for t in v {
851 values.push(AttributeValue::Deprecated(t.clone()));
852 }
853 }
854
855 attributes.extend(values);
856 attributes
857 }
858}
859
860#[cfg(test)]
861mod tests {
862 use super::*;
863 use crate::BgpkitParser;
864 use std::net::{Ipv4Addr, Ipv6Addr};
865 use std::str::FromStr;
866
867 #[test]
868 fn test_option_to_string() {
869 let o1 = Some(1);
870 let o2: Option<u32> = None;
871 assert_eq!(option_to_string(&o1), "1");
872 assert_eq!(option_to_string(&o2), "");
873 }
874
875 #[test]
876 fn test_record_to_elems() {
877 let url_table_dump_v1 = "https://data.ris.ripe.net/rrc00/2003.01/bview.20030101.0000.gz";
878 let url_table_dump_v2 = "https://data.ris.ripe.net/rrc00/2023.01/bview.20230101.0000.gz";
879 let url_bgp4mp = "https://data.ris.ripe.net/rrc00/2021.10/updates.20211001.0000.gz";
880
881 let mut elementor = Elementor::new();
882 let parser = BgpkitParser::new(url_table_dump_v1).unwrap();
883 let mut record_iter = parser.into_record_iter();
884 let record = record_iter.next().unwrap();
885 let elems = elementor.record_to_elems(record);
886 assert_eq!(elems.len(), 1);
887
888 let parser = BgpkitParser::new(url_table_dump_v2).unwrap();
889 let mut record_iter = parser.into_record_iter();
890 let peer_index_table = record_iter.next().unwrap();
891 let _elems = elementor.record_to_elems(peer_index_table);
892 let record = record_iter.next().unwrap();
893 let elems = elementor.record_to_elems(record);
894 assert!(!elems.is_empty());
895
896 let parser = BgpkitParser::new(url_bgp4mp).unwrap();
897 let mut record_iter = parser.into_record_iter();
898 let record = record_iter.next().unwrap();
899 let elems = elementor.record_to_elems(record);
900 assert!(!elems.is_empty());
901 }
902
903 #[test]
904 fn test_attributes_from_bgp_elem() {
905 let mut elem = BgpElem {
906 timestamp: 0.0,
907 elem_type: ElemType::ANNOUNCE,
908 peer_ip: IpAddr::from_str("10.0.0.1").unwrap(),
909 peer_asn: Asn::new_32bit(65000),
910 peer_bgp_id: None,
911 prefix: NetworkPrefix::from_str("10.0.1.0/24").unwrap(),
912 next_hop: Some(IpAddr::from_str("10.0.0.2").unwrap()),
913 as_path: Some(AsPath::from_sequence([65000, 65001, 65002])),
914 origin: Some(Origin::EGP),
915 origin_asns: Some(vec![Asn::new_32bit(65000)]),
916 local_pref: Some(100),
917 med: Some(200),
918 communities: Some(vec![
919 MetaCommunity::Plain(Community::NoAdvertise),
920 MetaCommunity::Extended(ExtendedCommunity::Raw([0, 0, 0, 0, 0, 0, 0, 0])),
921 MetaCommunity::Large(LargeCommunity {
922 global_admin: 0,
923 local_data: [0, 0],
924 }),
925 MetaCommunity::Ipv6Extended(Ipv6AddrExtCommunity {
926 community_type: ExtendedCommunityType::TransitiveTwoOctetAs,
927 subtype: 0,
928 global_admin: Ipv6Addr::from_str("2001:db8::").unwrap(),
929 local_admin: [0, 0],
930 }),
931 ]),
932 atomic: false,
933 aggr_asn: Some(Asn::new_32bit(65000)),
934 aggr_ip: Some(Ipv4Addr::from_str("10.2.0.0").unwrap()),
935 only_to_customer: Some(Asn::new_32bit(65000)),
936 unknown: Some(vec![AttrRaw {
937 code: AttrType::RESERVED.into(),
938 bytes: Bytes::new(),
939 }]),
940 deprecated: Some(vec![AttrRaw {
941 code: AttrType::RESERVED.into(),
942 bytes: Bytes::new(),
943 }]),
944 };
945
946 let _attributes = Attributes::from(&elem);
947 elem.elem_type = ElemType::WITHDRAW;
948 let _attributes = Attributes::from(&elem);
949 }
950
951 #[test]
952 fn test_get_relevant_attributes() {
953 let attributes = vec![
954 AttributeValue::Origin(Origin::IGP),
955 AttributeValue::As4Path(AsPath::from_sequence([65000, 65001, 65002])),
956 AttributeValue::NextHop(IpAddr::from_str("10.0.0.1").unwrap()),
957 AttributeValue::MultiExitDiscriminator(100),
958 AttributeValue::LocalPreference(200),
959 AttributeValue::AtomicAggregate,
960 AttributeValue::Aggregator {
961 asn: Asn::new_32bit(65000),
962 id: Ipv4Addr::from_str("10.0.0.1").unwrap(),
963 },
964 AttributeValue::Communities(vec![Community::NoExport]),
965 AttributeValue::ExtendedCommunities(vec![ExtendedCommunity::Raw([
966 0, 0, 0, 0, 0, 0, 0, 0,
967 ])]),
968 AttributeValue::LargeCommunities(vec![LargeCommunity {
969 global_admin: 0,
970 local_data: [0, 0],
971 }]),
972 AttributeValue::Ipv6AddressSpecificExtendedCommunities(vec![Ipv6AddrExtCommunity {
973 community_type: ExtendedCommunityType::TransitiveTwoOctetAs,
974 subtype: 0,
975 global_admin: Ipv6Addr::from_str("2001:db8::").unwrap(),
976 local_admin: [0, 0],
977 }]),
978 AttributeValue::MpReachNlri(Nlri::new_reachable(
979 NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
980 Some(IpAddr::from_str("10.0.0.1").unwrap()),
981 )),
982 AttributeValue::MpUnreachNlri(Nlri::new_unreachable(
983 NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
984 )),
985 AttributeValue::OnlyToCustomer(Asn::new_32bit(65000)),
986 AttributeValue::Unknown(AttrRaw {
987 code: AttrType::RESERVED.into(),
988 bytes: Bytes::new(),
989 }),
990 AttributeValue::Deprecated(AttrRaw {
991 code: AttrType::RESERVED.into(),
992 bytes: Bytes::new(),
993 }),
994 ]
995 .into_iter()
996 .map(Attribute::from)
997 .collect::<Vec<Attribute>>();
998
999 let attributes = Attributes::from(attributes);
1000
1001 let (
1002 _as_path,
1003 _as4_path, _origin,
1005 _next_hop,
1006 _local_pref,
1007 _med,
1008 _communities,
1009 _atomic,
1010 _aggregator,
1011 _announced,
1012 _withdrawn,
1013 _only_to_customer,
1014 _unknown,
1015 _deprecated,
1016 ) = get_relevant_attributes(attributes);
1017 }
1018
1019 #[test]
1020 fn test_next_hop_from_nlri() {
1021 let attributes = vec![AttributeValue::NextHop(
1022 IpAddr::from_str("10.0.0.1").unwrap(),
1023 )]
1024 .into_iter()
1025 .map(Attribute::from)
1026 .collect::<Vec<Attribute>>();
1027
1028 let attributes = Attributes::from(attributes);
1029
1030 let (
1031 _as_path,
1032 _as4_path, _origin,
1034 next_hop,
1035 _local_pref,
1036 _med,
1037 _communities,
1038 _atomic,
1039 _aggregator,
1040 _announced,
1041 _withdrawn,
1042 _only_to_customer,
1043 _unknown,
1044 _deprecated,
1045 ) = get_relevant_attributes(attributes);
1046
1047 assert_eq!(next_hop, Some(IpAddr::from_str("10.0.0.1").unwrap()));
1048
1049 let attributes = vec![AttributeValue::MpReachNlri(Nlri::new_reachable(
1050 NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
1051 Some(IpAddr::from_str("10.0.0.2").unwrap()),
1052 ))]
1053 .into_iter()
1054 .map(Attribute::from)
1055 .collect::<Vec<Attribute>>();
1056
1057 let attributes = Attributes::from(attributes);
1058
1059 let (
1060 _as_path,
1061 _as4_path, _origin,
1063 next_hop,
1064 _local_pref,
1065 _med,
1066 _communities,
1067 _atomic,
1068 _aggregator,
1069 _announced,
1070 _withdrawn,
1071 _only_to_customer,
1072 _unknown,
1073 _deprecated,
1074 ) = get_relevant_attributes(attributes);
1075
1076 assert_eq!(next_hop, Some(IpAddr::from_str("10.0.0.2").unwrap()));
1077 }
1078
1079 #[test]
1080 fn test_record_to_elems_iter_equivalence_tabledumpv2_small() {
1081 let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1083
1084 let mut elementor = Elementor::new();
1085 let parser = BgpkitParser::new(url).unwrap();
1086 let mut record_iter = parser.into_record_iter();
1087
1088 let peer_index_table = record_iter.next().unwrap();
1090 let _ = elementor.record_to_elems(peer_index_table);
1091
1092 let record = record_iter.next().unwrap();
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 fn test_record_to_elems_iter_equivalence_bgp4mp() {
1102 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
1103
1104 let mut elementor = Elementor::new();
1105 let parser = BgpkitParser::new(url).unwrap();
1106 let mut record_iter = parser.into_record_iter();
1107 let record = record_iter.next().unwrap();
1108
1109 let elems_vec = elementor.record_to_elems(record.clone());
1110 let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1111 assert_eq!(elems_vec, elems_iter);
1112 assert!(!elems_vec.is_empty());
1113 }
1114
1115 #[test]
1116 #[ignore = "requires large RIB file download"]
1117 fn test_record_to_elems_iter_equivalence_tabledumpv2() {
1118 let url = "https://data.ris.ripe.net/rrc00/2023.01/bview.20230101.0000.gz";
1119
1120 let mut elementor = Elementor::new();
1121 let parser = BgpkitParser::new(url).unwrap();
1122 let mut record_iter = parser.into_record_iter();
1123
1124 let peer_index_table = record_iter.next().unwrap();
1125 let _ = elementor.record_to_elems(peer_index_table);
1126
1127 let record = record_iter.next().unwrap();
1128 let elems_vec = elementor.record_to_elems(record.clone());
1129 let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1130 assert_eq!(elems_vec, elems_iter);
1131 assert!(!elems_vec.is_empty());
1132 }
1133
1134 #[test]
1135 fn test_record_to_elems_iter_tabledumpv2_with_peer_table() {
1136 let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1137
1138 let parser = BgpkitParser::new(url).unwrap();
1139 let mut record_iter = parser.into_record_iter();
1140
1141 let peer_index_table = record_iter.next().unwrap();
1142 let mut elementor = Elementor::with_peer_table(
1143 if let MrtMessage::TableDumpV2Message(TableDumpV2Message::PeerIndexTable(pit)) =
1144 peer_index_table.message
1145 {
1146 pit
1147 } else {
1148 panic!("Expected PeerIndexTable");
1149 },
1150 );
1151
1152 let record = record_iter.next().unwrap();
1153 let elems_vec = elementor.record_to_elems(record.clone());
1154 let elems_iter: Vec<BgpElem> = elementor.record_to_elems_iter(record).unwrap().collect();
1155 assert_eq!(elems_vec, elems_iter);
1156 assert!(!elems_vec.is_empty());
1157 }
1158
1159 #[test]
1160 fn test_record_to_elems_iter_error_unexpected_peer_index_table() {
1161 let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1162
1163 let elementor = Elementor::new();
1164 let parser = BgpkitParser::new(url).unwrap();
1165 let mut record_iter = parser.into_record_iter();
1166 let record = record_iter.next().unwrap();
1167
1168 let result = elementor.record_to_elems_iter(record);
1169 assert!(matches!(
1170 result,
1171 Err(ElemError::UnexpectedPeerIndexTable(_))
1172 ));
1173 }
1174
1175 #[test]
1176 fn test_record_to_elems_iter_error_missing_peer_table() {
1177 let url = "https://spaces.bgpkit.org/parser/rib-example-small.bz2";
1179
1180 let elementor = Elementor::new();
1181 let parser = BgpkitParser::new(url).unwrap();
1182 let mut record_iter = parser.into_record_iter();
1183
1184 let _peer_index_table = record_iter.next().unwrap();
1187
1188 let record = record_iter.next().unwrap();
1190 let result = elementor.record_to_elems_iter(record);
1191 assert!(matches!(result, Err(ElemError::MissingPeerTable)));
1192 }
1193
1194 #[test]
1195 fn test_bgp_to_elems_iter_equivalence() {
1196 let timestamp = 0.0;
1197 let peer_ip = IpAddr::from_str("10.0.0.1").unwrap();
1198 let peer_asn = Asn::new_32bit(65000);
1199
1200 let attributes = vec![
1201 AttributeValue::Origin(Origin::IGP),
1202 AttributeValue::AsPath(AsPath::from_sequence([65000, 65001, 65002])),
1203 AttributeValue::NextHop(peer_ip),
1204 ]
1205 .into_iter()
1206 .map(Attribute::from)
1207 .collect::<Vec<Attribute>>();
1208 let attributes = Attributes::from(attributes);
1209
1210 let announced_prefixes = vec![NetworkPrefix::from_str("10.0.0.0/24").unwrap()];
1211
1212 let bgp_message = BgpMessage::Update(BgpUpdateMessage {
1213 attributes,
1214 announced_prefixes,
1215 withdrawn_prefixes: vec![],
1216 });
1217
1218 let elems_vec =
1219 Elementor::bgp_to_elems(bgp_message.clone(), timestamp, &peer_ip, &peer_asn);
1220 let elems_iter: Vec<BgpElem> =
1221 Elementor::bgp_to_elems_iter(bgp_message, timestamp, &peer_ip, &peer_asn)
1222 .unwrap()
1223 .collect();
1224 assert_eq!(elems_vec, elems_iter);
1225 assert_eq!(elems_vec.len(), 1);
1226 }
1227
1228 #[test]
1229 fn test_bgp_to_elems_iter_non_update_messages() {
1230 use std::net::Ipv4Addr;
1231
1232 let timestamp = 0.0;
1233 let peer_ip = IpAddr::from_str("10.0.0.1").unwrap();
1234 let peer_asn = Asn::new_32bit(65000);
1235
1236 let open_msg = BgpOpenMessage {
1237 version: 4,
1238 asn: Asn::new_32bit(1),
1239 hold_time: 180,
1240 bgp_identifier: Ipv4Addr::new(192, 0, 2, 1),
1241 extended_length: false,
1242 opt_params: vec![],
1243 };
1244 assert!(Elementor::bgp_to_elems_iter(
1245 BgpMessage::Open(open_msg),
1246 timestamp,
1247 &peer_ip,
1248 &peer_asn
1249 )
1250 .is_none());
1251
1252 let notification_msg = BgpNotificationMessage {
1253 error: BgpError::Unknown(0, 0),
1254 data: vec![],
1255 };
1256 assert!(Elementor::bgp_to_elems_iter(
1257 BgpMessage::Notification(notification_msg),
1258 timestamp,
1259 &peer_ip,
1260 &peer_asn
1261 )
1262 .is_none());
1263
1264 assert!(Elementor::bgp_to_elems_iter(
1265 BgpMessage::KeepAlive,
1266 timestamp,
1267 &peer_ip,
1268 &peer_asn
1269 )
1270 .is_none());
1271 }
1272
1273 #[test]
1274 fn test_bgp_update_to_elems_iter_equivalence() {
1275 let timestamp = 0.0;
1276 let peer_ip = IpAddr::from_str("10.0.0.1").unwrap();
1277 let peer_asn = Asn::new_32bit(65000);
1278
1279 let attributes = vec![
1280 AttributeValue::Origin(Origin::IGP),
1281 AttributeValue::AsPath(AsPath::from_sequence([65000, 65001, 65002])),
1282 AttributeValue::NextHop(peer_ip),
1283 ]
1284 .into_iter()
1285 .map(Attribute::from)
1286 .collect::<Vec<Attribute>>();
1287 let attributes = Attributes::from(attributes);
1288
1289 let announced_prefixes = vec![NetworkPrefix::from_str("10.0.0.0/24").unwrap()];
1290 let withdrawn_prefixes = vec![NetworkPrefix::from_str("10.0.1.0/24").unwrap()];
1291
1292 let update = BgpUpdateMessage {
1293 attributes,
1294 announced_prefixes,
1295 withdrawn_prefixes,
1296 };
1297
1298 let elems_vec =
1299 Elementor::bgp_update_to_elems(update.clone(), timestamp, &peer_ip, &peer_asn);
1300 let elems_iter: Vec<BgpElem> =
1301 Elementor::bgp_update_to_elems_iter(update, timestamp, &peer_ip, &peer_asn).collect();
1302 assert_eq!(elems_vec, elems_iter);
1303 assert_eq!(elems_vec.len(), 2);
1304 }
1305
1306 #[test]
1307 fn test_record_elem_iter_size_hint() {
1308 use std::collections::HashMap;
1309
1310 let peer_table = PeerIndexTable {
1311 collector_bgp_id: BgpIdentifier::from_str("10.0.0.1").unwrap(),
1312 view_name: "".to_string(),
1313 id_peer_map: HashMap::new(),
1314 peer_ip_id_map: HashMap::new(),
1315 };
1316
1317 let entries: Vec<RibEntry> = vec![];
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(), (0, Some(0)));
1324
1325 let entries: Vec<RibEntry> = (0..5)
1326 .map(|i| RibEntry {
1327 peer_index: i as u16,
1328 originated_time: 0,
1329 path_id: None,
1330 attributes: Attributes::default(),
1331 })
1332 .collect();
1333 let iter = RecordElemIter::RibAfi {
1334 peer_table: &peer_table,
1335 prefix: NetworkPrefix::from_str("10.0.0.0/24").unwrap(),
1336 entries: entries.into_iter(),
1337 };
1338 assert_eq!(iter.size_hint(), (5, Some(5)));
1339 }
1340}