Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down
42 changes: 40 additions & 2 deletions src/bin/commands/search.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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>,
Expand All @@ -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(),
Expand All @@ -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(
Expand Down
43 changes: 40 additions & 3 deletions src/bin/commands/search_remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,24 @@ pub struct RemoteSearchFilters {
pub as_path: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub only_to_customer: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub next_hop: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub origin: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub local_pref: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub med: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub atomic_aggregate: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub aggr_asn: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub aggr_ip: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub peer_bgp_id: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub generic_filters: Vec<String>,
pub start_ts: String,
pub end_ts: String,
#[serde(skip_serializing_if = "Option::is_none")]
Expand Down Expand Up @@ -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"])
);
}
}
134 changes: 134 additions & 0 deletions src/lens/parse/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -279,6 +279,16 @@ pub struct ParseFilters {
#[cfg_attr(feature = "cli", clap(short = 'a', long))]
pub as_path: Option<String>,

/// 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<String>,

// --- 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,
Expand Down Expand Up @@ -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)> {
Expand Down Expand Up @@ -454,6 +475,7 @@ impl ParseFilters {

// --- v0.19 extended element filter validation ---
self.validate_extended_filters()?;
self.generic_filter_specs()?;

Ok(())
}
Expand Down Expand Up @@ -748,12 +770,76 @@ impl ParseFilters {
Ok(specs)
}

fn generic_filter_specs(&self) -> Result<Vec<(String, String)>> {
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::<Vec<_>>()
.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<Vec<Filter>> {
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)
}

Expand All @@ -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)
}
}
Expand Down Expand Up @@ -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());
}
}
3 changes: 2 additions & 1 deletion src/lens/search/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}
Expand Down
25 changes: 25 additions & 0 deletions src/server/search.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,10 @@ pub struct SearchStreamFilters {
pub aggr_ip: Option<String>,
#[serde(default)]
pub peer_bgp_id: Option<String>,
/// 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<String>,
/// Start timestamp (unix or human-readable). Required.
pub start_ts: String,
/// End timestamp (unix or human-readable). Required.
Expand Down Expand Up @@ -150,6 +154,7 @@ impl TryFrom<SearchStreamFilters> 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,
Expand Down Expand Up @@ -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()),
Expand All @@ -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]
Expand All @@ -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 {
Expand Down
Loading