1use crate::models::{BgpElem, MrtRecord};
5use log::warn;
6use std::io::{BufReader, Cursor, Read};
7pub use text_dump::{detect_text_dump, infer_timestamp_from_path, TextDumpElemIterator};
8
9#[macro_use]
10pub mod utils;
11pub mod bgp;
12pub mod bmp;
13pub mod filter;
14pub mod iters;
15pub mod mrt;
16pub mod rpki;
17pub mod text_dump;
18
19#[cfg(feature = "rislive")]
20pub mod rislive;
21
22pub(crate) use self::utils::*;
23
24pub use mrt::mrt_elem::{BgpUpdateElemIter, ElemError, Elementor, RecordElemIter};
25#[cfg(feature = "oneio")]
26use oneio::{get_cache_reader, get_reader, get_resumable_http_reader};
27
28pub use crate::error::{ParserError, ParserErrorWithBytes};
29pub use bmp::{parse_bmp_msg, parse_openbmp_header, parse_openbmp_msg};
30pub use filter::*;
31pub use iters::*;
32pub use mrt::*;
33
34#[cfg(feature = "rislive")]
35pub use rislive::messages::{
36 RisLiveClientMessage, RisLiveMeta, RisLiveRawFull, RisSubscribe, RisSubscribeSocketOptions,
37 RisSubscribeType,
38};
39#[cfg(feature = "rislive")]
40pub use rislive::{
41 parse_ris_live_message, parse_ris_live_message_json, parse_ris_live_message_raw,
42 parse_ris_live_message_raw_full,
43};
44
45pub struct BgpkitParser<R> {
46 reader: R,
47 core_dump: bool,
48 filters: Vec<Filter>,
49 options: ParserOptions,
50 text_dump_iter: Option<Box<dyn Iterator<Item = BgpElem> + Send>>,
53}
54
55pub(crate) struct ParserOptions {
56 show_warnings: bool,
57 warned_zebra_compat: bool,
58}
59impl Default for ParserOptions {
60 fn default() -> Self {
61 ParserOptions {
62 show_warnings: true,
63 warned_zebra_compat: false,
64 }
65 }
66}
67
68impl ParserOptions {
69 pub(crate) fn warn_zebra_compat_once(&mut self) {
70 if self.show_warnings && !self.warned_zebra_compat {
71 warn!(
72 "recovered shortened Zebra BGP4MP records with missing envelope fields; substituting IPv4 zero addresses and interface index 0 (further occurrences for this parser will not be logged)"
73 );
74 self.warned_zebra_compat = true;
75 }
76 }
77}
78
79#[cfg(feature = "oneio")]
80impl BgpkitParser<Box<dyn Read + Send>> {
81 pub fn new(path: &str) -> Result<Self, ParserErrorWithBytes> {
83 let reader = get_reader(path)?;
84 Ok(BgpkitParser {
85 reader,
86 core_dump: false,
87 filters: vec![],
88 options: ParserOptions::default(),
89 text_dump_iter: None,
90 })
91 }
92
93 pub fn new_resumable_http(path: &str) -> Result<Self, ParserErrorWithBytes> {
115 let reader = get_resumable_http_reader(path)?;
116 Ok(BgpkitParser {
117 reader,
118 core_dump: false,
119 filters: vec![],
120 options: ParserOptions::default(),
121 text_dump_iter: None,
122 })
123 }
124
125 pub fn new_cached(path: &str, cache_dir: &str) -> Result<Self, ParserErrorWithBytes> {
131 let file_name = path.rsplit('/').next().unwrap().to_string();
132 let new_file_name = format!(
133 "cache-{}",
134 add_suffix_to_filename(file_name.as_str(), crc32(path).as_str())
135 );
136 let reader = get_cache_reader(path, cache_dir, Some(new_file_name), false)?;
137 Ok(BgpkitParser {
138 reader,
139 core_dump: false,
140 filters: vec![],
141 options: ParserOptions::default(),
142 text_dump_iter: None,
143 })
144 }
145
146 pub fn new_text(path: &str) -> Result<Self, ParserErrorWithBytes> {
169 let timestamp = infer_timestamp_from_path(path).unwrap_or(0.0);
170 let reader = get_reader(path)?;
171 Self::from_text_reader_with_timestamp(reader, timestamp)
172 }
173
174 pub fn new_auto(path: &str) -> Result<Self, ParserErrorWithBytes> {
193 let reader = get_reader(path)?;
194 Self::from_auto_reader_with_timestamp(reader, infer_timestamp_from_path(path))
195 }
196}
197
198#[cfg(feature = "oneio")]
199fn add_suffix_to_filename(filename: &str, suffix: &str) -> String {
200 let mut parts: Vec<&str> = filename.split('.').collect(); if parts.len() > 1 {
202 let last_part = parts.pop().unwrap(); let new_last_part = format!("{suffix}.{last_part}"); parts.push(&new_last_part); parts.join(".") } else {
207 format!("{filename}.{suffix}")
209 }
210}
211
212impl<R: Read> BgpkitParser<R> {
213 pub fn from_reader(reader: R) -> Self {
215 BgpkitParser {
216 reader,
217 core_dump: false,
218 filters: vec![],
219 options: ParserOptions::default(),
220 text_dump_iter: None,
221 }
222 }
223
224 pub fn next_record(&mut self) -> Result<MrtRecord, ParserErrorWithBytes> {
226 if self.text_dump_iter.is_some() {
227 return Err(ParserError::Unsupported(
228 "text-dump parsers have no MRT record representation; iterate elements instead"
229 .to_string(),
230 )
231 .into());
232 }
233 let (record, used_zebra_compat) =
234 mrt::mrt_record::parse_mrt_record_with_zebra_compat(&mut self.reader)?;
235 if used_zebra_compat {
236 self.warn_zebra_compat_once();
237 }
238 Ok(record)
239 }
240}
241
242impl BgpkitParser<Box<dyn Read + Send>> {
243 pub fn from_text_reader(
247 reader: impl Read + Send + 'static,
248 ) -> Result<Self, ParserErrorWithBytes> {
249 Self::from_text_reader_with_timestamp(reader, 0.0)
250 }
251
252 pub fn from_text_reader_with_timestamp(
256 reader: impl Read + Send + 'static,
257 timestamp: f64,
258 ) -> Result<Self, ParserErrorWithBytes> {
259 let buf_reader = BufReader::new(reader);
260 let iter = TextDumpElemIterator::new(buf_reader, timestamp).map_err(ParserError::from)?;
261 Ok(BgpkitParser {
262 reader: Box::new(std::io::empty()),
263 core_dump: false,
264 filters: vec![],
265 options: ParserOptions::default(),
266 text_dump_iter: Some(Box::new(iter)),
267 })
268 }
269
270 pub fn from_auto_reader(
273 reader: impl Read + Send + 'static,
274 ) -> Result<Self, ParserErrorWithBytes> {
275 Self::from_auto_reader_with_timestamp(reader, None)
276 }
277
278 pub fn from_auto_reader_with_timestamp(
282 reader: impl Read + Send + 'static,
283 timestamp: Option<f64>,
284 ) -> Result<Self, ParserErrorWithBytes> {
285 let mut buf_reader = BufReader::new(reader);
286 let (is_text, head) = detect_text_dump(&mut buf_reader).map_err(ParserError::from)?;
287 if is_text {
288 let chained = BufReader::new(Cursor::new(head).chain(buf_reader));
289 let ts = timestamp.unwrap_or(0.0);
290 let iter = TextDumpElemIterator::new(chained, ts).map_err(ParserError::from)?;
291 Ok(BgpkitParser {
292 reader: Box::new(std::io::empty()),
293 core_dump: false,
294 filters: vec![],
295 options: ParserOptions::default(),
296 text_dump_iter: Some(Box::new(iter)),
297 })
298 } else {
299 Ok(BgpkitParser {
300 reader: Box::new(Cursor::new(head).chain(buf_reader)),
301 core_dump: false,
302 filters: vec![],
303 options: ParserOptions::default(),
304 text_dump_iter: None,
305 })
306 }
307 }
308}
309
310impl<R> BgpkitParser<R> {
311 pub(crate) fn warn_zebra_compat_once(&mut self) {
312 self.options.warn_zebra_compat_once();
313 }
314
315 pub fn enable_core_dump(self) -> Self {
316 BgpkitParser {
317 reader: self.reader,
318 core_dump: true,
319 filters: self.filters,
320 options: self.options,
321 text_dump_iter: self.text_dump_iter,
322 }
323 }
324
325 pub fn disable_warnings(self) -> Self {
326 let mut options = self.options;
327 options.show_warnings = false;
328 BgpkitParser {
329 reader: self.reader,
330 core_dump: self.core_dump,
331 filters: self.filters,
332 options,
333 text_dump_iter: self.text_dump_iter,
334 }
335 }
336
337 pub fn add_filter(
401 self,
402 filter_type: &str,
403 filter_value: &str,
404 ) -> Result<Self, ParserErrorWithBytes> {
405 let mut filters = self.filters;
406 filters.push(Filter::new(filter_type, filter_value)?);
407 Ok(BgpkitParser {
408 reader: self.reader,
409 core_dump: self.core_dump,
410 filters,
411 options: self.options,
412 text_dump_iter: self.text_dump_iter,
413 })
414 }
415
416 pub fn add_filters(mut self, filters: &[Filter]) -> Self {
436 self.filters.extend(filters.iter().cloned());
437 self
438 }
439
440 pub fn with_filters(mut self, filters: &[Filter]) -> Self {
468 self.filters = filters.to_vec();
469 self
470 }
471}
472
473#[cfg(test)]
474mod tests {
475 use super::*;
476 use crate::models::Asn;
477
478 #[test]
479 fn test_new_with_reader() {
480 let reader = oneio::get_reader("http://archive.routeviews.org/route-views.ny/bgpdata/2023.02/UPDATES/updates.20230215.0630.bz2").unwrap();
482 assert_eq!(
483 12683,
484 BgpkitParser::from_reader(reader).into_elem_iter().count()
485 );
486
487 let reader = oneio::get_reader("https://spaces.bgpkit.org/parser/update-example").unwrap();
489 assert_eq!(
490 8160,
491 BgpkitParser::from_reader(reader).into_elem_iter().count()
492 );
493 }
494
495 #[test]
496 fn test_new_resumable_http() {
497 let parser =
498 BgpkitParser::new_resumable_http("https://spaces.bgpkit.org/parser/update-example.gz")
499 .unwrap();
500 assert_eq!(8160, parser.into_elem_iter().count());
501 }
502
503 #[test]
504 fn test_new_cached_with_reader() {
505 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
506 let parser = BgpkitParser::new_cached(url, "/tmp/bgpkit-parser-tests")
507 .unwrap()
508 .enable_core_dump()
509 .disable_warnings();
510 let count = parser.into_elem_iter().count();
511 assert_eq!(8160, count);
512 let parser = BgpkitParser::new_cached(url, "/tmp/bgpkit-parser-tests").unwrap();
513 let count = parser.into_elem_iter().count();
514 assert_eq!(8160, count);
515 }
516
517 #[test]
518 fn test_add_suffix_to_filename() {
519 let filename = "example.txt";
521 let suffix = "suffix";
522 let result = add_suffix_to_filename(filename, suffix);
523 assert_eq!(result, "example.suffix.txt");
524
525 let filename = "example.tar.gz";
527 let suffix = "suffix";
528 let result = add_suffix_to_filename(filename, suffix);
529 assert_eq!(result, "example.tar.suffix.gz");
530
531 let filename = "example";
533 let suffix = "suffix";
534 let result = add_suffix_to_filename(filename, suffix);
535 assert_eq!(result, "example.suffix");
536
537 let filename = "";
539 let suffix = "suffix";
540 let result = add_suffix_to_filename(filename, suffix);
541 assert_eq!(result, ".suffix");
542
543 let filename = "example.txt";
545 let suffix = "";
546 let result = add_suffix_to_filename(filename, suffix);
547 assert_eq!(result, "example..txt");
548 }
549
550 #[test]
551 fn test_with_filters() {
552 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
553
554 let filters = vec![
556 Filter::new("peer_ip", "185.1.8.65").unwrap(),
557 Filter::new("type", "w").unwrap(),
558 ];
559
560 let parser = BgpkitParser::new(url).unwrap().with_filters(&filters);
562 let count = parser.into_elem_iter().count();
563
564 assert_eq!(count, 132);
566
567 let filters1 = vec![Filter::new("peer_ip", "185.1.8.65").unwrap()];
569 let filters2 = vec![Filter::new("peer_ip", "185.1.8.50").unwrap()];
570
571 let parser = BgpkitParser::new(url)
572 .unwrap()
573 .with_filters(&filters1)
574 .with_filters(&filters2); let count = parser.into_elem_iter().count();
576
577 assert_eq!(count, 1563);
579 }
580
581 #[test]
582 fn test_add_filters() {
583 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
584
585 let filters = vec![
587 Filter::new("peer_ip", "185.1.8.65").unwrap(),
588 Filter::new("type", "w").unwrap(),
589 ];
590
591 let parser = BgpkitParser::new(url).unwrap().add_filters(&filters);
593 let count = parser.into_elem_iter().count();
594
595 assert_eq!(count, 132);
597
598 let parser = BgpkitParser::new(url)
600 .unwrap()
601 .add_filter("peer_ip", "185.1.8.65")
602 .unwrap()
603 .add_filters(&[Filter::new("type", "w").unwrap()]);
604 let count = parser.into_elem_iter().count();
605 assert_eq!(count, 132);
606 }
607
608 #[test]
609 fn test_with_filters_empty() {
610 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
611
612 let parser = BgpkitParser::new(url).unwrap().with_filters(&[]);
614 let count = parser.into_elem_iter().count();
615
616 assert_eq!(count, 8160);
618 }
619
620 #[test]
621 fn test_add_filters_empty() {
622 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
623
624 let parser = BgpkitParser::new(url)
626 .unwrap()
627 .add_filter("peer_ip", "185.1.8.65")
628 .unwrap()
629 .add_filters(&[]);
630 let count = parser.into_elem_iter().count();
631
632 assert_eq!(count, 3393);
634 }
635
636 #[test]
637 fn test_with_filters_reuse() {
638 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
639
640 let filters = vec![
642 Filter::new("peer_ip", "185.1.8.65").unwrap(),
643 Filter::new("type", "w").unwrap(),
644 ];
645
646 let parser1 = BgpkitParser::new(url).unwrap().with_filters(&filters);
648 let count1 = parser1.into_elem_iter().count();
649
650 let parser2 = BgpkitParser::new(url).unwrap().with_filters(&filters);
651 let count2 = parser2.into_elem_iter().count();
652
653 assert_eq!(count1, 132);
655 assert_eq!(count2, 132);
656 }
657
658 #[test]
659 fn test_from_text_reader_inline() {
660 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
661Default local pref 100, local AS 65001\n\n\
662 Network Next Hop Metric LocPrf Weight Path\n\
663 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
664 let parser =
665 BgpkitParser::from_text_reader(dump.as_bytes()).expect("inline text-dump parse");
666 let elems: Vec<_> = parser.into_elem_iter().collect();
667 assert_eq!(elems.len(), 1);
668 assert_eq!(elems[0].prefix.prefix.to_string(), "1.0.0.0/24");
669 assert_eq!(elems[0].peer_ip.to_string(), "1.2.3.4");
670 assert_eq!(u32::from(elems[0].peer_asn), 65001);
671 assert_eq!(elems[0].origin_asns, Some(vec![Asn::from(13335u32)]));
672 }
673
674 #[test]
675 fn test_from_auto_reader_detects_text() {
676 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
677Default local pref 100, local AS 65001\n\n\
678 Network Next Hop Metric LocPrf Weight Path\n\
679 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
680 let parser = BgpkitParser::from_auto_reader(dump.as_bytes()).expect("auto-detect parse");
681 let elems: Vec<_> = parser.into_elem_iter().collect();
682 assert_eq!(elems.len(), 1);
683 assert_eq!(elems[0].origin_asns, Some(vec![Asn::from(13335u32)]));
684 }
685
686 #[test]
687 fn test_from_auto_reader_detects_mrt() {
688 let data: Vec<u8> = vec![0x00u8; 16];
691 let parser =
692 BgpkitParser::from_auto_reader(std::io::Cursor::new(data)).expect("auto-detect parse");
693 let count = parser.into_elem_iter().count();
694 assert_eq!(count, 0);
695 }
696
697 #[test]
698 fn test_text_dump_parser_with_filter() {
699 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
700Default local pref 100, local AS 65001\n\n\
701 Network Next Hop Metric LocPrf Weight Path\n\
702 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n\
703 *> 8.8.8.0/24 10.0.0.2 0 0 15169 i\n";
704 let parser = BgpkitParser::from_text_reader(dump.as_bytes())
705 .unwrap()
706 .add_filter("origin_asn", "13335")
707 .unwrap();
708 let elems: Vec<_> = parser.into_elem_iter().collect();
709 assert_eq!(elems.len(), 1);
710 assert_eq!(elems[0].prefix.prefix.to_string(), "1.0.0.0/24");
711 }
712
713 #[test]
714 fn test_text_dump_next_record_errors() {
715 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
716Default local pref 100, local AS 65001\n\n\
717 Network Next Hop Metric LocPrf Weight Path\n\
718 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
719 let mut parser =
720 BgpkitParser::from_text_reader(dump.as_bytes()).expect("inline text-dump parse");
721 assert!(parser.next_record().is_err());
722 }
723
724 #[test]
725 fn test_text_dump_record_iter_terminates() {
726 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
729Default local pref 100, local AS 65001\n\n\
730 Network Next Hop Metric LocPrf Weight Path\n\
731 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
732 let parser =
733 BgpkitParser::from_text_reader(dump.as_bytes()).expect("inline text-dump parse");
734 assert_eq!(parser.into_record_iter().count(), 0);
735 }
736}