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, RisSubscribe, RisSubscribeSocketOptions, RisSubscribeType,
37};
38#[cfg(feature = "rislive")]
39pub use rislive::{
40 parse_ris_live_message, parse_ris_live_message_json, parse_ris_live_message_raw,
41};
42
43pub struct BgpkitParser<R> {
44 reader: R,
45 core_dump: bool,
46 filters: Vec<Filter>,
47 options: ParserOptions,
48 text_dump_iter: Option<Box<dyn Iterator<Item = BgpElem> + Send>>,
51}
52
53pub(crate) struct ParserOptions {
54 show_warnings: bool,
55 warned_zebra_compat: bool,
56}
57impl Default for ParserOptions {
58 fn default() -> Self {
59 ParserOptions {
60 show_warnings: true,
61 warned_zebra_compat: false,
62 }
63 }
64}
65
66impl ParserOptions {
67 pub(crate) fn warn_zebra_compat_once(&mut self) {
68 if self.show_warnings && !self.warned_zebra_compat {
69 warn!(
70 "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)"
71 );
72 self.warned_zebra_compat = true;
73 }
74 }
75}
76
77#[cfg(feature = "oneio")]
78impl BgpkitParser<Box<dyn Read + Send>> {
79 pub fn new(path: &str) -> Result<Self, ParserErrorWithBytes> {
81 let reader = get_reader(path)?;
82 Ok(BgpkitParser {
83 reader,
84 core_dump: false,
85 filters: vec![],
86 options: ParserOptions::default(),
87 text_dump_iter: None,
88 })
89 }
90
91 pub fn new_resumable_http(path: &str) -> Result<Self, ParserErrorWithBytes> {
113 let reader = get_resumable_http_reader(path)?;
114 Ok(BgpkitParser {
115 reader,
116 core_dump: false,
117 filters: vec![],
118 options: ParserOptions::default(),
119 text_dump_iter: None,
120 })
121 }
122
123 pub fn new_cached(path: &str, cache_dir: &str) -> Result<Self, ParserErrorWithBytes> {
129 let file_name = path.rsplit('/').next().unwrap().to_string();
130 let new_file_name = format!(
131 "cache-{}",
132 add_suffix_to_filename(file_name.as_str(), crc32(path).as_str())
133 );
134 let reader = get_cache_reader(path, cache_dir, Some(new_file_name), false)?;
135 Ok(BgpkitParser {
136 reader,
137 core_dump: false,
138 filters: vec![],
139 options: ParserOptions::default(),
140 text_dump_iter: None,
141 })
142 }
143
144 pub fn new_text(path: &str) -> Result<Self, ParserErrorWithBytes> {
167 let timestamp = infer_timestamp_from_path(path).unwrap_or(0.0);
168 let reader = get_reader(path)?;
169 Self::from_text_reader_with_timestamp(reader, timestamp)
170 }
171
172 pub fn new_auto(path: &str) -> Result<Self, ParserErrorWithBytes> {
191 let reader = get_reader(path)?;
192 Self::from_auto_reader_with_timestamp(reader, infer_timestamp_from_path(path))
193 }
194}
195
196#[cfg(feature = "oneio")]
197fn add_suffix_to_filename(filename: &str, suffix: &str) -> String {
198 let mut parts: Vec<&str> = filename.split('.').collect(); if parts.len() > 1 {
200 let last_part = parts.pop().unwrap(); let new_last_part = format!("{suffix}.{last_part}"); parts.push(&new_last_part); parts.join(".") } else {
205 format!("{filename}.{suffix}")
207 }
208}
209
210impl<R: Read> BgpkitParser<R> {
211 pub fn from_reader(reader: R) -> Self {
213 BgpkitParser {
214 reader,
215 core_dump: false,
216 filters: vec![],
217 options: ParserOptions::default(),
218 text_dump_iter: None,
219 }
220 }
221
222 pub fn next_record(&mut self) -> Result<MrtRecord, ParserErrorWithBytes> {
224 if self.text_dump_iter.is_some() {
225 return Err(ParserError::Unsupported(
226 "text-dump parsers have no MRT record representation; iterate elements instead"
227 .to_string(),
228 )
229 .into());
230 }
231 let (record, used_zebra_compat) =
232 mrt::mrt_record::parse_mrt_record_with_zebra_compat(&mut self.reader)?;
233 if used_zebra_compat {
234 self.warn_zebra_compat_once();
235 }
236 Ok(record)
237 }
238}
239
240impl BgpkitParser<Box<dyn Read + Send>> {
241 pub fn from_text_reader(
245 reader: impl Read + Send + 'static,
246 ) -> Result<Self, ParserErrorWithBytes> {
247 Self::from_text_reader_with_timestamp(reader, 0.0)
248 }
249
250 pub fn from_text_reader_with_timestamp(
254 reader: impl Read + Send + 'static,
255 timestamp: f64,
256 ) -> Result<Self, ParserErrorWithBytes> {
257 let buf_reader = BufReader::new(reader);
258 let iter = TextDumpElemIterator::new(buf_reader, timestamp).map_err(ParserError::from)?;
259 Ok(BgpkitParser {
260 reader: Box::new(std::io::empty()),
261 core_dump: false,
262 filters: vec![],
263 options: ParserOptions::default(),
264 text_dump_iter: Some(Box::new(iter)),
265 })
266 }
267
268 pub fn from_auto_reader(
271 reader: impl Read + Send + 'static,
272 ) -> Result<Self, ParserErrorWithBytes> {
273 Self::from_auto_reader_with_timestamp(reader, None)
274 }
275
276 pub fn from_auto_reader_with_timestamp(
280 reader: impl Read + Send + 'static,
281 timestamp: Option<f64>,
282 ) -> Result<Self, ParserErrorWithBytes> {
283 let mut buf_reader = BufReader::new(reader);
284 let (is_text, head) = detect_text_dump(&mut buf_reader).map_err(ParserError::from)?;
285 if is_text {
286 let chained = BufReader::new(Cursor::new(head).chain(buf_reader));
287 let ts = timestamp.unwrap_or(0.0);
288 let iter = TextDumpElemIterator::new(chained, ts).map_err(ParserError::from)?;
289 Ok(BgpkitParser {
290 reader: Box::new(std::io::empty()),
291 core_dump: false,
292 filters: vec![],
293 options: ParserOptions::default(),
294 text_dump_iter: Some(Box::new(iter)),
295 })
296 } else {
297 Ok(BgpkitParser {
298 reader: Box::new(Cursor::new(head).chain(buf_reader)),
299 core_dump: false,
300 filters: vec![],
301 options: ParserOptions::default(),
302 text_dump_iter: None,
303 })
304 }
305 }
306}
307
308impl<R> BgpkitParser<R> {
309 pub(crate) fn warn_zebra_compat_once(&mut self) {
310 self.options.warn_zebra_compat_once();
311 }
312
313 pub fn enable_core_dump(self) -> Self {
314 BgpkitParser {
315 reader: self.reader,
316 core_dump: true,
317 filters: self.filters,
318 options: self.options,
319 text_dump_iter: self.text_dump_iter,
320 }
321 }
322
323 pub fn disable_warnings(self) -> Self {
324 let mut options = self.options;
325 options.show_warnings = false;
326 BgpkitParser {
327 reader: self.reader,
328 core_dump: self.core_dump,
329 filters: self.filters,
330 options,
331 text_dump_iter: self.text_dump_iter,
332 }
333 }
334
335 pub fn add_filter(
399 self,
400 filter_type: &str,
401 filter_value: &str,
402 ) -> Result<Self, ParserErrorWithBytes> {
403 let mut filters = self.filters;
404 filters.push(Filter::new(filter_type, filter_value)?);
405 Ok(BgpkitParser {
406 reader: self.reader,
407 core_dump: self.core_dump,
408 filters,
409 options: self.options,
410 text_dump_iter: self.text_dump_iter,
411 })
412 }
413
414 pub fn add_filters(mut self, filters: &[Filter]) -> Self {
434 self.filters.extend(filters.iter().cloned());
435 self
436 }
437
438 pub fn with_filters(mut self, filters: &[Filter]) -> Self {
466 self.filters = filters.to_vec();
467 self
468 }
469}
470
471#[cfg(test)]
472mod tests {
473 use super::*;
474 use crate::models::Asn;
475
476 #[test]
477 fn test_new_with_reader() {
478 let reader = oneio::get_reader("http://archive.routeviews.org/route-views.ny/bgpdata/2023.02/UPDATES/updates.20230215.0630.bz2").unwrap();
480 assert_eq!(
481 12683,
482 BgpkitParser::from_reader(reader).into_elem_iter().count()
483 );
484
485 let reader = oneio::get_reader("https://spaces.bgpkit.org/parser/update-example").unwrap();
487 assert_eq!(
488 8160,
489 BgpkitParser::from_reader(reader).into_elem_iter().count()
490 );
491 }
492
493 #[test]
494 fn test_new_resumable_http() {
495 let parser =
496 BgpkitParser::new_resumable_http("https://spaces.bgpkit.org/parser/update-example.gz")
497 .unwrap();
498 assert_eq!(8160, parser.into_elem_iter().count());
499 }
500
501 #[test]
502 fn test_new_cached_with_reader() {
503 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
504 let parser = BgpkitParser::new_cached(url, "/tmp/bgpkit-parser-tests")
505 .unwrap()
506 .enable_core_dump()
507 .disable_warnings();
508 let count = parser.into_elem_iter().count();
509 assert_eq!(8160, count);
510 let parser = BgpkitParser::new_cached(url, "/tmp/bgpkit-parser-tests").unwrap();
511 let count = parser.into_elem_iter().count();
512 assert_eq!(8160, count);
513 }
514
515 #[test]
516 fn test_add_suffix_to_filename() {
517 let filename = "example.txt";
519 let suffix = "suffix";
520 let result = add_suffix_to_filename(filename, suffix);
521 assert_eq!(result, "example.suffix.txt");
522
523 let filename = "example.tar.gz";
525 let suffix = "suffix";
526 let result = add_suffix_to_filename(filename, suffix);
527 assert_eq!(result, "example.tar.suffix.gz");
528
529 let filename = "example";
531 let suffix = "suffix";
532 let result = add_suffix_to_filename(filename, suffix);
533 assert_eq!(result, "example.suffix");
534
535 let filename = "";
537 let suffix = "suffix";
538 let result = add_suffix_to_filename(filename, suffix);
539 assert_eq!(result, ".suffix");
540
541 let filename = "example.txt";
543 let suffix = "";
544 let result = add_suffix_to_filename(filename, suffix);
545 assert_eq!(result, "example..txt");
546 }
547
548 #[test]
549 fn test_with_filters() {
550 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
551
552 let filters = vec![
554 Filter::new("peer_ip", "185.1.8.65").unwrap(),
555 Filter::new("type", "w").unwrap(),
556 ];
557
558 let parser = BgpkitParser::new(url).unwrap().with_filters(&filters);
560 let count = parser.into_elem_iter().count();
561
562 assert_eq!(count, 132);
564
565 let filters1 = vec![Filter::new("peer_ip", "185.1.8.65").unwrap()];
567 let filters2 = vec![Filter::new("peer_ip", "185.1.8.50").unwrap()];
568
569 let parser = BgpkitParser::new(url)
570 .unwrap()
571 .with_filters(&filters1)
572 .with_filters(&filters2); let count = parser.into_elem_iter().count();
574
575 assert_eq!(count, 1563);
577 }
578
579 #[test]
580 fn test_add_filters() {
581 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
582
583 let filters = vec![
585 Filter::new("peer_ip", "185.1.8.65").unwrap(),
586 Filter::new("type", "w").unwrap(),
587 ];
588
589 let parser = BgpkitParser::new(url).unwrap().add_filters(&filters);
591 let count = parser.into_elem_iter().count();
592
593 assert_eq!(count, 132);
595
596 let parser = BgpkitParser::new(url)
598 .unwrap()
599 .add_filter("peer_ip", "185.1.8.65")
600 .unwrap()
601 .add_filters(&[Filter::new("type", "w").unwrap()]);
602 let count = parser.into_elem_iter().count();
603 assert_eq!(count, 132);
604 }
605
606 #[test]
607 fn test_with_filters_empty() {
608 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
609
610 let parser = BgpkitParser::new(url).unwrap().with_filters(&[]);
612 let count = parser.into_elem_iter().count();
613
614 assert_eq!(count, 8160);
616 }
617
618 #[test]
619 fn test_add_filters_empty() {
620 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
621
622 let parser = BgpkitParser::new(url)
624 .unwrap()
625 .add_filter("peer_ip", "185.1.8.65")
626 .unwrap()
627 .add_filters(&[]);
628 let count = parser.into_elem_iter().count();
629
630 assert_eq!(count, 3393);
632 }
633
634 #[test]
635 fn test_with_filters_reuse() {
636 let url = "https://spaces.bgpkit.org/parser/update-example.gz";
637
638 let filters = vec![
640 Filter::new("peer_ip", "185.1.8.65").unwrap(),
641 Filter::new("type", "w").unwrap(),
642 ];
643
644 let parser1 = BgpkitParser::new(url).unwrap().with_filters(&filters);
646 let count1 = parser1.into_elem_iter().count();
647
648 let parser2 = BgpkitParser::new(url).unwrap().with_filters(&filters);
649 let count2 = parser2.into_elem_iter().count();
650
651 assert_eq!(count1, 132);
653 assert_eq!(count2, 132);
654 }
655
656 #[test]
657 fn test_from_text_reader_inline() {
658 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
659Default local pref 100, local AS 65001\n\n\
660 Network Next Hop Metric LocPrf Weight Path\n\
661 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
662 let parser =
663 BgpkitParser::from_text_reader(dump.as_bytes()).expect("inline text-dump parse");
664 let elems: Vec<_> = parser.into_elem_iter().collect();
665 assert_eq!(elems.len(), 1);
666 assert_eq!(elems[0].prefix.prefix.to_string(), "1.0.0.0/24");
667 assert_eq!(elems[0].peer_ip.to_string(), "1.2.3.4");
668 assert_eq!(u32::from(elems[0].peer_asn), 65001);
669 assert_eq!(elems[0].origin_asns, Some(vec![Asn::from(13335u32)]));
670 }
671
672 #[test]
673 fn test_from_auto_reader_detects_text() {
674 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
675Default local pref 100, local AS 65001\n\n\
676 Network Next Hop Metric LocPrf Weight Path\n\
677 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
678 let parser = BgpkitParser::from_auto_reader(dump.as_bytes()).expect("auto-detect parse");
679 let elems: Vec<_> = parser.into_elem_iter().collect();
680 assert_eq!(elems.len(), 1);
681 assert_eq!(elems[0].origin_asns, Some(vec![Asn::from(13335u32)]));
682 }
683
684 #[test]
685 fn test_from_auto_reader_detects_mrt() {
686 let data: Vec<u8> = vec![0x00u8; 16];
689 let parser =
690 BgpkitParser::from_auto_reader(std::io::Cursor::new(data)).expect("auto-detect parse");
691 let count = parser.into_elem_iter().count();
692 assert_eq!(count, 0);
693 }
694
695 #[test]
696 fn test_text_dump_parser_with_filter() {
697 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
698Default local pref 100, local AS 65001\n\n\
699 Network Next Hop Metric LocPrf Weight Path\n\
700 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n\
701 *> 8.8.8.0/24 10.0.0.2 0 0 15169 i\n";
702 let parser = BgpkitParser::from_text_reader(dump.as_bytes())
703 .unwrap()
704 .add_filter("origin_asn", "13335")
705 .unwrap();
706 let elems: Vec<_> = parser.into_elem_iter().collect();
707 assert_eq!(elems.len(), 1);
708 assert_eq!(elems[0].prefix.prefix.to_string(), "1.0.0.0/24");
709 }
710
711 #[test]
712 fn test_text_dump_next_record_errors() {
713 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
714Default local pref 100, local AS 65001\n\n\
715 Network Next Hop Metric LocPrf Weight Path\n\
716 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
717 let mut parser =
718 BgpkitParser::from_text_reader(dump.as_bytes()).expect("inline text-dump parse");
719 assert!(parser.next_record().is_err());
720 }
721
722 #[test]
723 fn test_text_dump_record_iter_terminates() {
724 let dump = "BGP table version is 1, local router ID is 1.2.3.4, vrf id 0\n\
727Default local pref 100, local AS 65001\n\n\
728 Network Next Hop Metric LocPrf Weight Path\n\
729 *> 1.0.0.0/24 10.0.0.1 0 0 13335 i\n";
730 let parser =
731 BgpkitParser::from_text_reader(dump.as_bytes()).expect("inline text-dump parse");
732 assert_eq!(parser.into_record_iter().count(), 0);
733 }
734}