diff --git a/CHANGELOG.md b/CHANGELOG.md index 835053c..4353822 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,13 @@ All notable changes to this project will be documented in this file. ### New Features +* Added repeatable `--filter key=value|key!=value` expressions to `monocle parse` + and `monocle search`. Generic filters use bgpkit-parser's validation and matching + semantics, including IP-family filtering and parser-native regular expressions; + timestamp keys are rejected in favor of Monocle's `--start-ts`, `--end-ts`, and + `--duration` options. Generic filters are retained when `ParseFilters` is + serialized, and local and remote search forward all extended element filters + and generic filters through the SSE request schema. * `monocle parse` now supports route-views `sh ip bgp` snapshots (e.g. `oix-full-snapshot-*.bz2`). These dumps omit the Cisco `BGP table version` / `local AS` preamble, so `peer_ip` and `peer_asn` diff --git a/src/bin/commands/search.rs b/src/bin/commands/search.rs index 8de264e..93b018e 100644 --- a/src/bin/commands/search.rs +++ b/src/bin/commands/search.rs @@ -1482,10 +1482,31 @@ mod tests { let result = url_to_cache_path(&cache_dir, collector, url); assert_eq!(result, None); } + + #[test] + fn test_remote_search_time_bounds_expand_duration() { + let filters = SearchFilters { + parse_filters: monocle::lens::parse::ParseFilters { + start_ts: Some("2026-01-01T00:00:00Z".to_string()), + duration: Some("1h".to_string()), + ..Default::default() + }, + ..Default::default() + }; + + let (start_ts, end_ts) = remote_search_time_bounds(&filters).expect("valid time bounds"); + assert_eq!(start_ts, "1767225600"); + assert_eq!(end_ts, "1767229200"); + } } /// Wrapper to convert local SearchFilters to wire RemoteSearchFilters and run /// the async remote search client on a tokio runtime. +fn remote_search_time_bounds(filters: &SearchFilters) -> anyhow::Result<(String, String)> { + let (start_ts, end_ts) = filters.parse_filters.parse_start_end_strings()?; + Ok((start_ts.to_string(), end_ts.to_string())) +} + fn run_remote_search_wrapper( url: &str, auth_token: Option<&str>, @@ -1494,6 +1515,14 @@ fn run_remote_search_wrapper( output_format: OutputFormat, time_format: TimestampFormat, ) { + let (start_ts, end_ts) = match remote_search_time_bounds(filters) { + Ok(bounds) => bounds, + Err(error) => { + eprintln!("ERROR: failed to resolve remote search time range: {error}"); + std::process::exit(1); + } + }; + // Convert internal SearchFilters to wire RemoteSearchFilters let wire = RemoteSearchFilters { prefix: filters.parse_filters.prefix.clone(), @@ -1514,8 +1543,17 @@ fn run_remote_search_wrapper( }), as_path: filters.parse_filters.as_path.clone(), only_to_customer: filters.parse_filters.only_to_customer.clone(), - start_ts: filters.parse_filters.start_ts.clone().unwrap_or_default(), - end_ts: filters.parse_filters.end_ts.clone().unwrap_or_default(), + next_hop: filters.parse_filters.next_hop.clone(), + origin: filters.parse_filters.origin.clone(), + local_pref: filters.parse_filters.local_pref.clone(), + med: filters.parse_filters.med.clone(), + atomic_aggregate: filters.parse_filters.atomic_aggregate, + aggr_asn: filters.parse_filters.aggr_asn.clone(), + aggr_ip: filters.parse_filters.aggr_ip.clone(), + peer_bgp_id: filters.parse_filters.peer_bgp_id.clone(), + generic_filters: filters.parse_filters.generic_filters.clone(), + start_ts, + end_ts, collector: filters.collector.clone(), project: filters.project.clone(), dump_type: Some( diff --git a/src/bin/commands/search_remote.rs b/src/bin/commands/search_remote.rs index 16d2859..31b011b 100644 --- a/src/bin/commands/search_remote.rs +++ b/src/bin/commands/search_remote.rs @@ -48,6 +48,24 @@ pub struct RemoteSearchFilters { pub as_path: Option, #[serde(skip_serializing_if = "Option::is_none")] pub only_to_customer: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub next_hop: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub origin: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub local_pref: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub med: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub atomic_aggregate: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub aggr_asn: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub aggr_ip: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub peer_bgp_id: Option, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub generic_filters: Vec, pub start_ts: String, pub end_ts: String, #[serde(skip_serializing_if = "Option::is_none")] @@ -249,16 +267,35 @@ mod tests { use super::*; #[test] - fn test_remote_filters_serialize_only_to_customer() { + fn test_remote_filters_serialize_extended_and_generic_filters() { let filters = RemoteSearchFilters { only_to_customer: Some("6777".to_string()), + next_hop: Some("2001:db8::1".to_string()), + origin: Some("igp".to_string()), + local_pref: Some("100".to_string()), + med: Some("50".to_string()), + atomic_aggregate: Some(true), + aggr_asn: Some("64496".to_string()), + aggr_ip: Some("192.0.2.1".to_string()), + peer_bgp_id: Some("192.0.2.2".to_string()), + generic_filters: vec!["ip_version=ipv6".to_string()], start_ts: "1".to_string(), end_ts: "2".to_string(), ..Default::default() }; let json = serde_json::to_value(&filters).unwrap(); assert_eq!(json["only_to_customer"], "6777"); - // Unset optional dimensions are omitted from the wire payload - assert!(json.get("next_hop").is_none()); + assert_eq!(json["next_hop"], "2001:db8::1"); + assert_eq!(json["origin"], "igp"); + assert_eq!(json["local_pref"], "100"); + assert_eq!(json["med"], "50"); + assert_eq!(json["atomic_aggregate"], true); + assert_eq!(json["aggr_asn"], "64496"); + assert_eq!(json["aggr_ip"], "192.0.2.1"); + assert_eq!(json["peer_bgp_id"], "192.0.2.2"); + assert_eq!( + json["generic_filters"], + serde_json::json!(["ip_version=ipv6"]) + ); } } diff --git a/src/lens/parse/mod.rs b/src/lens/parse/mod.rs index 36335b3..b8df4bb 100644 --- a/src/lens/parse/mod.rs +++ b/src/lens/parse/mod.rs @@ -279,6 +279,16 @@ pub struct ParseFilters { #[cfg_attr(feature = "cli", clap(short = 'a', long))] pub as_path: Option, + /// Apply a bgpkit-parser filter expression (`key=value` or `key!=value`). + /// May be specified multiple times. Time filter keys are rejected; use + /// `--start-ts`, `--end-ts`, and `--duration` instead. + #[cfg_attr( + feature = "cli", + clap(long = "filter", value_name = "KEY=VALUE", action = clap::ArgAction::Append) + )] + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub generic_filters: Vec, + // --- bgpkit-parser extended element filters --- // Each accepts the literal value, or `*` (present) / `!*` (absent) for // optional fields (only_to_customer, next_hop, origin, local_pref, med, @@ -335,6 +345,17 @@ pub struct ParseFilters { type FilterSpec = (&'static str, String); +const TIME_FILTER_KEYS: [&str; 4] = ["start_ts", "end_ts", "ts_start", "ts_end"]; +const MULTI_VALUE_FILTER_KEYS: [&str; 7] = [ + "origin_asns", + "prefixes", + "prefixes_super", + "prefixes_sub", + "prefixes_super_sub", + "peer_ips", + "peer_asns", +]; + impl ParseFilters { /// Parse start and end time strings into Unix timestamps pub fn parse_start_end_strings(&self) -> Result<(i64, i64)> { @@ -454,6 +475,7 @@ impl ParseFilters { // --- v0.19 extended element filter validation --- self.validate_extended_filters()?; + self.generic_filter_specs()?; Ok(()) } @@ -748,12 +770,76 @@ impl ParseFilters { Ok(specs) } + fn generic_filter_specs(&self) -> Result> { + self.generic_filters + .iter() + .map(|expression| { + let (filter_type, filter_value) = Self::parse_generic_filter_expression(expression)?; + if TIME_FILTER_KEYS.contains(&filter_type.as_str()) { + return Err(anyhow!( + "Invalid --filter '{}': time filter '{}' is not supported; use --start-ts, --end-ts, or --duration instead", + expression, + filter_type + )); + } + Filter::new(&filter_type, &filter_value).map_err(|error| { + anyhow!("Invalid --filter '{}': {}", expression, error) + })?; + Ok((filter_type, filter_value)) + }) + .collect() + } + + fn parse_generic_filter_expression(expression: &str) -> Result<(String, String)> { + let (filter_type, filter_value, negated) = + if let Some((filter_type, filter_value)) = expression.split_once("!=") { + (filter_type.trim(), filter_value.trim(), true) + } else if let Some((filter_type, filter_value)) = expression.split_once('=') { + (filter_type.trim(), filter_value.trim(), false) + } else { + return Err(anyhow!( + "Invalid --filter '{}': expression must contain '=' or '!='", + expression + )); + }; + + if filter_type.is_empty() { + return Err(anyhow!( + "Invalid --filter '{}': filter key cannot be empty", + expression + )); + } + if filter_value.is_empty() { + return Err(anyhow!( + "Invalid --filter '{}': filter value cannot be empty", + expression + )); + } + + let filter_value = if negated && MULTI_VALUE_FILTER_KEYS.contains(&filter_type) { + filter_value + .split(',') + .map(|value| format!("!{}", value.trim())) + .collect::>() + .join(",") + } else if negated { + format!("!{filter_value}") + } else { + filter_value.to_string() + }; + + Ok((filter_type.to_string(), filter_value)) + } + /// Convert filters into BgpElem predicates using bgpkit-parser's canonical semantics. pub fn to_filters(&self) -> Result> { let mut filters = Vec::new(); for (filter_type, filter_value) in self.filter_specs()? { filters.push(Filter::new(filter_type, &filter_value)?); } + for (filter_type, filter_value) in self.generic_filter_specs()? { + filters.push(Filter::new(&filter_type, &filter_value)?); + } Ok(filters) } @@ -767,6 +853,9 @@ impl ParseFilters { for (filter_type, filter_value) in self.filter_specs()? { parser = parser.add_filter(filter_type, &filter_value)?; } + for (filter_type, filter_value) in self.generic_filter_specs()? { + parser = parser.add_filter(&filter_type, &filter_value)?; + } Ok(parser) } } @@ -1429,4 +1518,49 @@ mod tests { } assert_eq!(actual.len(), 9); } + + #[test] + fn test_generic_filters_use_parser_semantics() { + let filters = ParseFilters { + generic_filters: vec![ + "ip_version=ipv6".to_string(), + "origin_asns!=13335,15169".to_string(), + ], + ..Default::default() + }; + + let actual = filters.to_filters().expect("filter conversion failed"); + assert!(actual.contains(&Filter::new("ip_version", "ipv6").expect("valid IP filter"))); + assert!(actual.contains( + &Filter::new("origin_asns", "!13335,!15169").expect("valid negated origin filter") + )); + } + + #[test] + fn test_generic_filters_reject_time_filter_keys() { + let filters = ParseFilters { + generic_filters: vec!["start_ts=2026-01-01T00:00:00Z".to_string()], + ..Default::default() + }; + + assert!(filters.validate().is_err()); + } + + #[test] + fn test_generic_filters_serialize_when_present() { + let filters = ParseFilters { + generic_filters: vec!["ip_version=ipv6".to_string()], + ..Default::default() + }; + + let serialized = serde_json::to_value(&filters).expect("filters should serialize"); + assert_eq!( + serialized["generic_filters"], + serde_json::json!(["ip_version=ipv6"]) + ); + assert!(serde_json::to_value(ParseFilters::default()) + .expect("empty filters should serialize") + .get("generic_filters") + .is_none()); + } } diff --git a/src/lens/search/mod.rs b/src/lens/search/mod.rs index eb36626..73fc0ec 100644 --- a/src/lens/search/mod.rs +++ b/src/lens/search/mod.rs @@ -306,8 +306,9 @@ impl SearchFilters { Ok(broker) } - /// Validate the filters + /// Validate source-selection and parser filters before searching. pub fn validate(&self) -> Result<()> { + self.parse_filters.validate()?; let _ = self.parse_filters.parse_start_end_strings()?; Ok(()) } diff --git a/src/server/search.rs b/src/server/search.rs index 0bfe891..1f32944 100644 --- a/src/server/search.rs +++ b/src/server/search.rs @@ -91,6 +91,10 @@ pub struct SearchStreamFilters { pub aggr_ip: Option, #[serde(default)] pub peer_bgp_id: Option, + /// Additional bgpkit-parser filter expressions (`key=value` or `key!=value`). + /// Time filter keys are rejected; use the required start/end timestamps. + #[serde(default)] + pub generic_filters: Vec, /// Start timestamp (unix or human-readable). Required. pub start_ts: String, /// End timestamp (unix or human-readable). Required. @@ -150,6 +154,7 @@ impl TryFrom for SearchFilters { end_ts: Some(f.end_ts), duration: None, as_path: f.as_path, + generic_filters: f.generic_filters, only_to_customer: f.only_to_customer, next_hop: f.next_hop, origin: f.origin, @@ -654,6 +659,8 @@ mod tests { fn test_search_stream_filters_conversion() { let wire = SearchStreamFilters { prefix: vec!["1.1.1.0/24".to_string()], + next_hop: Some("192.0.2.1".to_string()), + generic_filters: vec!["ip_version=ipv4".to_string()], start_ts: "2024-01-01T00:00:00Z".to_string(), end_ts: "2024-01-01T00:10:00Z".to_string(), collector: Some("rrc00".to_string()), @@ -664,6 +671,11 @@ mod tests { let filters: SearchFilters = wire.try_into().expect("conversion should succeed"); assert_eq!(filters.collector, Some("rrc00".to_string())); assert_eq!(filters.parse_filters.prefix, vec!["1.1.1.0/24"]); + assert_eq!(filters.parse_filters.next_hop.as_deref(), Some("192.0.2.1")); + assert_eq!( + filters.parse_filters.generic_filters, + vec!["ip_version=ipv4"] + ); } #[test] @@ -679,6 +691,19 @@ mod tests { assert!(result.is_err()); } + #[test] + fn test_search_stream_rejects_generic_time_filter() { + let wire = SearchStreamFilters { + generic_filters: vec!["start_ts=2026-01-01T00:00:00Z".to_string()], + start_ts: "2026-01-01T00:00:00Z".to_string(), + end_ts: "2026-01-01T00:10:00Z".to_string(), + ..Default::default() + }; + + let filters: SearchFilters = wire.try_into().expect("conversion should succeed"); + assert!(filters.validate().is_err()); + } + #[test] fn test_search_stream_event_to_sse() { let event = SearchStreamEvent::Started(SearchStarted {