From 11c2842d79f5558e932639e0b685ac3368980ded Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 25 Sep 2026 12:50:47 -0500 Subject: [PATCH 01/10] Add tenant-aware EVPN monitoring for symmetric IRB --- README.md | 1 + docs/evpn.md | 93 +++++++ docs/index.md | 1 + docs/rib-query-api.md | 3 + src/payload.rs | 15 +- src/roto_runtime/runtime.rs | 15 ++ src/roto_runtime/types.rs | 29 ++- src/units/bmp_tcp_out/bmp_builder.rs | 112 ++++++++ src/units/bmp_tcp_out/client_handler.rs | 201 ++++++++++++++ src/units/mrt_file_in/unit.rs | 2 + src/units/rib_unit/evpn.rs | 333 ++++++++++++++++++++++++ src/units/rib_unit/http_ng.rs | 153 ++++++++++- src/units/rib_unit/mod.rs | 1 + src/units/rib_unit/rib.rs | 144 +++++++++- 14 files changed, 1097 insertions(+), 6 deletions(-) create mode 100644 docs/evpn.md create mode 100644 src/units/rib_unit/evpn.rs diff --git a/README.md b/README.md index f8b83dd..ddaca5d 100644 --- a/README.md +++ b/README.md @@ -16,6 +16,7 @@ Rotonda version from which it was forked, Netom adds: - [active TCP/TLS BMP input](docs/bmp-tcp-in.md) for pulling exporter feeds; - BMP restreaming with an initial RIB dump followed by live updates; +- [EVPN monitoring](docs/evpn.md) with tenant-aware route queries for symmetric IRB; - bounded buffers, streaming full-RIB exports, and slow-consumer protection; - stronger BMP peer lifecycle, reconnect, withdrawal, and memory handling; - TLS and access controls for BMP consumers; diff --git a/docs/evpn.md b/docs/evpn.md new file mode 100644 index 0000000..dc2ed43 --- /dev/null +++ b/docs/evpn.md @@ -0,0 +1,93 @@ +# EVPN monitoring + +Netom collects L2VPN EVPN (AFI 25, SAFI 70) from BGP and BMP, including +ADD-PATH sessions. Configure `L2VpnEvpn` in a BGP peer's `protocols` list; +BMP peers use the capabilities in their exported Peer Up messages. + +EVPN routes live in a separate RIB. Overlapping tenant prefixes do not collide +with each other or with the global unicast table. Identity includes the route +type, wire-format route distinguisher (RD), route-specific key, and ingress +(including the ADD-PATH child). Type 2 identity excludes ESI and labels; +type 5 identity excludes ESI, gateway, and label. Changes to forwarding fields +replace the same route, and withdrawals match even when labels differ. +Peer Down, family-scoped withdrawals, and ingress cleanup include EVPN. + +## Symmetric IRB + +Inspect type 2 MAC/IP advertisements and type 5 IP prefix advertisements +together to monitor the MAC-VRF and IP-VRF views of a tenant. The decoder +exposes the RD, Ethernet tag, ESI, MAC, IPv4/IPv6 host or prefix, gateway, +and label fields. Both type 2 labels are preserved for deployments that +advertise a MAC-VRF label and an IP-VRF label. + +The API also decodes route targets, the EVPN Router's MAC extended community, +and the MP_REACH next hop. Use a route target to select the routes associated +with a tenant across multiple advertising RDs. An RD distinguishes routes; +it is not a tenant membership or import-policy identifier. Route targets are +reported as advertised; Netom does not simulate VRF import policy or recursive +forwarding resolution. + +Label fields are returned as **raw 24-bit integers**. For VXLAN these are +VNIs. For MPLS, decode the label-stack entry according to the encapsulation; +do not interpret the raw field as an MPLS label number. The API does not +infer which VNI is an L2 VNI or L3 VNI from its value alone. + +## Query API + +`GET /api/v1/ribs/l2vpnevpn/routes` returns `{"data": [...]}`. Each row has: + +- `route`: `nlri`, retained `attributes`, `ingress_id`, `ltime`, and `active`; +- `overlay`: `route_targets`, `router_mac`, and `next_hop`; +- `source_ingress_id` and `path_id`: original session and optional ADD-PATH ID. + +The `nlri.raw` byte array includes the route type and length header, allowing +inspection of fields not decoded by the API. Types 1, 3, and 4 are retained, +with RD and applicable tag/ESI fields decoded. Unknown route types with an RD +are preserved as opaque records and matched by their complete NLRI. + +All filters below are optional and combine with AND: + +| Parameter | Meaning | +| --- | --- | +| `rd` | Exact RD, e.g. `65000:100` or `192.0.2.1:100` | +| `route_target` | Exact advertised RT, e.g. `65000:100` | +| `route_type` | Numeric EVPN type, e.g. `2` or `5` | +| `vni` | Match either raw label field (for VXLAN monitoring) | +| `prefix` | Exact type 2 host or type 5 prefix | +| `ingress_id` | Exact stored ingress ID, including ADD-PATH child IDs | +| `include_withdrawn` | `true` includes retained withdrawn records; default `false` | + +Unknown parameters are rejected. Withdrawn records are available only when +the RIB is configured to retain withdrawn attributes. A withdrawal preserves +the last announced attributes and forwarding fields. + +```sh +curl -G http://127.0.0.1:8080/api/v1/ribs/l2vpnevpn/routes \ + --data-urlencode 'route_target=65000:100' \ + --data-urlencode 'route_type=5' + +curl -G http://127.0.0.1:8080/api/v1/ribs/l2vpnevpn/routes \ + --data-urlencode 'rd=192.0.2.1:100' \ + --data-urlencode 'prefix=10.0.0.0/24' +``` + +The endpoint produces buffered JSON and takes a snapshot of the EVPN table; +large tables require memory proportional to the stored records and response. +It shares the concurrent query limit with other RIB queries. + +## Pipeline support and limits + +BMP output includes EVPN in live rebuilt updates and initial RIB dumps, +retaining next hops, communities, and ADD-PATH IDs. Synthetic Peer Up messages +and End-of-RIB markers include EVPN when advertised by the source peer. +BGP4MP MRT updates use the same decoder; TABLE_DUMP_V2 EVPN import is not +implemented. The existing CLI route commands and ClickHouse route schema do +not expose EVPN; use the HTTP API or BMP output. + +Roto route filters can use `is_evpn()` and `evpn_rd()`. `fmt_prefix()` returns +the type 2 host or type 5 prefix; EVPN routes without an IP prefix return +`0.0.0.0/0`, so guard IP-only policies with `is_evpn()`. + +Wire formats follow [RFC 7432](https://www.rfc-editor.org/rfc/rfc7432), +[RFC 9135](https://www.rfc-editor.org/rfc/rfc9135), and +[RFC 9136](https://www.rfc-editor.org/rfc/rfc9136). diff --git a/docs/index.md b/docs/index.md index ea2193d..819762b 100644 --- a/docs/index.md +++ b/docs/index.md @@ -38,6 +38,7 @@ clickhouse rib-query-api addpath-flowspec-api +evpn best-path-selection ``` diff --git a/docs/rib-query-api.md b/docs/rib-query-api.md index 1904c30..d815296 100644 --- a/docs/rib-query-api.md +++ b/docs/rib-query-api.md @@ -15,8 +15,11 @@ in this API changes state. ## Endpoints +EVPN uses a separate [tenant-aware endpoint](evpn.md). + | Endpoint | Returns | | --- | --- | +| `GET /api/v1/ribs/l2vpnevpn/routes` | EVPN routes, with RD/RT/VNI filters (see [EVPN](evpn.md)) | | `GET /api/v1/ribs/ipv4unicast/routes/{addr}/{len}` | every route for one prefix | | `GET /api/v1/ribs/ipv6unicast/routes/{addr}/{len}` | | | `GET /api/v1/ribs/ipv4unicast/routes` | the whole table (see [Whole-table dumps](#whole-table-dumps)) | diff --git a/src/payload.rs b/src/payload.rs index 3c989d9..1f71b0f 100644 --- a/src/payload.rs +++ b/src/payload.rs @@ -60,6 +60,7 @@ pub enum RotondaRoute { routecore::bgp::nlri::afisafi::Ipv6FlowSpecNlri, RotondaPaMap, ), + L2VpnEvpn(crate::units::rib_unit::evpn::EvpnNlri, RotondaPaMap), // TODO support all routecore AfiSafiTypes } @@ -70,6 +71,7 @@ impl Serialize for RotondaRoute { { let mut s = serializer.serialize_struct("Route", 2)?; match self { + RotondaRoute::L2VpnEvpn(n, _) => s.serialize_field("evpn", n), RotondaRoute::Ipv4Unicast(n, _) => s.serialize_field("prefix", n), RotondaRoute::Ipv6Unicast(n, _) => s.serialize_field("prefix", n), RotondaRoute::Ipv4Multicast(n, _) => { @@ -98,6 +100,7 @@ impl RotondaRoute { &self, ) -> routecore::bgp::path_attributes::OwnedPathAttributes { match self { + RotondaRoute::L2VpnEvpn(_, p) => p.path_attributes(), RotondaRoute::Ipv4Unicast(_, p) => p.path_attributes(), RotondaRoute::Ipv6Unicast(_, p) => p.path_attributes(), RotondaRoute::Ipv4Multicast(_, p) => p.path_attributes(), @@ -109,6 +112,7 @@ impl RotondaRoute { pub fn rotonda_pamap(&self) -> &RotondaPaMap { match self { + RotondaRoute::L2VpnEvpn(_, p) => p, RotondaRoute::Ipv4Unicast(_, p) => p, RotondaRoute::Ipv6Unicast(_, p) => p, RotondaRoute::Ipv4Multicast(_, p) => p, @@ -120,6 +124,7 @@ impl RotondaRoute { pub fn rotonda_pamap_mut(&mut self) -> &mut RotondaPaMap { match self { + RotondaRoute::L2VpnEvpn(_, ref mut p) => p, RotondaRoute::Ipv4Unicast(_, ref mut p) => p, RotondaRoute::Ipv6Unicast(_, ref mut p) => p, RotondaRoute::Ipv4Multicast(_, ref mut p) => p, @@ -129,15 +134,20 @@ impl RotondaRoute { } } - /// The prefix under which this route is keyed in the RIB. + /// Prefix exposed to IP-oriented filters. EVPN storage uses its own RD-scoped key. /// /// For unicast/multicast this is the NLRI prefix. For FlowSpec it is /// the destination-prefix component when one is usable as a key (always /// for IPv4; for IPv6 only when the pattern offset is 0), otherwise the /// family default route (`0.0.0.0/0` / `::/0`). roto scripts, the HTTP /// API and the store all derive the key through this one helper. + /// For EVPN this is the type-2 host or type-5 prefix; routes without + /// an IP prefix return 0.0.0.0/0. Use is_evpn() before IP-only policy. pub fn index_prefix(&self) -> inetnum::addr::Prefix { match self { + RotondaRoute::L2VpnEvpn(n, _) => { + n.prefix.unwrap_or_else(|| "0.0.0.0/0".parse().unwrap()) + } RotondaRoute::Ipv4Unicast(n, _) => n.prefix(), RotondaRoute::Ipv6Unicast(n, _) => n.prefix(), RotondaRoute::Ipv4Multicast(n, _) => n.prefix(), @@ -168,6 +178,9 @@ impl RotondaRoute { impl fmt::Display for RotondaRoute { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { + RotondaRoute::L2VpnEvpn(n, ..) => { + write!(f, "RR-EVPN {} {}", n.route_type, n.rd) + } RotondaRoute::Ipv4Unicast(p, ..) => { write!(f, "RR-Ipv4Unicast {}", p) } diff --git a/src/roto_runtime/runtime.rs b/src/roto_runtime/runtime.rs index 2eae915..030cc67 100644 --- a/src/roto_runtime/runtime.rs +++ b/src/roto_runtime/runtime.rs @@ -368,6 +368,21 @@ pub fn create_runtime() -> Result { rr.index_prefix().to_string().into() } + /// Whether this route belongs to the EVPN address family. + #[roto_method(rt, MutRotondaRoute, is_evpn)] + fn rr_is_evpn(rr: Val) -> bool { + matches!(&*rr.borrow(), RotondaRoute::L2VpnEvpn(..)) + } + + /// EVPN route distinguisher, or an empty string for other families. + #[roto_method(rt, MutRotondaRoute, evpn_rd)] + fn rr_evpn_rd(rr: Val) -> Arc { + match &*rr.borrow() { + RotondaRoute::L2VpnEvpn(n, _) => n.rd.clone().into(), + _ => "".into(), + } + } + /// Whether this `RotondaRoute` is a FlowSpec rule (SAFI 133) #[roto_method(rt, MutRotondaRoute, is_flowspec)] fn rr_is_flowspec(rr: Val) -> bool { diff --git a/src/roto_runtime/types.rs b/src/roto_runtime/types.rs index fd3eb29..da6b210 100644 --- a/src/roto_runtime/types.rs +++ b/src/roto_runtime/types.rs @@ -746,6 +746,31 @@ pub(crate) fn convert_nlri>( ) } + Nlri::L2VpnEvpn(n) => { + use routecore::bgp::nlri::afisafi::NlriCompose; + let mut raw = Vec::new(); + n.compose(&mut raw).map_err(|_| ())?; + ( + RotondaRoute::L2VpnEvpn( + crate::units::rib_unit::evpn::EvpnNlri::parse(&raw)?, + pamap, + ), + None, + ) + } + Nlri::L2VpnEvpnAddpath(n) => { + use routecore::bgp::nlri::afisafi::NlriCompose; + let mut raw = Vec::new(); + n.compose(&mut raw).map_err(|_| ())?; + ( + RotondaRoute::L2VpnEvpn( + crate::units::rib_unit::evpn::EvpnNlri::parse(&raw[4..])?, + pamap, + ), + Some(n.path_id()), + ) + } + Nlri::Ipv4MplsUnicast(..) | Nlri::Ipv4MplsUnicastAddpath(..) | Nlri::Ipv4MplsVpnUnicast(..) @@ -757,9 +782,7 @@ pub(crate) fn convert_nlri>( | Nlri::Ipv6MplsVpnUnicast(..) | Nlri::Ipv6MplsVpnUnicastAddpath(..) | Nlri::L2VpnVpls(..) - | Nlri::L2VpnVplsAddpath(..) - | Nlri::L2VpnEvpn(..) - | Nlri::L2VpnEvpnAddpath(..) => { + | Nlri::L2VpnVplsAddpath(..) => { note_unsupported_nlri(nlri.nlri_type()); debug!( "NLRI type {:?} not yet supported in RotondaRoute: {}", diff --git a/src/units/bmp_tcp_out/bmp_builder.rs b/src/units/bmp_tcp_out/bmp_builder.rs index d2c2afe..99d26b8 100644 --- a/src/units/bmp_tcp_out/bmp_builder.rs +++ b/src/units/bmp_tcp_out/bmp_builder.rs @@ -267,6 +267,7 @@ impl PeerInfo { AfiSafiType::Ipv6Unicast => (2u16, 1u8), AfiSafiType::Ipv4FlowSpec => (1u16, 133u8), AfiSafiType::Ipv6FlowSpec => (2u16, 133u8), + AfiSafiType::L2VpnEvpn => (25u16, 70u8), _ => return false, }; Capabilities(&self.peer_capabilities).iter().any(|cap| { @@ -285,6 +286,7 @@ impl PeerInfo { AfiSafiType::Ipv6Unicast, AfiSafiType::Ipv4FlowSpec, AfiSafiType::Ipv6FlowSpec, + AfiSafiType::L2VpnEvpn, ] .into_iter() .filter(|afisafi| self.supports_afisafi(*afisafi)) @@ -458,6 +460,7 @@ fn build_bgp_open( AfiSafiType::Ipv6Unicast => (2u16, 1u8), AfiSafiType::Ipv4FlowSpec => (1u16, 133u8), AfiSafiType::Ipv6FlowSpec => (2u16, 133u8), + AfiSafiType::L2VpnEvpn => (25u16, 70u8), _ => continue, }; caps.push(1); @@ -491,6 +494,7 @@ fn build_bgp_open( AfiSafiType::Ipv6Unicast => (2u16, 1u8), AfiSafiType::Ipv4FlowSpec => (1u16, 133u8), AfiSafiType::Ipv6FlowSpec => (2u16, 133u8), + AfiSafiType::L2VpnEvpn => (25u16, 70u8), _ => continue, }; caps.extend_from_slice(&afi.to_be_bytes()); @@ -820,6 +824,15 @@ pub fn build_route_monitoring_from_route( path_id, ); } + RotondaRoute::L2VpnEvpn(nlri, pamap) => { + return build_evpn_route_monitoring( + peer, + &nlri.raw, + pamap, + is_withdrawal, + path_id, + ); + } RotondaRoute::Ipv6FlowSpec(nlri, pamap) => { return build_flowspec_route_monitoring( peer, @@ -835,6 +848,48 @@ pub fn build_route_monitoring_from_route( build_route_monitoring(peer, prefix, pamap, is_withdrawal, path_id) } +/// Rebuild MP-BGP framing while retaining EVPN next hop and communities. +pub fn build_evpn_route_monitoring( + peer: &PeerInfo, + raw: &[u8], + pamap: &RotondaPaMap, + withdrawn: bool, + path_id: Option, +) -> Option> { + let (mut attrs, next_hop) = filter_raw_path_attributes(pamap); + let mut value = vec![0, 25, 70]; + if withdrawn { + attrs.clear(); + } else { + let nh = next_hop?; + value.push(u8::try_from(nh.len()).ok()?); + value.extend_from_slice(&nh); + value.push(0); + } + if let Some(pid) = path_id { + value.extend_from_slice(&pid.to_be_bytes()); + } + value.extend_from_slice(raw); + attrs.extend_from_slice(&[0x90, if withdrawn { 15 } else { 14 }]); + attrs.extend_from_slice(&(value.len() as u16).to_be_bytes()); + attrs.extend_from_slice(&value); + let len = 23 + attrs.len(); + if len > MAX_BGP_UPDATE_LEN { + return None; + } + let total = BMP_COMMON_HEADER_LEN + BMP_PER_PEER_HEADER_LEN + len; + let mut out = Vec::with_capacity(total); + write_common_header(&mut out, BMP_MSG_ROUTE_MONITORING, total as u32); + write_per_peer_header(&mut out, peer, None); + out.extend_from_slice(&BGP_MARKER); + out.extend_from_slice(&(len as u16).to_be_bytes()); + out.push(BGP_MSG_UPDATE); + out.extend_from_slice(&0u16.to_be_bytes()); + out.extend_from_slice(&(attrs.len() as u16).to_be_bytes()); + out.extend_from_slice(&attrs); + Some(out) +} + /// Append raw FlowSpec NLRI bytes with their RFC 8955 §4.1 length header: /// one byte for lengths < 240, else two bytes `0xFnnn` (max 4095). An /// RFC 7911 path id, when present, precedes the length header. @@ -1725,6 +1780,7 @@ pub fn build_end_of_rib_marker( afisafi: AfiSafiType, ) -> Option> { match afisafi { + AfiSafiType::L2VpnEvpn => Some(build_eor_mp_unreach(peer, afisafi)), AfiSafiType::Ipv4Unicast => Some(build_eor_ipv4(peer)), AfiSafiType::Ipv6Unicast | AfiSafiType::Ipv4FlowSpec @@ -3125,6 +3181,62 @@ mod tests { assert_eq!(pd_from_msg(&build_peer_down(&peer)), real_rd); } + #[test] + fn evpn_bmp_roundtrip_plain_and_addpath() { + use crate::roto_runtime::types::{ + explode_announcements, explode_withdrawals, + }; + use routecore::bgp::message::{SessionConfig, UpdateMessage}; + let peer = agg_test_peer(); + let mut raw = vec![5, 34]; + raw.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, 7]); + raw.extend_from_slice(&[0; 14]); + raw.extend_from_slice(&[24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, 42]); + let pamap = RotondaPaMap::new( + routecore::bgp::path_attributes::OwnedPathAttributes::new( + routecore::bgp::message::update::PduParseInfo::modern(), + vec![ + 0x40, 1, 1, 0, 0x80, 14, 9, 0, 25, 70, 4, 192, 0, 2, 1, 0, + ], + ), + ); + for pid in [None, Some(17)] { + let mut sc = SessionConfig::modern(); + if pid.is_some() { + sc.add_addpath_rxtx(AfiSafiType::L2VpnEvpn); + } + for withdrawn in [false, true] { + let bmp = build_evpn_route_monitoring( + &peer, &raw, &pamap, withdrawn, pid, + ) + .unwrap(); + let bgp = bytes::Bytes::copy_from_slice( + &bmp[BMP_COMMON_HEADER_LEN + BMP_PER_PEER_HEADER_LEN..], + ); + let update = UpdateMessage::from_octets(bgp, &sc).unwrap(); + let routes = if withdrawn { + explode_withdrawals(&update) + } else { + explode_announcements(&update) + } + .unwrap(); + assert_eq!(routes.len(), 1); + assert_eq!(routes[0].1.map(|p| p.0), pid); + match &routes[0].0 { + RotondaRoute::L2VpnEvpn(n, attrs) => { + assert_eq!(n.raw, raw); + assert_eq!(n.rd, "1:7"); + if !withdrawn { + assert_eq!(crate::units::rib_unit::evpn::EvpnAttributes::decode(attrs).next_hop, + Some("192.0.2.1".parse().unwrap())); + } + } + other => panic!("unexpected route {other:?}"), + } + } + } + } + fn agg_test_peer() -> PeerInfo { PeerInfo { peer_type: PeerType::GlobalInstance, diff --git a/src/units/bmp_tcp_out/client_handler.rs b/src/units/bmp_tcp_out/client_handler.rs index 20e7f23..7c14e46 100644 --- a/src/units/bmp_tcp_out/client_handler.rs +++ b/src/units/bmp_tcp_out/client_handler.rs @@ -477,6 +477,63 @@ pub async fn perform_initial_dump( true }) }; + // EVPN uses RD-scoped keys instead of the IP prefix tree. + if !client_gone { + for record in + rib_for_walk.evpn_records().into_iter().filter(|r| r.active) + { + let source = ingress_register_for_walk.get(record.ingress_id); + let (ingress_id, path_id) = match source { + Some(ref info) + if info.ingress_type + == Some(IngressType::BgpPath) => + { + ( + info.parent_ingress.unwrap_or(record.ingress_id), + info.path_id, + ) + } + _ => (record.ingress_id, None), + }; + let Some(info) = ingress_register_for_walk.get(ingress_id) + else { + continue; + }; + let pi = build_peer_info_for_emit( + &info, + &ingress_register_for_walk, + forward_router_info, + fan_in_peer_distinguisher, + ); + if !aggregator.has_peer(ingress_id) { + if msg_tx + .blocking_send(( + bmp_builder::build_peer_up(&pi, false), + 0, + )) + .is_err() + { + client_gone = true; + break; + } + aggregator.insert_peer(ingress_id, pi.clone()); + discovered.push((ingress_id, pi.clone())); + } + if let Some(msg) = bmp_builder::build_evpn_route_monitoring( + &pi, + &record.nlri.raw, + &record.attributes, + false, + path_id, + ) { + if msg_tx.blocking_send((msg, 1)).is_err() { + client_gone = true; + break; + } + *routes_per_ingress.entry(ingress_id).or_insert(0) += 1; + } + } + } let walk_result = match (walk_result, fs_walk_result) { (Ok(unicast), Ok(flowspec)) => Ok(unicast + flowspec), (Err(e), _) | (_, Err(e)) => Err(e), @@ -1309,6 +1366,150 @@ mod tests { ); } + #[tokio::test] + async fn evpn_dump_replays_active_addpath_and_eor() { + use crate::{ + payload::{RotondaPaMap, RotondaRoute}, + roto_runtime::Ctx, + }; + use rotonda_store::prefix_record::RouteStatus; + use std::sync::Mutex; + let register: Arc = Default::default(); + let rib = Arc::new( + Rib::new( + register.clone(), + None, + Arc::new(Mutex::new(Ctx::empty())), + ) + .unwrap(), + ); + let session = register.register(); + register.update_info( + session, + IngressInfo::new() + .with_ingress_type(IngressType::BgpViaBmp) + .with_state(IngressState::Connected) + .with_remote_addr("192.0.2.1".parse::().unwrap()) + .with_remote_asn(Asn::from_u32(65000)) + .with_remote_capabilities(vec![1, 4, 0, 25, 0, 70]) + .with_addpath_families(vec![0, 25, 70, 3]), + ); + let mut children = Vec::new(); + for pid in [11u32, 22] { + let child = register.register(); + register.update_info( + child, + IngressInfo::new() + .with_ingress_type(IngressType::BgpPath) + .with_state(IngressState::Connected) + .with_parent_ingress(session) + .with_path_id(pid), + ); + children.push(child); + let mut raw = vec![5, 34]; + raw.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, 7]); + raw.extend_from_slice(&[0; 14]); + raw.extend_from_slice(&[24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, 42]); + let attrs = RotondaPaMap::new( + routecore::bgp::path_attributes::OwnedPathAttributes::new( + routecore::bgp::message::update::PduParseInfo::modern(), + vec![ + 0x40, 1, 1, 0, 0x80, 14, 9, 0, 25, 70, 4, 192, 0, 2, + 1, 0, + ], + ), + ); + let route = RotondaRoute::L2VpnEvpn( + crate::units::rib_unit::evpn::EvpnNlri::parse(&raw).unwrap(), + attrs, + ); + rib.insert(&route, RouteStatus::Active, 0, child, true, false) + .unwrap(); + } + assert!(rib.reap_idle_path_children(Default::default()).is_empty()); + rib.withdraw_for_ingress( + children[1], + Some(routecore::bgp::types::AfiSafiType::L2VpnEvpn), + true, + ); + let (tx, mut rx) = tokio::sync::mpsc::channel(1024); + let client = Arc::new(ClientState::new( + "127.0.0.1:0".parse().unwrap(), + tx, + 10000, + 10 * 1024 * 1024, + )); + let gate = crate::comms::Gate::default(); + let metrics = Arc::new(BmpTcpOutMetrics::new(&gate)); + let reporter = Arc::new(BmpTcpOutStatusReporter::new( + "evpn-test", + metrics.clone(), + )); + assert!( + perform_initial_dump( + &client, + &rib, + ®ister, + "test", + "test", + false, + FanInPeerDistinguisher::Off, + &metrics, + &reporter + ) + .await + ); + let mut kinds = Vec::new(); + while let Ok(msg) = rx.try_recv() { + let kind = classify(&msg); + if kind == MsgKind::Route(25, 70) { + use routecore::bgp::message::{SessionConfig, UpdateMessage}; + let mut sc = SessionConfig::modern(); + sc.add_addpath_rxtx( + routecore::bgp::types::AfiSafiType::L2VpnEvpn, + ); + let update = UpdateMessage::from_octets( + bytes::Bytes::copy_from_slice(&msg[48..]), + &sc, + ) + .unwrap(); + let routes = + crate::roto_runtime::types::explode_announcements( + &update, + ) + .unwrap(); + assert_eq!(routes.len(), 1); + assert_eq!(routes[0].1.unwrap().0, 11); + } + kinds.push(kind); + } + assert_eq!( + kinds.iter().filter(|k| **k == MsgKind::PeerUp).count(), + 1 + ); + assert_eq!( + kinds + .iter() + .filter(|k| **k == MsgKind::Route(25, 70)) + .count(), + 1 + ); + assert_eq!( + kinds.iter().filter(|k| **k == MsgKind::Eor(25, 70)).count(), + 1 + ); + assert!( + kinds + .iter() + .position(|k| *k == MsgKind::Route(25, 70)) + .unwrap() + < kinds + .iter() + .position(|k| *k == MsgKind::Eor(25, 70)) + .unwrap() + ); + } + /// Full dump sequence with a unicast+flowspec peer and a flowspec-only /// peer: Initiation -> one Peer Up per peer -> unicast routes -> /// flowspec routes (never mixed per UPDATE) -> 4 EoRs per peer, and diff --git a/src/units/mrt_file_in/unit.rs b/src/units/mrt_file_in/unit.rs index dd94aad..1dc8a25 100644 --- a/src/units/mrt_file_in/unit.rs +++ b/src/units/mrt_file_in/unit.rs @@ -979,6 +979,7 @@ fn route_afisafi(route: &RotondaRoute) -> AfiSafiType { RotondaRoute::Ipv6Multicast(..) => AfiSafiType::Ipv6Multicast, RotondaRoute::Ipv4FlowSpec(..) => AfiSafiType::Ipv4FlowSpec, RotondaRoute::Ipv6FlowSpec(..) => AfiSafiType::Ipv6FlowSpec, + RotondaRoute::L2VpnEvpn(..) => AfiSafiType::L2VpnEvpn, } } @@ -993,6 +994,7 @@ fn is_supported_afisafi(afisafi: AfiSafiType) -> bool { | AfiSafiType::Ipv6Unicast | AfiSafiType::Ipv4FlowSpec | AfiSafiType::Ipv6FlowSpec + | AfiSafiType::L2VpnEvpn ) } diff --git a/src/units/rib_unit/evpn.rs b/src/units/rib_unit/evpn.rs new file mode 100644 index 0000000..3803d4f --- /dev/null +++ b/src/units/rib_unit/evpn.rs @@ -0,0 +1,333 @@ +//! EVPN monitoring records. Route identity includes the RD and never uses +//! the tenant IP prefix as a global unicast key. +use inetnum::addr::Prefix; +use serde::Serialize; +use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + +#[derive(Clone, Debug, Eq, PartialEq, Serialize)] +pub struct EvpnNlri { + pub route_type: u8, + pub rd: String, + pub ethernet_tag: Option, + pub esi: Option, + pub mac: Option, + pub prefix: Option, + pub gateway: Option, + /// Raw 24-bit label fields: VXLAN interprets these as VNIs. + pub labels: Vec, + pub raw: Vec, + #[serde(skip)] + pub(crate) key: Vec, +} + +fn hex(raw: &[u8]) -> String { + raw.iter().map(|b| format!("{b:02x}")).collect() +} +fn ip(raw: &[u8]) -> Result { + match raw.len() { + 4 => { + Ok(Ipv4Addr::from(<[u8; 4]>::try_from(raw).map_err(|_| ())?) + .into()) + } + 16 => { + Ok(Ipv6Addr::from(<[u8; 16]>::try_from(raw).map_err(|_| ())?) + .into()) + } + _ => Err(()), + } +} + +impl EvpnNlri { + pub fn parse(raw: &[u8]) -> Result { + if raw.len() < 10 || raw[1] as usize + 2 != raw.len() { + return Err(()); + } + let b = &raw[2..]; + let rd = match u16::from_be_bytes([b[0], b[1]]) { + 0 => format!( + "{}:{}", + u16::from_be_bytes([b[2], b[3]]), + u32::from_be_bytes(b[4..8].try_into().unwrap()) + ), + 1 => format!( + "{}:{}", + Ipv4Addr::from(<[u8; 4]>::try_from(&b[2..6]).unwrap()), + u16::from_be_bytes([b[6], b[7]]) + ), + 2 => format!( + "{}:{}", + u32::from_be_bytes(b[2..6].try_into().unwrap()), + u16::from_be_bytes([b[6], b[7]]) + ), + _ => hex(&b[..8]), + }; + let mut n = Self { + route_type: raw[0], + rd, + ethernet_tag: None, + esi: None, + mac: None, + prefix: None, + gateway: None, + labels: vec![], + raw: raw.to_vec(), + key: raw.to_vec(), + }; + match n.route_type { + 1 => { + if b.len() != 25 { + return Err(()); + } + n.esi = Some(hex(&b[8..18])); + n.ethernet_tag = + Some(u32::from_be_bytes(b[18..22].try_into().unwrap())); + n.labels.push(u32::from_be_bytes([0, b[22], b[23], b[24]])); + n.key = vec![1]; + n.key.extend_from_slice(&b[..22]); + } + 3 => { + if b.len() < 13 + || !matches!((b[12], b.len()), (32, 17) | (128, 29)) + { + return Err(()); + } + n.ethernet_tag = + Some(u32::from_be_bytes(b[8..12].try_into().unwrap())); + } + 4 => { + if b.len() < 19 + || !matches!((b[18], b.len()), (32, 23) | (128, 35)) + { + return Err(()); + } + n.esi = Some(hex(&b[8..18])); + } + _ => (), + } + if matches!(n.route_type, 2 | 5) { + if b.len() < 25 { + return Err(()); + } + n.esi = Some(hex(&b[8..18])); + n.ethernet_tag = + Some(u32::from_be_bytes(b[18..22].try_into().unwrap())); + n.key = vec![n.route_type]; + n.key.extend_from_slice(&b[..8]); + n.key.extend_from_slice(&b[18..22]); + let labels_at; + if n.route_type == 2 { + if b.len() < 33 || b[22] != 48 { + return Err(()); + } + let ip_len = match b[29] { + 0 => 0, + 32 => 4, + 128 => 16, + _ => return Err(()), + }; + labels_at = 30 + ip_len; + if b.len() != labels_at + 3 && b.len() != labels_at + 6 { + return Err(()); + } + n.mac = Some( + b[23..29] + .iter() + .map(|v| format!("{v:02x}")) + .collect::>() + .join(":"), + ); + if ip_len != 0 { + n.prefix = Some( + Prefix::new(ip(&b[30..labels_at])?, b[29]) + .map_err(|_| ())?, + ); + } + n.key.extend_from_slice(&b[22..labels_at]); + } else { + // RFC 9136 carries a full 4/16-octet address, regardless of prefix length. + let size = match b.len() { + 34 => 4, + 58 => 16, + _ => return Err(()), + }; + let prefix = Prefix::new(ip(&b[23..23 + size])?, b[22]) + .map_err(|_| ())?; + n.prefix = Some(prefix); + n.gateway = Some(ip(&b[23 + size..23 + 2 * size])?); + labels_at = 23 + 2 * size; + n.key.extend_from_slice(prefix.to_string().as_bytes()); + } + for l in b[labels_at..].chunks_exact(3) { + n.labels.push(u32::from_be_bytes([0, l[0], l[1], l[2]])); + } + } + Ok(n) + } +} + +#[derive(Clone, Serialize)] +pub struct EvpnRecord { + pub ingress_id: crate::ingress::IngressId, + pub ltime: u64, + pub active: bool, + pub nlri: EvpnNlri, + pub attributes: crate::payload::RotondaPaMap, +} + +#[cfg(test)] +mod tests { + use super::*; + fn type5(rd: u8, label: u8) -> Vec { + let mut b = vec![5, 34]; + b.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, rd]); + b.extend_from_slice(&[0; 14]); + b.push(24); + b.extend_from_slice(&[10, 0, 0, 0, 0, 0, 0, 0, 0, 0, label]); + b + } + #[test] + fn tenant_identity_and_label_changes() { + let a = EvpnNlri::parse(&type5(1, 10)).unwrap(); + let b = EvpnNlri::parse(&type5(2, 10)).unwrap(); + let c = EvpnNlri::parse(&type5(1, 20)).unwrap(); + assert_eq!(a.prefix, b.prefix); + assert_ne!(a.key, b.key); + assert_eq!(a.key, c.key); + assert_eq!(a.labels, vec![10]); + } + #[test] + fn evpn_mac_ip_two_vnis_and_ipv6_prefix() { + let mut raw = vec![2, 52]; + raw.extend_from_slice(&[0; 22]); + raw.push(48); + raw.extend_from_slice(&[0, 1, 2, 3, 4, 5]); + raw.push(128); + raw.extend_from_slice(&Ipv6Addr::LOCALHOST.octets()); + raw.extend_from_slice(&[0, 0, 10, 0, 0, 20]); + let route = EvpnNlri::parse(&raw).unwrap(); + assert_eq!(route.labels, vec![10, 20]); + assert_eq!(route.mac.as_deref(), Some("00:01:02:03:04:05")); + assert_eq!(route.prefix.unwrap().to_string(), "::1/128"); + raw[10] = 1; // ESI is not part of RT-2 route identity. + raw[53] = 30; + assert_eq!(route.key, EvpnNlri::parse(&raw).unwrap().key); + + let mut raw = vec![5, 58]; + raw.extend_from_slice(&[0; 22]); + raw.push(64); + raw.extend_from_slice( + &"2001:db8::".parse::().unwrap().octets(), + ); + raw.extend_from_slice(&Ipv6Addr::UNSPECIFIED.octets()); + raw.extend_from_slice(&[0, 1, 0]); + let route = EvpnNlri::parse(&raw).unwrap(); + assert_eq!(route.prefix.unwrap().to_string(), "2001:db8::/64"); + assert_eq!(route.labels, vec![256]); + } + + #[test] + fn evpn_overlay_metadata() { + use routecore::bgp::{ + message::update::PduParseInfo, + path_attributes::OwnedPathAttributes, + }; + let attrs = + crate::payload::RotondaPaMap::new(OwnedPathAttributes::new( + PduParseInfo::modern(), + vec![ + 0xc0, 16, 24, 0, 2, 0xfd, 0xe8, 0, 0, 0, 10, 2, 2, 0, 1, + 0, 0, 0, 20, 6, 3, 0, 1, 2, 3, 4, 5, + ], + )); + let overlay = EvpnAttributes::decode(&attrs); + assert_eq!(overlay.route_targets, vec!["65000:10", "65536:20"]); + assert_eq!(overlay.router_mac.as_deref(), Some("00:01:02:03:04:05")); + } + + #[test] + fn truncated_routes_are_rejected() { + let b = type5(1, 10); + for len in 0..b.len() { + assert!(EvpnNlri::parse(&b[..len]).is_err()); + } + } +} + +/// Extended communities used when correlating MAC-VRF and IP-VRF routes. +#[derive(Default, Serialize)] +pub struct EvpnAttributes { + pub route_targets: Vec, + pub router_mac: Option, + pub next_hop: Option, +} +impl EvpnAttributes { + pub fn decode(attributes: &crate::payload::RotondaPaMap) -> Self { + let mut out = Self::default(); + let raw = attributes.as_ref(); + let mut b = raw.get(2..).unwrap_or_default(); + while b.len() >= 3 { + let (header, len) = if b[0] & 0x10 != 0 { + if b.len() < 4 { + break; + } + (4, u16::from_be_bytes([b[2], b[3]]) as usize) + } else { + (3, b[2] as usize) + }; + if b.len() < header + len { + break; + } + let v = &b[header..header + len]; + match b[1] { + 16 => { + for c in v.chunks_exact(8) { + if c[1] == 2 && c[0] <= 2 { + let rt = match c[0] { + 0 => format!( + "{}:{}", + u16::from_be_bytes([c[2], c[3]]), + u32::from_be_bytes( + c[4..8].try_into().unwrap() + ) + ), + 1 => format!( + "{}:{}", + Ipv4Addr::from( + <[u8; 4]>::try_from(&c[2..6]) + .unwrap() + ), + u16::from_be_bytes([c[6], c[7]]) + ), + _ => format!( + "{}:{}", + u32::from_be_bytes( + c[2..6].try_into().unwrap() + ), + u16::from_be_bytes([c[6], c[7]]) + ), + }; + out.route_targets.push(rt); + } + if c[0..2] == [6, 3] { + out.router_mac = Some( + c[2..] + .iter() + .map(|v| format!("{v:02x}")) + .collect::>() + .join(":"), + ); + } + } + } + 14 if v.len() >= 4 && v[..3] == [0, 25, 70] => { + if let Some(nh) = v.get(4..4 + v[3] as usize) { + out.next_hop = ip(nh).ok(); + } + } + _ => (), + } + b = &b[header + len..]; + } + out + } +} diff --git a/src/units/rib_unit/http_ng.rs b/src/units/rib_unit/http_ng.rs index a089104..af699d8 100644 --- a/src/units/rib_unit/http_ng.rs +++ b/src/units/rib_unit/http_ng.rs @@ -39,6 +39,7 @@ use crate::{ /// Add ingress register specific endpoints to a HTTP API pub fn register_routes(router: &mut Api) { + router.add_get("/ribs/l2vpnevpn/routes", search_evpn); router.add_get( "/ribs/ipv4unicast/routes/{prefix}/{prefix_len}", search_ipv4unicast, @@ -889,7 +890,9 @@ async fn run_best_path( Ok(( [("content-type", OutputFormat::Json.content_type())], serde_json::to_string(&body).map_err(|e| { - ApiError::InternalServerError(format!("serialization failed: {e}")) + ApiError::InternalServerError(format!( + "serialization failed: {e}" + )) })?, )) } @@ -971,3 +974,151 @@ async fn best_path_ipv6_addr( ) .await } + +#[derive(Default, Deserialize)] +#[serde(deny_unknown_fields)] +struct EvpnFilter { + rd: Option, + route_target: Option, + route_type: Option, + vni: Option, + prefix: Option, + ingress_id: Option, + #[serde(default)] + include_withdrawn: bool, +} + +async fn search_evpn( + Query(filter): Query, + state: State, +) -> Result { + let rib = load_rib(&state)?; + let permit = super::rib::DumpGuard::try_enter().ok_or_else(|| { + ApiError::ServiceUnavailable("too many concurrent RIB queries".into()) + })?; + let data = tokio::task::spawn_blocking(move || { + let _permit = permit; + evpn_rows(&rib, &filter) + }) + .await + .map_err(|e| ApiError::InternalServerError(e.to_string()))?; + let body = serde_json::to_vec(&serde_json::json!({"data": data})) + .map_err(|e| ApiError::InternalServerError(e.to_string()))?; + Ok( + ([("content-type", OutputFormat::Json.content_type())], body) + .into_response(), + ) +} + +fn evpn_rows( + rib: &super::rib::Rib, + filter: &EvpnFilter, +) -> Vec { + let mut records = rib.evpn_records(); + records.sort_by(|a, b| { + (&a.nlri.key, a.ingress_id).cmp(&(&b.nlri.key, b.ingress_id)) + }); + let mut rows = Vec::new(); + for record in records { + let n = &record.nlri; + if (!record.active && !filter.include_withdrawn) + || filter.rd.as_ref().is_some_and(|v| v != &n.rd) + || filter.route_type.is_some_and(|v| v != n.route_type) + || filter.vni.is_some_and(|v| !n.labels.contains(&v)) + || filter.prefix.is_some_and(|v| Some(v) != n.prefix) + || filter.ingress_id.is_some_and(|v| v != record.ingress_id) + { + continue; + } + let overlay = super::evpn::EvpnAttributes::decode(&record.attributes); + if filter + .route_target + .as_ref() + .is_some_and(|v| !overlay.route_targets.contains(v)) + { + continue; + } + let source = rib.ingress_register.get(record.ingress_id); + let path_source = source + .as_ref() + .filter(|s| s.ingress_type == Some(IngressType::BgpPath)); + rows.push(serde_json::json!({ + "route": record, + "overlay": overlay, + "source_ingress_id": path_source.and_then(|s| s.parent_ingress).unwrap_or(record.ingress_id), + "path_id": path_source.and_then(|s| s.path_id), + })); + } + rows +} + +#[cfg(test)] +mod evpn_tests { + use super::*; + use crate::{ + payload::{RotondaPaMap, RotondaRoute}, + roto_runtime::Ctx, + }; + use rotonda_store::prefix_record::RouteStatus; + use std::sync::{Arc, Mutex}; + #[test] + fn evpn_api_filters_overlapping_tenants() { + let rib = super::super::rib::Rib::new( + Default::default(), + None, + Arc::new(Mutex::new(Ctx::empty())), + ) + .unwrap(); + for tenant in [1u8, 2] { + let mut raw = vec![5, 34]; + raw.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, tenant]); + raw.extend_from_slice(&[0; 14]); + raw.extend_from_slice(&[ + 24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, tenant, + ]); + let attributes = RotondaPaMap::new( + routecore::bgp::path_attributes::OwnedPathAttributes::new( + routecore::bgp::message::update::PduParseInfo::modern(), + vec![0xc0, 16, 8, 0, 2, 0xfd, 0xe8, 0, 0, 0, tenant], + ), + ); + let route = RotondaRoute::L2VpnEvpn( + super::super::evpn::EvpnNlri::parse(&raw).unwrap(), + attributes, + ); + rib.insert(&route, RouteStatus::Active, 0, 10, true, false) + .unwrap(); + } + assert_eq!(evpn_rows(&rib, &EvpnFilter::default()).len(), 2); + let filter: EvpnFilter = serde_json::from_value(serde_json::json!({ + "route_target": "65000:1", "prefix": "10.0.0.0/24", "route_type": 5, + "vni": 1, "rd": "1:1", "ingress_id": 10 + })).unwrap(); + let rows = evpn_rows(&rib, &filter); + assert_eq!(rows.len(), 1); + assert_eq!(rows[0]["route"]["nlri"]["rd"], "1:1"); + assert_eq!(rows[0]["source_ingress_id"], 10); + let filter = EvpnFilter { + vni: Some(2), + ..filter + }; + assert!(evpn_rows(&rib, &filter).is_empty()); + rib.withdraw_for_ingress(10, Some(AfiSafiType::L2VpnEvpn), true); + assert!(evpn_rows(&rib, &EvpnFilter::default()).is_empty()); + assert_eq!( + evpn_rows( + &rib, + &EvpnFilter { + include_withdrawn: true, + ..Default::default() + } + ) + .len(), + 2 + ); + assert!(serde_json::from_value::( + serde_json::json!({"best_path": true}) + ) + .is_err()); + } +} diff --git a/src/units/rib_unit/mod.rs b/src/units/rib_unit/mod.rs index 1bd4b0d..840a1f2 100644 --- a/src/units/rib_unit/mod.rs +++ b/src/units/rib_unit/mod.rs @@ -1,3 +1,4 @@ +pub mod evpn; mod http_ng; pub use http_ng::{Include, QueryFilter}; mod metrics; diff --git a/src/units/rib_unit/rib.rs b/src/units/rib_unit/rib.rs index e952817..55f8e69 100644 --- a/src/units/rib_unit/rib.rs +++ b/src/units/rib_unit/rib.rs @@ -160,6 +160,7 @@ type RotoHttpFilter = roto::TypedFunc< #[derive(Clone)] pub struct Rib { + evpn: Arc, IngressId), super::evpn::EvpnRecord>>>, unicast: Arc>, multicast: Arc>, /// FlowSpec rules (SAFI 133, v4+v6 in the one dual-family store), keyed @@ -241,6 +242,7 @@ fn reset_peer_gauge(mui: IngressId, family: Option) { AfiSafiType::Ipv6Multicast => (2, 2), AfiSafiType::Ipv4FlowSpec => (1, 133), AfiSafiType::Ipv6FlowSpec => (2, 133), + AfiSafiType::L2VpnEvpn => (25, 70), // Families we never store; nothing was ever counted. _ => return, }; @@ -274,6 +276,7 @@ impl Rib { Ok(Rib { unicast: Arc::new(Some(Store::try_default()?)), multicast: Arc::new(Some(Store::try_default()?)), + evpn: Arc::default(), flowspec: Arc::new(Some(FlowSpecStore::try_default()?)), flowspec_rule_counts, ingress_register, @@ -554,6 +557,53 @@ impl Rib { deduplicate_path_attributes: bool, ) -> Result { let res = match val { + RotondaRoute::L2VpnEvpn(n, attributes) => { + let mut records = + self.evpn.lock().unwrap_or_else(|e| e.into_inner()); + let key = (n.key.clone(), ingress_id); + let existed = records.contains_key(&key); + let was_active = records.get(&key).is_some_and(|r| r.active); + let active = route_status == RouteStatus::Active; + if let Some(gauge) = peer_gauge(ingress_id) { + if active && !was_active { + gauge.add_adj_rib_in((25, 70), 1); + } + if !active && was_active { + gauge.sub_adj_rib_in((25, 70), 1); + } + } + if active { + records.insert( + key, + super::evpn::EvpnRecord { + ingress_id, + ltime, + active, + nlri: n.clone(), + attributes: if deduplicate_path_attributes { + attributes + .dedup_with(&self.path_attribute_interner) + } else { + attributes.clone() + }, + }, + ); + } else if retain_withdrawn_attributes { + if let Some(record) = records.get_mut(&key) { + record.active = false; + record.ltime = ltime; + } + } else { + records.remove(&key); + } + Ok(UpsertReport { + cas_count: 0, + prefix_new: !existed && active, + mui_new: !existed && active, + mui_count: usize::from(active), + }) + } + RotondaRoute::Ipv4Unicast(n, ..) => self.insert_prefix( &n.prefix(), Multicast(false), @@ -1271,6 +1321,26 @@ impl Rib { self.compact_withdrawn_attributes_for_ingresses(ids); } + { + let mut records = + self.evpn.lock().unwrap_or_else(|e| e.into_inner()); + let evpn_ingresses: HashSet<_> = ids + .iter() + .filter(|(_, family)| { + family.is_none() + || *family == Some(AfiSafiType::L2VpnEvpn) + }) + .map(|(id, _)| *id) + .collect(); + records.retain(|(_, ingress), record| { + if evpn_ingresses.contains(ingress) { + record.active = false; + retain_withdrawn_attributes + } else { + true + } + }); + } for (ingress_id, specific_afisafi) in ids { debug!("withdraw_for_ingress for {ingress_id}"); reset_peer_gauge(*ingress_id, *specific_afisafi); @@ -1383,6 +1453,7 @@ impl Rib { } } + Some(AfiSafiType::L2VpnEvpn) => (), // handled above afisafi => { // Reachable for families we never store (they are // dropped at the TryFrom chokepoint), so log instead @@ -1407,7 +1478,20 @@ impl Rib { /// retained previous session before reusing an id. Synthesized BMP peers /// that mint a fresh ingress id every session must take this path so /// mark-withdraw does not leak one record slot per announced prefix. + pub fn evpn_records(&self) -> Vec { + self.evpn + .lock() + .unwrap_or_else(|e| e.into_inner()) + .values() + .cloned() + .collect() + } + pub fn remove_for_ingresses(&self, ids: &[IngressId]) { + self.evpn + .lock() + .unwrap_or_else(|e| e.into_inner()) + .retain(|(_, ingress), _| !ids.contains(ingress)); if ids.is_empty() { return; } @@ -1513,6 +1597,7 @@ impl Rib { Some(AfiSafiType::Ipv6FlowSpec) => { v6_fs.insert(*id); } + Some(AfiSafiType::L2VpnEvpn) => (), // separate EVPN store afisafi => { warn!( "no support to compact withdrawn attributes for {:?} yet", @@ -1807,7 +1892,14 @@ impl Rib { // chunks under short-lived guards, the way the jsonl dump does it: // one guard held across the whole table would pin concurrent churn // garbage for the length of the walk. - let mut active_muis: HashSet = HashSet::new(); + let mut active_muis: HashSet = self + .evpn + .lock() + .unwrap_or_else(|e| e.into_inner()) + .values() + .filter(|r| r.active) + .map(|r| r.ingress_id) + .collect(); for store in [self.unicast.as_ref(), self.multicast.as_ref()] .into_iter() .flatten() @@ -4129,6 +4221,56 @@ mod tests { // ------------ FlowSpec store ------------------------------------------ + #[test] + fn evpn_tenants_paths_and_peer_lifecycle() { + use super::super::evpn::EvpnNlri; + fn route(rd: u8, vni: u8) -> RotondaRoute { + let mut raw = vec![5, 34]; + raw.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, rd]); + raw.extend_from_slice(&[0; 14]); + raw.extend_from_slice(&[24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, vni]); + RotondaRoute::L2VpnEvpn( + EvpnNlri::parse(&raw).unwrap(), + RotondaPaMap::empty_path_attributes(), + ) + } + let rib = test_rib(); + for (rd, ingress) in [(1, 1), (2, 1), (1, 2)] { + rib.insert( + &route(rd, 10), + RouteStatus::Active, + 0, + ingress, + true, + false, + ) + .unwrap(); + } + assert_eq!(rib.evpn_records().len(), 3); + rib.insert(&route(1, 20), RouteStatus::Active, 1, 1, true, false) + .unwrap(); + assert_eq!(rib.evpn_records().len(), 3); + rib.insert(&route(1, 0), RouteStatus::Withdrawn, 2, 1, true, false) + .unwrap(); + let rows = rib.evpn_records(); + assert_eq!(rows.iter().filter(|r| r.active).count(), 2); + assert_eq!( + rows.iter().find(|r| !r.active).unwrap().nlri.labels, + vec![20] + ); + rib.withdraw_for_ingress(1, Some(AfiSafiType::Ipv4Unicast), false); + assert_eq!(rib.evpn_records().len(), 3); + rib.withdraw_for_ingress(1, Some(AfiSafiType::L2VpnEvpn), false); + assert_eq!(rib.evpn_records().len(), 1); + rib.insert(&route(2, 30), RouteStatus::Active, 3, 1, false, false) + .unwrap(); + assert_eq!(rib.evpn_records().len(), 2); + rib.remove_for_ingresses(&[2]); + assert_eq!(rib.evpn_records().len(), 1); + rib.withdraw_for_ingress(1, None, false); + assert!(rib.evpn_records().is_empty()); + } + fn test_rib() -> Rib { Rib::new(Default::default(), None, Arc::new(Mutex::new(Ctx::empty()))) .unwrap() From b3529738f5a7b1432b30d14708d6c6b64c55e232 Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 25 Sep 2026 13:14:41 -0500 Subject: [PATCH 02/10] Add EVPN route queries to netom-cli --- doc/netom-cli.1 | 7 + docs/cli.md | 29 ++++ docs/evpn.md | 5 +- src/bin/netom-cli/commands/evpn.rs | 238 +++++++++++++++++++++++++++++ src/bin/netom-cli/commands/mod.rs | 1 + src/bin/netom-cli/tree.rs | 144 ++++++++++++++++- test-data/cli/evpn-routes.json | 5 + 7 files changed, 425 insertions(+), 4 deletions(-) create mode 100644 src/bin/netom-cli/commands/evpn.rs create mode 100644 test-data/cli/evpn-routes.json diff --git a/doc/netom-cli.1 b/doc/netom-cli.1 index 888bdb2..ee0ff53 100644 --- a/doc/netom-cli.1 +++ b/doc/netom-cli.1 @@ -179,6 +179,13 @@ as above; the other filters are not implemented for FlowSpec. .B show ipv6 bgp ... As above, for IPv6. .TP +.B show evpn +EVPN routes with RD, MAC/prefix, VNIs, next hop, route targets, peer/path, +and active/withdrawn state. Optional filters, in order: rd VALUE or +route-target VALUE; route-type NUMBER; one of vni NUMBER, prefix PREFIX, +or ingress ID; include-withdrawn; detail. Detail displays all returned +fields and attributes. With --json, emits the buffered data-array response. +.TP .B show bmp routers Monitored routers feeding BMP to this daemon. .TP diff --git a/docs/cli.md b/docs/cli.md index ef3f7f6..366145f 100644 --- a/docs/cli.md +++ b/docs/cli.md @@ -61,6 +61,35 @@ Commands: exit Exit the CLI ``` +## EVPN routes + +`show evpn` queries the [EVPN RIB](evpn.md), including both IPv4 and IPv6 +routes. The table shows route type, RD, MAC/prefix, both VNI fields, next hop, +route targets, source peer/path ID, and active/withdrawn state. + +```sh +netom-cli show evpn +netom-cli show evpn route-target 65000:100 route-type 5 +netom-cli show evpn rd 192.0.2.1:100 vni 50000 +netom-cli show evpn prefix 2001:db8::/64 +netom-cli show evpn ingress 101 include-withdrawn detail +netom-cli --json show evpn route-target 65000:100 +``` + +Combine filters in this order, omitting stages as needed: + +1. `rd ` or `route-target `. +2. `route-type <1-255>` (type 2 MAC/IP, type 5 IP prefix). +3. One of `vni <0-16777215>`, `prefix
`, or `ingress `. +4. `include-withdrawn`, then `detail`. + +`ingress` matches the stored ingress ID, including an ADD-PATH child ID. +`include-withdrawn` requires retained withdrawn records on the daemon. +`detail` displays all returned fields, including ESI, gateway, Router's MAC, +raw NLRI, and attributes. `--json` passes through the original buffered +`{"data": [...]}` response. VNI columns contain raw 24-bit label fields; +interpret them as VNIs for VXLAN, not as decoded MPLS label numbers. + ## Finding the daemon In order: `--url`, `$NETOM_URL`, `-c `, `./netom.conf`, diff --git a/docs/evpn.md b/docs/evpn.md index dc2ed43..b2c306d 100644 --- a/docs/evpn.md +++ b/docs/evpn.md @@ -81,8 +81,9 @@ BMP output includes EVPN in live rebuilt updates and initial RIB dumps, retaining next hops, communities, and ADD-PATH IDs. Synthetic Peer Up messages and End-of-RIB markers include EVPN when advertised by the source peer. BGP4MP MRT updates use the same decoder; TABLE_DUMP_V2 EVPN import is not -implemented. The existing CLI route commands and ClickHouse route schema do -not expose EVPN; use the HTTP API or BMP output. +implemented. Use [`netom-cli show evpn`](cli.md#evpn-routes) for interactive +inspection. +The ClickHouse route schema does not expose EVPN. Roto route filters can use `is_evpn()` and `evpn_rd()`. `fmt_prefix()` returns the type 2 host or type 5 prefix; EVPN routes without an IP prefix return diff --git a/src/bin/netom-cli/commands/evpn.rs b/src/bin/netom-cli/commands/evpn.rs new file mode 100644 index 0000000..2ef7ff3 --- /dev/null +++ b/src/bin/netom-cli/commands/evpn.rs @@ -0,0 +1,238 @@ +//! EVPN route queries, independent of the daemon library. +use std::io::Write; + +use crate::{ + error::CliError, + render::{fmt, left, Col, Table}, + session::Session, + tree::{Captures, Flag, Value}, +}; + +/// Percent-encode a query value without adding an HTTP dependency to the CLI. +fn encode(value: &str) -> String { + value + .bytes() + .map(|b| { + if b.is_ascii_alphanumeric() || b"-._~".contains(&b) { + (b as char).to_string() + } else { + format!("%{b:02X}") + } + }) + .collect() +} + +fn query_path(c: &Captures) -> String { + let mut params = Vec::new(); + for arg in &c.args { + let (key, value) = match arg { + Value::Rd(v) => ("rd", v.clone()), + Value::RouteTarget(v) => ("route_target", v.clone()), + Value::RouteType(v) => ("route_type", v.to_string()), + Value::Vni(v) => ("vni", v.to_string()), + Value::Prefix(ip, len) => ("prefix", format!("{ip}/{len}")), + Value::IngressId(v) => ("ingress_id", v.to_string()), + _ => continue, + }; + params.push(format!("{key}={}", encode(&value))); + } + if c.flags.contains(&Flag::IncludeWithdrawn) { + params.push("include_withdrawn=true".into()); + } + let mut path = "/api/v1/ribs/l2vpnevpn/routes".to_string(); + if !params.is_empty() { + path.push('?'); + path.push_str(¶ms.join("&")); + } + path +} + +pub fn routes(session: &mut Session, c: &Captures) -> Result<(), CliError> { + let path = query_path(c); + if session.json { + return session.passthrough(&path); + } + let body = session.get(&path)?.body_string()?; + let mut out = session.writer(); + render(&mut out, &body, c.detail())?; + out.finish()?; + Ok(()) +} + +static COLS: &[Col] = &[ + left("Type", 4), + left("RD", 10), + left("MAC / Prefix", 18), + left("VNI(s)", 6), + left("Next hop", 15), + left("Route targets", 13), + left("Peer / Path", 11), + left("State", 6), +]; + +fn scalar(v: &serde_json::Value) -> String { + match v { + serde_json::Value::Null => "-".into(), + serde_json::Value::String(s) => s.clone(), + _ => v.to_string(), + } +} +fn list(v: &serde_json::Value) -> String { + v.as_array() + .filter(|a| !a.is_empty()) + .map(|a| a.iter().map(scalar).collect::>().join(", ")) + .unwrap_or_else(|| "-".into()) +} + +fn render( + out: &mut W, + body: &str, + detail: bool, +) -> Result<(), CliError> { + let value: serde_json::Value = + serde_json::from_str(body).map_err(|e| { + CliError::Transport(format!("Invalid EVPN JSON: {e}")) + })?; + let rows = value["data"].as_array().ok_or_else(|| { + CliError::Transport( + "Invalid EVPN response: missing data array".into(), + ) + })?; + if detail { + for row in rows { + writeln!(out, "{}", serde_json::to_string_pretty(row).unwrap())?; + } + } else { + let mut table = Table::fit(out, COLS); + for row in rows { + let route = &row["route"]; + let nlri = &route["nlri"]; + let overlay = &row["overlay"]; + let destination = [nlri["mac"].as_str(), nlri["prefix"].as_str()] + .into_iter() + .flatten() + .collect::>() + .join(" / "); + let mut peer = scalar(&row["source_ingress_id"]); + if let Some(pid) = row["path_id"].as_u64() { + peer.push_str(&format!(" path {pid}")); + } + table.row(&[ + scalar(&nlri["route_type"]), + scalar(&nlri["rd"]), + if destination.is_empty() { + "-".into() + } else { + destination + }, + list(&nlri["labels"]), + scalar(&overlay["next_hop"]), + list(&overlay["route_targets"]), + peer, + match route["active"].as_bool() { + Some(true) => "active", + Some(false) => "withdrawn", + None => "-", + } + .into(), + ])?; + } + table.finish()?; + } + writeln!(out, "\nTotal EVPN routes {}", fmt::count(rows.len() as u64))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::tree; + const BODY: &str = include_str!(concat!( + env!("CARGO_MANIFEST_DIR"), + "/test-data/cli/evpn-routes.json" + )); + + #[test] + fn evpn_query_uses_api_names_and_encodes_ipv6() { + let (_, captures) = tree::resolve("show evpn route-target 65000:10 route-type 5 prefix 2001:db8::/64 include-withdrawn detail").unwrap(); + let path = query_path(&captures); + assert_eq!(path, "/api/v1/ribs/l2vpnevpn/routes?route_target=65000%3A10&route_type=5&prefix=2001%3Adb8%3A%3A%2F64&include_withdrawn=true"); + assert!(captures.detail()); + let (_, captures) = + tree::resolve("show evpn rd 192.0.2.1:20 vni 50000").unwrap(); + assert_eq!( + query_path(&captures), + "/api/v1/ribs/l2vpnevpn/routes?rd=192.0.2.1%3A20&vni=50000" + ); + let (_, captures) = tree::resolve("show evpn ingress 101").unwrap(); + assert_eq!( + query_path(&captures), + "/api/v1/ribs/l2vpnevpn/routes?ingress_id=101" + ); + assert_eq!( + query_path(&Captures::default()), + "/api/v1/ribs/l2vpnevpn/routes" + ); + } + + #[test] + fn evpn_table_and_detail_preserve_overlay_information() { + let mut out = Vec::new(); + render(&mut out, BODY, false).unwrap(); + let text = String::from_utf8(out).unwrap(); + for expected in [ + "65000:10", + "00:11:22:33:44:55 / 10.0.0.1/32", + "10010, 50000", + "7 path 11", + "2001:db8::/64", + "withdrawn", + "Total EVPN routes 3", + ] { + assert!(text.contains(expected), "missing {expected}: {text}"); + } + let mut out = Vec::new(); + render(&mut out, BODY, true).unwrap(); + let text = String::from_utf8(out).unwrap(); + assert!(text.contains("aa:bb:cc:dd:ee:ff")); + assert!(text.contains("\"ingress_id\": 101")); + assert!(text.contains("\"gateway\": \"::\"")); + assert!(!text.contains("VNI(s)")); + } + + #[test] + fn evpn_empty_and_invalid_responses() { + let mut out = Vec::new(); + render(&mut out, r#"{"data":[]}"#, false).unwrap(); + assert!(String::from_utf8(out) + .unwrap() + .contains("Total EVPN routes 0")); + assert!(render(&mut Vec::new(), "not json", false).is_err()); + assert!(render(&mut Vec::new(), r#"{"data":{}}"#, false).is_err()); + } + + #[test] + fn evpn_completion_validation_and_abbreviation() { + assert!(tree::resolve("sh ev rd 65000:1 route-ty 2 vni 100 detail") + .is_ok()); + assert!(tree::candidates("show evpn ") + .iter() + .any(|c| c.insert == "route-target")); + assert!(tree::candidates("show evpn rd 65000:1 ") + .iter() + .any(|c| c.insert == "vni")); + for command in [ + "show evpn rd", + "show evpn rd 65536:65536", + "show evpn rd bad&x=y:1", + "show evpn route-type 256", + "show evpn route-type 0", + "show evpn vni 16777216", + "show evpn source bmp", + "show evpn best", + ] { + assert!(tree::resolve(command).is_err(), "accepted {command}"); + } + assert!(tree::all_commands().iter().any(|(c, _)| c == "show evpn")); + } +} diff --git a/src/bin/netom-cli/commands/mod.rs b/src/bin/netom-cli/commands/mod.rs index 9aff962..fb2a470 100644 --- a/src/bin/netom-cli/commands/mod.rs +++ b/src/bin/netom-cli/commands/mod.rs @@ -6,4 +6,5 @@ pub mod bgp; pub mod bmp; +pub mod evpn; pub mod system; diff --git a/src/bin/netom-cli/tree.rs b/src/bin/netom-cli/tree.rs index 955c89c..202bca5 100644 --- a/src/bin/netom-cli/tree.rs +++ b/src/bin/netom-cli/tree.rs @@ -33,6 +33,10 @@ pub enum ArgKind { Asn, /// A standard community: `65000:100`, `0x1a2b3c4d`, or a well-known name. Community, + Rd, + RouteTarget, + RouteType, + Vni, } impl ArgKind { @@ -44,6 +48,9 @@ impl ArgKind { ArgKind::IngressId => "<0-4294967295>", ArgKind::Asn => "<1-4294967295>", ArgKind::Community => "", + ArgKind::Rd | ArgKind::RouteTarget => "", + ArgKind::RouteType => "<1-255>", + ArgKind::Vni => "<0-16777215>", } } @@ -56,6 +63,38 @@ impl ArgKind { let max = if addr.is_ipv4() { 32 } else { 128 }; (len <= max).then_some(Value::Prefix(addr, len)) } + ArgKind::Rd | ArgKind::RouteTarget => { + let (admin, assigned) = tok.split_once(':')?; + let assigned: u32 = assigned.parse().ok()?; + let admin = if let Ok(ip) = + admin.parse::() + { + if assigned > u16::MAX as u32 { + return None; + } + ip.to_string() + } else { + let asn: u32 = admin.parse().ok()?; + if asn > u16::MAX as u32 && assigned > u16::MAX as u32 { + return None; + } + asn.to_string() + }; + let value = format!("{admin}:{assigned}"); + Some(if self == ArgKind::Rd { + Value::Rd(value) + } else { + Value::RouteTarget(value) + }) + } + ArgKind::RouteType => { + let n: u8 = tok.parse().ok()?; + (n > 0).then_some(Value::RouteType(n)) + } + ArgKind::Vni => { + let n: u32 = tok.parse().ok()?; + (n <= 0xff_ffff).then_some(Value::Vni(n)) + } ArgKind::Ip => tok.parse().ok().map(Value::Ip), ArgKind::IngressId => tok.parse().ok().map(Value::IngressId), ArgKind::Asn => { @@ -102,6 +141,10 @@ pub enum Value { IngressId(u32), Asn(u32), Community(String), + Rd(String), + RouteTarget(String), + RouteType(u8), + Vni(u32), } /// Static context a node contributes when traversed, so that one subtree can @@ -118,6 +161,7 @@ pub enum Flag { /// Print every attribute of each path rather than one table row. Set by /// the `detail` keyword; like `best` it rides the shared filter subtree. Detail, + IncludeWithdrawn, } #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -671,6 +715,13 @@ pub static ROOT: &[Node] = &[ ]; static SHOW: &[Node] = &[ + Node { + kw: Kw::Lit("evpn"), + help: "EVPN routes and tenant overlays", + set: None, + run: Some(commands::evpn::routes), + children: EVPN_TENANT, + }, lit!("ip", "IPv4 information", SHOW_IP), lit!("ipv6", "IPv6 information", SHOW_IPV6), lit!("bmp", "BMP monitoring information", SHOW_BMP), @@ -825,8 +876,12 @@ static BGP_BEST_ADDR: &[Node] = &[Node { /// The filters narrow the candidate set: "what would win if only BGP-learned /// routes existed". `neighbors routes` is deliberately not among them -- /// narrowing the candidates to one peer makes the decision process vacuous. -static BGP_BEST_FILTERS: &[Node] = - &[FILTER_SOURCE, FILTER_INGRESS, FILTER_ORIGIN_AS, FILTER_COMMUNITY]; +static BGP_BEST_FILTERS: &[Node] = &[ + FILTER_SOURCE, + FILTER_INGRESS, + FILTER_ORIGIN_AS, + FILTER_COMMUNITY, +]; static BGP_SUMMARY: &[Node] = &[ Node { @@ -875,6 +930,91 @@ static BGP_FLOWSPEC: &[Node] = &[ FILTER_INGRESS, ]; +// Acyclic filter stages keep command enumeration finite. Filters may be +// omitted, but when combined follow tenant -> type -> scope -> output. +macro_rules! evpn_arg { + ($word:expr, $help:expr, $kind:ident, $next:expr) => { + Node { + kw: Kw::Lit($word), + help: $help, + set: None, + run: None, + children: &[Node { + kw: Kw::Arg(ArgKind::$kind), + help: $help, + set: None, + run: Some(commands::evpn::routes), + children: $next, + }], + } + }; +} +const EVPN_DETAIL: Node = Node { + kw: Kw::Lit("detail"), + help: "All EVPN fields and path attributes", + set: Some(Flag::Detail), + run: Some(commands::evpn::routes), + children: &[], +}; +const EVPN_WITHDRAWN: Node = Node { + kw: Kw::Lit("include-withdrawn"), + help: "Include retained withdrawn routes", + set: Some(Flag::IncludeWithdrawn), + run: Some(commands::evpn::routes), + children: &[EVPN_DETAIL], +}; +static EVPN_OUTPUT: &[Node] = &[EVPN_WITHDRAWN, EVPN_DETAIL]; +const EVPN_VNI: Node = + evpn_arg!("vni", "Match either VXLAN VNI field", Vni, EVPN_OUTPUT); +const EVPN_PREFIX: Node = evpn_arg!( + "prefix", + "Exact MAC/IP host or IP prefix", + Prefix, + EVPN_OUTPUT +); +const EVPN_INGRESS: Node = evpn_arg!( + "ingress", + "Exact stored ingress or ADD-PATH child", + IngressId, + EVPN_OUTPUT +); +static EVPN_SCOPE: &[Node] = &[ + EVPN_VNI, + EVPN_PREFIX, + EVPN_INGRESS, + EVPN_WITHDRAWN, + EVPN_DETAIL, +]; +const EVPN_TYPE: Node = evpn_arg!( + "route-type", + "EVPN route type (2 MAC/IP, 5 IP prefix)", + RouteType, + EVPN_SCOPE +); +static EVPN_TYPE_SCOPE: &[Node] = &[ + EVPN_TYPE, + EVPN_VNI, + EVPN_PREFIX, + EVPN_INGRESS, + EVPN_WITHDRAWN, + EVPN_DETAIL, +]; +static EVPN_TENANT: &[Node] = &[ + evpn_arg!("rd", "Route distinguisher", Rd, EVPN_TYPE_SCOPE), + evpn_arg!( + "route-target", + "Advertised tenant route target", + RouteTarget, + EVPN_TYPE_SCOPE + ), + EVPN_TYPE, + EVPN_VNI, + EVPN_PREFIX, + EVPN_INGRESS, + EVPN_WITHDRAWN, + EVPN_DETAIL, +]; + static SHOW_BMP: &[Node] = &[ leaf!("routers", "Monitored routers", commands::bmp::routers), lit!("router", "One monitored router", SHOW_BMP_ROUTER), diff --git a/test-data/cli/evpn-routes.json b/test-data/cli/evpn-routes.json new file mode 100644 index 0000000..03c57e3 --- /dev/null +++ b/test-data/cli/evpn-routes.json @@ -0,0 +1,5 @@ +{"data":[ + {"route":{"ingress_id":101,"active":true,"ltime":1,"nlri":{"route_type":2,"rd":"65000:10","ethernet_tag":0,"esi":"00000000000000000000","mac":"00:11:22:33:44:55","prefix":"10.0.0.1/32","gateway":null,"labels":[10010,50000],"raw":[]},"attributes":{}},"overlay":{"route_targets":["65000:10","65000:50000"],"router_mac":"aa:bb:cc:dd:ee:ff","next_hop":"192.0.2.1"},"source_ingress_id":7,"path_id":11}, + {"route":{"ingress_id":102,"active":false,"ltime":2,"nlri":{"route_type":5,"rd":"192.0.2.2:20","ethernet_tag":0,"esi":"00000000000000000000","mac":null,"prefix":"2001:db8::/64","gateway":"::","labels":[50020],"raw":[]},"attributes":{}},"overlay":{"route_targets":["65000:20"],"router_mac":null,"next_hop":"2001:db8::2"},"source_ingress_id":8,"path_id":null}, + {"route":{"ingress_id":103,"active":true,"ltime":3,"nlri":{"route_type":3,"rd":"65000:30","mac":null,"prefix":null,"labels":[],"raw":[]},"attributes":{}},"overlay":{"route_targets":[],"router_mac":null,"next_hop":null},"source_ingress_id":9,"path_id":null} +]} From b601bd33d7873654c0bf804ae4b54cba1a029ca5 Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:26:46 -0500 Subject: [PATCH 03/10] Box EVPN NLRI to keep route and payload storage compact --- src/payload.rs | 85 ++++++++++++++++++++++++- src/roto_runtime/types.rs | 8 ++- src/units/bmp_tcp_out/client_handler.rs | 5 +- src/units/rib_unit/http_ng.rs | 2 +- src/units/rib_unit/rib.rs | 4 +- 5 files changed, 96 insertions(+), 8 deletions(-) diff --git a/src/payload.rs b/src/payload.rs index 1f71b0f..098541a 100644 --- a/src/payload.rs +++ b/src/payload.rs @@ -60,10 +60,30 @@ pub enum RotondaRoute { routecore::bgp::nlri::afisafi::Ipv6FlowSpecNlri, RotondaPaMap, ), - L2VpnEvpn(crate::units::rib_unit::evpn::EvpnNlri, RotondaPaMap), + // Keep the large EVPN representation out of every route and queue slot. + // Boxing adds one allocation only for EVPN; other families stay compact. + L2VpnEvpn(Box, RotondaPaMap), // TODO support all routecore AfiSafiTypes } +impl RotondaRoute { + /// Owned EVPN heap storage, excluding shared path attributes. + fn evpn_heap_bytes(&self) -> usize { + match self { + Self::L2VpnEvpn(n, _) => { + std::mem::size_of_val(n.as_ref()) + + n.rd.capacity() + + n.esi.as_ref().map_or(0, String::capacity) + + n.mac.as_ref().map_or(0, String::capacity) + + n.labels.capacity() * std::mem::size_of::() + + n.raw.capacity() + + n.key.capacity() + } + _ => 0, + } + } +} + impl Serialize for RotondaRoute { fn serialize(&self, serializer: S) -> Result where @@ -695,9 +715,15 @@ impl Update { use std::mem::size_of; let base = size_of::(); match self { - Update::Single(_) => base, + Update::Single(payload) => { + base + payload.rx_value.evpn_heap_bytes() + } Update::Bulk(payloads) => { base + payloads.len() * size_of::() + + payloads + .iter() + .map(|p| p.rx_value.evpn_heap_bytes()) + .sum::() } Update::Withdraw(..) => base, Update::WithdrawBulk(items) => { @@ -829,3 +855,58 @@ mod interner_tests { assert_eq!(interner.sweep_shard(usize::MAX), (0, 0)); } } + +#[cfg(all(test, target_pointer_width = "64"))] +mod layout_tests { + #[test] + fn evpn_buffer_accounting_includes_owned_heap() { + use super::*; + let mut raw = vec![5, 34]; + raw.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, 1]); + raw.extend_from_slice(&[0; 14]); + raw.extend_from_slice(&[24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, 10]); + let nlri = + crate::units::rib_unit::evpn::EvpnNlri::parse(&raw).unwrap(); + let heap = std::mem::size_of_val(&nlri) + + nlri.rd.capacity() + + nlri.esi.as_ref().map_or(0, String::capacity) + + nlri.mac.as_ref().map_or(0, String::capacity) + + nlri.labels.capacity() * 4 + + nlri.raw.capacity() + + nlri.key.capacity(); + let payload = Payload::new( + RotondaRoute::L2VpnEvpn( + Box::new(nlri), + RotondaPaMap::empty_path_attributes(), + ), + None, + 1, + RouteStatus::Active, + ); + let single_payload = payload.clone(); + let single_heap = single_payload.rx_value.evpn_heap_bytes(); + let single = Update::Single(single_payload); + assert_eq!( + single.shallow_bytes(), + std::mem::size_of::() + single_heap + ); + let cloned = payload.clone(); + let bulk_heap = cloned.rx_value.evpn_heap_bytes() + heap; + let bulk = Update::Bulk(Box::new(smallvec![cloned, payload])); + assert_eq!( + bulk.shallow_bytes(), + std::mem::size_of::() + + 2 * std::mem::size_of::() + + bulk_heap + ); + } + + #[test] + fn route_and_payload_sizes() { + let route = std::mem::size_of::(); + let payload = std::mem::size_of::(); + eprintln!("RotondaRoute: {route} bytes; Payload: {payload} bytes"); + assert_eq!(route, 64, "EVPN must not inflate common route storage"); + assert_eq!(payload, 96, "EVPN must not inflate payload storage"); + } +} diff --git a/src/roto_runtime/types.rs b/src/roto_runtime/types.rs index da6b210..3517b75 100644 --- a/src/roto_runtime/types.rs +++ b/src/roto_runtime/types.rs @@ -752,7 +752,9 @@ pub(crate) fn convert_nlri>( n.compose(&mut raw).map_err(|_| ())?; ( RotondaRoute::L2VpnEvpn( - crate::units::rib_unit::evpn::EvpnNlri::parse(&raw)?, + Box::new(crate::units::rib_unit::evpn::EvpnNlri::parse( + &raw, + )?), pamap, ), None, @@ -764,7 +766,9 @@ pub(crate) fn convert_nlri>( n.compose(&mut raw).map_err(|_| ())?; ( RotondaRoute::L2VpnEvpn( - crate::units::rib_unit::evpn::EvpnNlri::parse(&raw[4..])?, + Box::new(crate::units::rib_unit::evpn::EvpnNlri::parse( + &raw[4..], + )?), pamap, ), Some(n.path_id()), diff --git a/src/units/bmp_tcp_out/client_handler.rs b/src/units/bmp_tcp_out/client_handler.rs index 7c14e46..b2aaade 100644 --- a/src/units/bmp_tcp_out/client_handler.rs +++ b/src/units/bmp_tcp_out/client_handler.rs @@ -1420,7 +1420,10 @@ mod tests { ), ); let route = RotondaRoute::L2VpnEvpn( - crate::units::rib_unit::evpn::EvpnNlri::parse(&raw).unwrap(), + Box::new( + crate::units::rib_unit::evpn::EvpnNlri::parse(&raw) + .unwrap(), + ), attrs, ); rib.insert(&route, RouteStatus::Active, 0, child, true, false) diff --git a/src/units/rib_unit/http_ng.rs b/src/units/rib_unit/http_ng.rs index af699d8..042f431 100644 --- a/src/units/rib_unit/http_ng.rs +++ b/src/units/rib_unit/http_ng.rs @@ -1083,7 +1083,7 @@ mod evpn_tests { ), ); let route = RotondaRoute::L2VpnEvpn( - super::super::evpn::EvpnNlri::parse(&raw).unwrap(), + Box::new(super::super::evpn::EvpnNlri::parse(&raw).unwrap()), attributes, ); rib.insert(&route, RouteStatus::Active, 0, 10, true, false) diff --git a/src/units/rib_unit/rib.rs b/src/units/rib_unit/rib.rs index 55f8e69..f301f7f 100644 --- a/src/units/rib_unit/rib.rs +++ b/src/units/rib_unit/rib.rs @@ -579,7 +579,7 @@ impl Rib { ingress_id, ltime, active, - nlri: n.clone(), + nlri: n.as_ref().clone(), attributes: if deduplicate_path_attributes { attributes .dedup_with(&self.path_attribute_interner) @@ -4230,7 +4230,7 @@ mod tests { raw.extend_from_slice(&[0; 14]); raw.extend_from_slice(&[24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, vni]); RotondaRoute::L2VpnEvpn( - EvpnNlri::parse(&raw).unwrap(), + Box::new(EvpnNlri::parse(&raw).unwrap()), RotondaPaMap::empty_path_attributes(), ) } From f29cc1f222bed3d54006f4188af2400ceaeb6b8b Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:33:49 -0500 Subject: [PATCH 04/10] Require explicit runtime enablement of EVPN monitoring --- docs/configuration.md | 10 +++ docs/evpn.md | 24 +++++++- etc/netom.conf | 5 ++ src/config.rs | 38 +++++++++++- src/roto_runtime/types.rs | 81 +++++++++++++++++++++++-- src/units/bgp_tcp_in/peer_config.rs | 20 +++++- src/units/bmp_tcp_out/bmp_builder.rs | 13 ++-- src/units/bmp_tcp_out/client_handler.rs | 7 +-- src/units/rib_unit/http_ng.rs | 5 ++ 9 files changed, 180 insertions(+), 23 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index 4cc0927..9e1b575 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -43,6 +43,16 @@ Component-specific settings belong under that component's table, such as placing a global setting at the end of the file would put it inside the last component instead. +## Enable EVPN monitoring + +The global runtime setting `enable_evpn` defaults to `false`. Set +`enable_evpn = true` above the first component table and restart Netom to +opt in to EVPN ingestion, BGP capabilities, and API queries. Reloads cannot +change this setting. A peer's `protocols = ["L2VpnEvpn"]` alone does not enable +EVPN. The opt-in avoids accidentally retaining additional EVPN routing state +and incurring its query memory and CPU costs; see [EVPN monitoring](evpn.md) +for behavior and configuration details. + ## Example files The [annotated configuration](../etc/netom.conf) is maintained with the diff --git a/docs/evpn.md b/docs/evpn.md index b2c306d..5c7eb6a 100644 --- a/docs/evpn.md +++ b/docs/evpn.md @@ -1,6 +1,28 @@ # EVPN monitoring -Netom collects L2VPN EVPN (AFI 25, SAFI 70) from BGP and BMP, including +EVPN monitoring is **disabled by default**. Enable it explicitly in the +TOML global settings, above all `[units.*]` and `[targets.*]` tables: + +```toml +enable_evpn = true +``` + +Restart Netom after changing this flag; a reload that changes it is rejected. +This keeps negotiated BGP capabilities and retained routing state consistent. +When disabled, Netom omits EVPN from BGP MP and ADD-PATH capabilities, drops +EVPN announcements and withdrawals from BGP/BMP route conversion through the +unsupported-NLRI accounting path, and returns HTTP 503 with an enablement +message for EVPN queries. Other families continue to work. Verbatim BMP +forwarding is independent of monitoring and can still carry EVPN updates. + +EVPN retains a separate collection of routes and attributes, and queries incur +additional copying, allocation, and serialization costs. These costs grow with +routes, paths, tenants, and concurrent queries. Requiring an explicit opt-in +keeps deployments that only monitor other families from accidentally taking +on this state and query load. Enabling EVPN is not a memory or query-cost limit; +size the deployment for its routing state and workload. + +With the flag enabled, Netom collects L2VPN EVPN (AFI 25, SAFI 70) from BGP and BMP, including ADD-PATH sessions. Configure `L2VpnEvpn` in a BGP peer's `protocols` list; BMP peers use the capabilities in their exported Peer Up messages. diff --git a/etc/netom.conf b/etc/netom.conf index 5a3696c..c094c9d 100644 --- a/etc/netom.conf +++ b/etc/netom.conf @@ -20,6 +20,11 @@ log_target = "stderr" # "stderr", "file" or "syslog" http_listen = ["[::]:8080"] +# EVPN monitoring is an explicit opt-in because its retained routes and API +# queries add memory and CPU costs. Defaults to false; changes need a restart. +# Also configure L2VpnEvpn in BGP peer protocols when using BGP input. +# enable_evpn = true + ### 2. Component Definitions diff --git a/src/config.rs b/src/config.rs index 0b10688..a91a608 100644 --- a/src/config.rs +++ b/src/config.rs @@ -14,7 +14,7 @@ use serde_with::serde_as; use std::cell::RefCell; use std::net::SocketAddr; use std::path::{Path, PathBuf}; -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; use std::{borrow, error, fmt, fs, io, ops}; use toml::{Spanned, Value}; @@ -22,6 +22,14 @@ use toml::{Spanned, Value}; const ARG_CONFIG: &str = "config"; +// Process-wide, immutable after successful startup. Changing this on reload +// would leave negotiated sessions and retained RIB state inconsistent. +static EVPN_ENABLED: OnceLock = OnceLock::new(); + +pub(crate) fn evpn_enabled() -> bool { + EVPN_ENABLED.get().copied().unwrap_or(false) +} + //------------ Config -------------------------------------------------------- /// The complete Netom configuration. @@ -38,6 +46,11 @@ const ARG_CONFIG: &str = "config"; #[derive(Deserialize)] #[serde(deny_unknown_fields)] pub struct Config { + /// Opt in to EVPN monitoring and its additional retained-state/query costs. + /// Defaults to false; changes require a daemon restart. + #[serde(default)] + pub enable_evpn: bool, + /// Location of the .roto script containing all user defined filters. pub roto_script: Option, @@ -146,7 +159,16 @@ impl Config { trace!("{}", config_file.to_string()); } + if EVPN_ENABLED + .get() + .is_some_and(|enabled| *enabled != self.enable_evpn) + { + error!("Changing enable_evpn requires a daemon restart"); + return Err(Terminate::error()); + } manager.prepare(&self, &config_file)?; + // Units are started only after finalise returns successfully. + let _ = EVPN_ENABLED.set(self.enable_evpn); // Pass the config file path, as well as the processed config, back to // the caller so that they can monitor it for changes while the @@ -560,6 +582,20 @@ impl AsRef for ConfigPath { mod tests { use super::*; + #[test] + fn evpn_requires_explicit_enablement() { + for (setting, expected) in [ + ("", false), + ("enable_evpn = false\n", false), + ("enable_evpn = true\n", true), + ] { + let text = format!("{setting}[units]\n[targets]\n"); + let config = + Config::from_bytes(text.as_bytes(), None::<&Path>).unwrap(); + assert_eq!(config.enable_evpn, expected); + } + } + fn mk(toml: &str) -> ConfigFile { ConfigFile::new(toml.as_bytes().to_vec(), Source::default()).unwrap() } diff --git a/src/roto_runtime/types.rs b/src/roto_runtime/types.rs index 3517b75..6b84e25 100644 --- a/src/roto_runtime/types.rs +++ b/src/roto_runtime/types.rs @@ -487,13 +487,11 @@ impl OutputStreamMessage { //--- Unsupported NLRI drop accounting --------------------------------------- -/// How often to emit a rolled-up summary of NLRI dropped because their type has -/// no [`RotondaRoute`] representation. +/// How often to emit a rolled-up summary of unsupported or disabled NLRI. const UNSUPPORTED_NLRI_SUMMARY_INTERVAL: Duration = Duration::from_secs(60); -/// Process-global accounting for NLRI dropped because their [`NlriType`] has no -/// [`RotondaRoute`] representation: genuinely unsupported families (MPLS-VPN, -/// EVPN, RouteTarget, ...). ADD-PATH NLRI of the stored families +/// Process-global accounting for unsupported families (MPLS-VPN, RouteTarget, +/// ...) and EVPN when disabled by configuration. ADD-PATH NLRI of enabled families /// (unicast/multicast/flowspec) are *not* counted here — [`convert_nlri`] /// strips their path id into a per-path child ingress and stores them. Such /// NLRI parse fine in routecore but are dropped before any RIB in the @@ -676,14 +674,27 @@ fn note_unsupported_nlri(nlri_type: NlriType) { /// the stripped [`PathId`] is returned alongside so the caller can resolve a /// per-(session, path_id) child ingress to store the route under. Plain /// variants return `None`. NLRI types with no `RotondaRoute` representation -/// (MPLS/VPN/EVPN/VPLS/RouteTarget) are counted and dropped as before. +/// (MPLS/VPN/VPLS/RouteTarget) are counted and dropped as before. EVPN +/// requires the global `enable_evpn` opt-in. pub(crate) fn convert_nlri>( nlri: Nlri, pamap: RotondaPaMap, +) -> Result<(RotondaRoute, Option), ()> { + convert_nlri_with_evpn(nlri, pamap, crate::config::evpn_enabled()) +} + +fn convert_nlri_with_evpn>( + nlri: Nlri, + pamap: RotondaPaMap, + enable_evpn: bool, ) -> Result<(RotondaRoute, Option), ()> { use routecore::bgp::nlri::afisafi::Addpath; let res = match nlri { + Nlri::L2VpnEvpn(..) | Nlri::L2VpnEvpnAddpath(..) if !enable_evpn => { + note_unsupported_nlri(nlri.nlri_type()); + return Err(()); + } Nlri::Ipv4Unicast(n) => (RotondaRoute::Ipv4Unicast(n, pamap), None), Nlri::Ipv4Multicast(n) => { (RotondaRoute::Ipv4Multicast(n, pamap), None) @@ -824,6 +835,26 @@ pub(crate) fn explode_announcements( Ok(res) } +// Exporter round-trip tests explicitly decode EVPN without changing the +// process-wide setting (which would race with default-off ingestion tests). +#[cfg(test)] +pub(crate) fn decode_evpn_test_update( + update: &UpdateMessage, + withdrawn: bool, +) -> Vec<(RotondaRoute, Option)> { + let attributes = + RotondaPaMap::new(update.path_attributes().unwrap().into()); + let decode = |nlri: Result, _>| { + convert_nlri_with_evpn(nlri.unwrap(), attributes.clone(), true) + .unwrap() + }; + if withdrawn { + update.withdrawals().unwrap().map(decode).collect() + } else { + update.announcements().unwrap().map(decode).collect() + } +} + /// Explode a BGP UPDATE's withdrawals into storable routes; see /// [`explode_announcements`] for the path-id component. pub(crate) fn explode_withdrawals( @@ -870,6 +901,44 @@ impl PeerId { mod tests { use super::*; + #[test] + fn evpn_conversion_requires_opt_in_including_addpath() { + use octseq::Parser; + use routecore::bgp::nlri::afisafi::{ + L2VpnEvpnAddpathNlri, L2VpnEvpnNlri, NlriParse, + }; + + // Opaque route type with an RD: valid, and independent of forwarding fields. + let raw = [99, 8, 0, 0, 0, 1, 0, 0, 0, 2]; + let mut addpath_raw = vec![0, 0, 0, 7]; + addpath_raw.extend_from_slice(&raw); + for enabled in [false, true] { + let plain = L2VpnEvpnNlri::<&[u8]>::parse(&mut Parser::from_ref( + &raw.as_slice(), + )) + .unwrap(); + let addpath = L2VpnEvpnAddpathNlri::<&[u8]>::parse( + &mut Parser::from_ref(&addpath_raw.as_slice()), + ) + .unwrap(); + let result = convert_nlri_with_evpn( + Nlri::L2VpnEvpn(plain), + RotondaPaMap::default(), + enabled, + ); + assert_eq!(result.is_ok(), enabled); + let result = convert_nlri_with_evpn( + Nlri::L2VpnEvpnAddpath(addpath), + RotondaPaMap::default(), + enabled, + ); + assert_eq!(result.is_ok(), enabled); + if enabled { + assert_eq!(result.unwrap().1, Some(PathId(7))); + } + } + } + #[test] fn unsupported_nlri_counter_increments_and_renders() { let m = UnsupportedNlriMetrics::default(); diff --git a/src/units/bgp_tcp_in/peer_config.rs b/src/units/bgp_tcp_in/peer_config.rs index c2464c6..d65fa8c 100644 --- a/src/units/bgp_tcp_in/peer_config.rs +++ b/src/units/bgp_tcp_in/peer_config.rs @@ -333,11 +333,27 @@ impl BgpConfig for CombinedConfig { } fn protocols(&self) -> Vec { - self.peer_config.protocols.clone() + self.peer_config + .protocols + .iter() + .copied() + .filter(|family| { + *family != AfiSafiType::L2VpnEvpn + || crate::config::evpn_enabled() + }) + .collect() } fn addpath(&self) -> Vec { - self.peer_config.addpath.clone() + self.peer_config + .addpath + .iter() + .copied() + .filter(|family| { + *family != AfiSafiType::L2VpnEvpn + || crate::config::evpn_enabled() + }) + .collect() } fn extended_messages(&self) -> bool { diff --git a/src/units/bmp_tcp_out/bmp_builder.rs b/src/units/bmp_tcp_out/bmp_builder.rs index 99d26b8..f296f08 100644 --- a/src/units/bmp_tcp_out/bmp_builder.rs +++ b/src/units/bmp_tcp_out/bmp_builder.rs @@ -3183,9 +3183,6 @@ mod tests { #[test] fn evpn_bmp_roundtrip_plain_and_addpath() { - use crate::roto_runtime::types::{ - explode_announcements, explode_withdrawals, - }; use routecore::bgp::message::{SessionConfig, UpdateMessage}; let peer = agg_test_peer(); let mut raw = vec![5, 34]; @@ -3214,12 +3211,10 @@ mod tests { &bmp[BMP_COMMON_HEADER_LEN + BMP_PER_PEER_HEADER_LEN..], ); let update = UpdateMessage::from_octets(bgp, &sc).unwrap(); - let routes = if withdrawn { - explode_withdrawals(&update) - } else { - explode_announcements(&update) - } - .unwrap(); + let routes = + crate::roto_runtime::types::decode_evpn_test_update( + &update, withdrawn, + ); assert_eq!(routes.len(), 1); assert_eq!(routes[0].1.map(|p| p.0), pid); match &routes[0].0 { diff --git a/src/units/bmp_tcp_out/client_handler.rs b/src/units/bmp_tcp_out/client_handler.rs index b2aaade..5848fd9 100644 --- a/src/units/bmp_tcp_out/client_handler.rs +++ b/src/units/bmp_tcp_out/client_handler.rs @@ -1477,10 +1477,9 @@ mod tests { ) .unwrap(); let routes = - crate::roto_runtime::types::explode_announcements( - &update, - ) - .unwrap(); + crate::roto_runtime::types::decode_evpn_test_update( + &update, false, + ); assert_eq!(routes.len(), 1); assert_eq!(routes[0].1.unwrap().0, 11); } diff --git a/src/units/rib_unit/http_ng.rs b/src/units/rib_unit/http_ng.rs index 042f431..1aa5705 100644 --- a/src/units/rib_unit/http_ng.rs +++ b/src/units/rib_unit/http_ng.rs @@ -992,6 +992,11 @@ async fn search_evpn( Query(filter): Query, state: State, ) -> Result { + if !crate::config::evpn_enabled() { + return Err(ApiError::ServiceUnavailable( + "EVPN support is disabled; set enable_evpn = true and restart Netom".into(), + )); + } let rib = load_rib(&state)?; let permit = super::rib::DumpGuard::try_enter().ok_or_else(|| { ApiError::ServiceUnavailable("too many concurrent RIB queries".into()) From 147a88fe705c2cc536d62f05a050ec89ee7ef71f Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:37:03 -0500 Subject: [PATCH 05/10] Use filtered shared EVPN snapshots and guard response serialization --- src/units/bmp_tcp_out/client_handler.rs | 4 +- src/units/rib_unit/http_ng.rs | 64 +++++++++------- src/units/rib_unit/rib.rs | 97 +++++++++++++++++++------ 3 files changed, 114 insertions(+), 51 deletions(-) diff --git a/src/units/bmp_tcp_out/client_handler.rs b/src/units/bmp_tcp_out/client_handler.rs index 5848fd9..cb18e9d 100644 --- a/src/units/bmp_tcp_out/client_handler.rs +++ b/src/units/bmp_tcp_out/client_handler.rs @@ -479,9 +479,7 @@ pub async fn perform_initial_dump( }; // EVPN uses RD-scoped keys instead of the IP prefix tree. if !client_gone { - for record in - rib_for_walk.evpn_records().into_iter().filter(|r| r.active) - { + for record in rib_for_walk.evpn_records_matching(|r| r.active) { let source = ingress_register_for_walk.get(record.ingress_id); let (ingress_id, path_id) = match source { Some(ref info) diff --git a/src/units/rib_unit/http_ng.rs b/src/units/rib_unit/http_ng.rs index 1aa5705..f30cca0 100644 --- a/src/units/rib_unit/http_ng.rs +++ b/src/units/rib_unit/http_ng.rs @@ -1001,40 +1001,48 @@ async fn search_evpn( let permit = super::rib::DumpGuard::try_enter().ok_or_else(|| { ApiError::ServiceUnavailable("too many concurrent RIB queries".into()) })?; - let data = tokio::task::spawn_blocking(move || { + let body = tokio::task::spawn_blocking(move || { let _permit = permit; - evpn_rows(&rib, &filter) + #[derive(serde::Serialize)] + struct Response { + data: Vec, + } + serde_json::to_vec(&Response { + data: evpn_rows(&rib, &filter), + }) }) .await + .map_err(|e| ApiError::InternalServerError(e.to_string()))? .map_err(|e| ApiError::InternalServerError(e.to_string()))?; - let body = serde_json::to_vec(&serde_json::json!({"data": data})) - .map_err(|e| ApiError::InternalServerError(e.to_string()))?; Ok( ([("content-type", OutputFormat::Json.content_type())], body) .into_response(), ) } -fn evpn_rows( - rib: &super::rib::Rib, - filter: &EvpnFilter, -) -> Vec { - let mut records = rib.evpn_records(); +#[derive(serde::Serialize)] +struct EvpnRow { + route: std::sync::Arc, + overlay: super::evpn::EvpnAttributes, + source_ingress_id: IngressId, + path_id: Option, +} + +fn evpn_rows(rib: &super::rib::Rib, filter: &EvpnFilter) -> Vec { + let mut records = rib.evpn_records_matching(|record| { + let n = &record.nlri; + (record.active || filter.include_withdrawn) + && filter.rd.as_ref().is_none_or(|v| v == &n.rd) + && filter.route_type.is_none_or(|v| v == n.route_type) + && filter.vni.is_none_or(|v| n.labels.contains(&v)) + && filter.prefix.is_none_or(|v| Some(v) == n.prefix) + && filter.ingress_id.is_none_or(|v| v == record.ingress_id) + }); records.sort_by(|a, b| { (&a.nlri.key, a.ingress_id).cmp(&(&b.nlri.key, b.ingress_id)) }); let mut rows = Vec::new(); for record in records { - let n = &record.nlri; - if (!record.active && !filter.include_withdrawn) - || filter.rd.as_ref().is_some_and(|v| v != &n.rd) - || filter.route_type.is_some_and(|v| v != n.route_type) - || filter.vni.is_some_and(|v| !n.labels.contains(&v)) - || filter.prefix.is_some_and(|v| Some(v) != n.prefix) - || filter.ingress_id.is_some_and(|v| v != record.ingress_id) - { - continue; - } let overlay = super::evpn::EvpnAttributes::decode(&record.attributes); if filter .route_target @@ -1047,12 +1055,14 @@ fn evpn_rows( let path_source = source .as_ref() .filter(|s| s.ingress_type == Some(IngressType::BgpPath)); - rows.push(serde_json::json!({ - "route": record, - "overlay": overlay, - "source_ingress_id": path_source.and_then(|s| s.parent_ingress).unwrap_or(record.ingress_id), - "path_id": path_source.and_then(|s| s.path_id), - })); + rows.push(EvpnRow { + source_ingress_id: path_source + .and_then(|s| s.parent_ingress) + .unwrap_or(record.ingress_id), + path_id: path_source.and_then(|s| s.path_id), + route: record, + overlay, + }); } rows } @@ -1101,8 +1111,8 @@ mod evpn_tests { })).unwrap(); let rows = evpn_rows(&rib, &filter); assert_eq!(rows.len(), 1); - assert_eq!(rows[0]["route"]["nlri"]["rd"], "1:1"); - assert_eq!(rows[0]["source_ingress_id"], 10); + assert_eq!(rows[0].route.nlri.rd, "1:1"); + assert_eq!(rows[0].source_ingress_id, 10); let filter = EvpnFilter { vni: Some(2), ..filter diff --git a/src/units/rib_unit/rib.rs b/src/units/rib_unit/rib.rs index f301f7f..336a614 100644 --- a/src/units/rib_unit/rib.rs +++ b/src/units/rib_unit/rib.rs @@ -160,7 +160,9 @@ type RotoHttpFilter = roto::TypedFunc< #[derive(Clone)] pub struct Rib { - evpn: Arc, IngressId), super::evpn::EvpnRecord>>>, + evpn: Arc< + Mutex, IngressId), Arc>>, + >, unicast: Arc>, multicast: Arc>, /// FlowSpec rules (SAFI 133, v4+v6 in the one dual-family store), keyed @@ -575,7 +577,7 @@ impl Rib { if active { records.insert( key, - super::evpn::EvpnRecord { + Arc::new(super::evpn::EvpnRecord { ingress_id, ltime, active, @@ -586,10 +588,11 @@ impl Rib { } else { attributes.clone() }, - }, + }), ); } else if retain_withdrawn_attributes { if let Some(record) = records.get_mut(&key) { + let record = Arc::make_mut(record); record.active = false; record.ltime = ltime; } @@ -1334,7 +1337,9 @@ impl Rib { .collect(); records.retain(|(_, ingress), record| { if evpn_ingresses.contains(ingress) { - record.active = false; + if retain_withdrawn_attributes && record.active { + Arc::make_mut(record).active = false; + } retain_withdrawn_attributes } else { true @@ -1470,23 +1475,32 @@ impl Rib { self.recount_flowspec_rules(); } - /// Physically remove every record for these ingress ids from the RIB, - /// reclaiming their memory — as opposed to `withdraw_for_ingresses`, which - /// only marks them withdrawn (a status bit) and keeps the records around. - /// - /// This is used both for ids that are gone for good and to discard a - /// retained previous session before reusing an id. Synthesized BMP peers - /// that mint a fresh ingress id every session must take this path so - /// mark-withdraw does not leak one record slot per announced prefix. - pub fn evpn_records(&self) -> Vec { + /// Capture a coherent snapshot of matching shared records. The predicate + /// runs under the writer mutex and must be inexpensive; decoding, sorting, + /// serialization and I/O belong after this method returns. Only matching + /// records allocate snapshot slots (one pointer each), with no deep copies. + /// Later mutations use copy-on-write while a snapshot retains a record. + pub fn evpn_records_matching( + &self, + matches: impl Fn(&super::evpn::EvpnRecord) -> bool, + ) -> Vec> { self.evpn .lock() .unwrap_or_else(|e| e.into_inner()) .values() + .filter(|record| matches(record)) .cloned() .collect() } + /// Physically remove every record for these ingress ids from the RIB, + /// reclaiming their memory — as opposed to `withdraw_for_ingresses`, which + /// only marks them withdrawn (a status bit) and keeps the records around. + /// + /// This is used both for ids that are gone for good and to discard a + /// retained previous session before reusing an id. Synthesized BMP peers + /// that mint a fresh ingress id every session must take this path so + /// mark-withdraw does not leak one record slot per announced prefix. pub fn remove_for_ingresses(&self, ids: &[IngressId]) { self.evpn .lock() @@ -4246,29 +4260,70 @@ mod tests { ) .unwrap(); } - assert_eq!(rib.evpn_records().len(), 3); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 3); + let snapshot = rib.evpn_records_matching(|r| { + r.ingress_id == 1 && r.nlri.rd == "1:1" + }); + assert_eq!(snapshot.len(), 1); + let same = rib.evpn_records_matching(|r| { + r.ingress_id == 1 && r.nlri.rd == "1:1" + }); + assert!(Arc::ptr_eq(&snapshot[0], &same[0])); + drop(same); + rib.withdraw_for_ingress(1, Some(AfiSafiType::L2VpnEvpn), true); + assert!(snapshot[0].active); + assert_eq!(rib.evpn_records_matching(|r| r.active).len(), 1); + rib.insert(&route(2, 10), RouteStatus::Active, 0, 1, true, false) + .unwrap(); rib.insert(&route(1, 20), RouteStatus::Active, 1, 1, true, false) .unwrap(); - assert_eq!(rib.evpn_records().len(), 3); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 3); + let updated_snapshot = rib.evpn_records_matching(|r| { + r.ingress_id == 1 && r.nlri.rd == "1:1" + }); rib.insert(&route(1, 0), RouteStatus::Withdrawn, 2, 1, true, false) .unwrap(); - let rows = rib.evpn_records(); + assert!(updated_snapshot[0].active); + assert_eq!(updated_snapshot[0].ltime, 1); + assert!(snapshot[0].active); + assert_eq!(snapshot[0].nlri.labels, vec![10]); + let withdrawn_snapshot = rib.evpn_records_matching(|r| !r.active); + assert_eq!(withdrawn_snapshot.len(), 1); + let rows = rib.evpn_records_matching(|_| true); assert_eq!(rows.iter().filter(|r| r.active).count(), 2); assert_eq!( rows.iter().find(|r| !r.active).unwrap().nlri.labels, vec![20] ); rib.withdraw_for_ingress(1, Some(AfiSafiType::Ipv4Unicast), false); - assert_eq!(rib.evpn_records().len(), 3); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 3); rib.withdraw_for_ingress(1, Some(AfiSafiType::L2VpnEvpn), false); - assert_eq!(rib.evpn_records().len(), 1); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 1); rib.insert(&route(2, 30), RouteStatus::Active, 3, 1, false, false) .unwrap(); - assert_eq!(rib.evpn_records().len(), 2); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 2); rib.remove_for_ingresses(&[2]); - assert_eq!(rib.evpn_records().len(), 1); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 1); rib.withdraw_for_ingress(1, None, false); - assert!(rib.evpn_records().is_empty()); + assert!(rib.evpn_records_matching(|_| true).is_empty()); + assert!(!withdrawn_snapshot[0].active); + assert_eq!(withdrawn_snapshot[0].nlri.labels, vec![20]); + assert!(snapshot[0].active); + } + + #[test] + fn evpn_shared_record_sizes() { + use super::super::evpn::EvpnRecord; + println!("EvpnRecord={} inline map value -> shared map value={}, Arc counters={} bytes per allocation; RotondaRoute={}, Payload={}", + std::mem::size_of::(), + std::mem::size_of::>(), + 2 * std::mem::size_of::(), + std::mem::size_of::(), + std::mem::size_of::()); + assert_eq!( + std::mem::size_of::>(), + std::mem::size_of::() + ); } fn test_rib() -> Rib { From f56f8e753ca68940409fbc7cda2fb8d7a9017199 Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:41:07 -0500 Subject: [PATCH 06/10] Optimize batched EVPN ingress cleanup --- src/units/rib_unit/rib.rs | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/src/units/rib_unit/rib.rs b/src/units/rib_unit/rib.rs index 336a614..3e463c1 100644 --- a/src/units/rib_unit/rib.rs +++ b/src/units/rib_unit/rib.rs @@ -1502,14 +1502,18 @@ impl Rib { /// that mint a fresh ingress id every session must take this path so /// mark-withdraw does not leak one record slot per announced prefix. pub fn remove_for_ingresses(&self, ids: &[IngressId]) { - self.evpn - .lock() - .unwrap_or_else(|e| e.into_inner()) - .retain(|(_, ingress), _| !ids.contains(ingress)); if ids.is_empty() { return; } + // Build membership outside the EVPN lock: O(I) setup, then O(R) + // expected scan cost rather than searching all I ids for each record. + let ingress_ids: HashSet = ids.iter().copied().collect(); + self.evpn + .lock() + .unwrap_or_else(|e| e.into_inner()) + .retain(|(_, ingress), _| !ingress_ids.contains(ingress)); + // `remove_mui` clears the per-store `withdrawn_muis_bmin` bitmap (via // mark_mui_as_active), the same CAS that livelocks under concurrent // writers, so it must hold `withdraw_lock` just like the mark path. @@ -4261,6 +4265,8 @@ mod tests { .unwrap(); } assert_eq!(rib.evpn_records_matching(|_| true).len(), 3); + rib.remove_for_ingresses(&[]); + assert_eq!(rib.evpn_records_matching(|_| true).len(), 3); let snapshot = rib.evpn_records_matching(|r| { r.ingress_id == 1 && r.nlri.rd == "1:1" }); @@ -4302,7 +4308,7 @@ mod tests { rib.insert(&route(2, 30), RouteStatus::Active, 3, 1, false, false) .unwrap(); assert_eq!(rib.evpn_records_matching(|_| true).len(), 2); - rib.remove_for_ingresses(&[2]); + rib.remove_for_ingresses(&[2, 999, 2]); assert_eq!(rib.evpn_records_matching(|_| true).len(), 1); rib.withdraw_for_ingress(1, None, false); assert!(rib.evpn_records_matching(|_| true).is_empty()); From ccd5a342d278e8abe4400f75c34ed6ec284c1f2d Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:43:14 -0500 Subject: [PATCH 07/10] Fix EVPN withdrawal metrics using prior record existence --- src/units/rib_unit/rib.rs | 51 ++++++++++++++++++++++++++++++++++++++- 1 file changed, 50 insertions(+), 1 deletion(-) diff --git a/src/units/rib_unit/rib.rs b/src/units/rib_unit/rib.rs index 3e463c1..7a02837 100644 --- a/src/units/rib_unit/rib.rs +++ b/src/units/rib_unit/rib.rs @@ -603,7 +603,9 @@ impl Rib { cas_count: 0, prefix_new: !existed && active, mui_new: !existed && active, - mui_count: usize::from(active), + // Withdrawals report prior existence so downstream + // metrics can distinguish known and unannounced routes. + mui_count: usize::from(active || existed), }) } @@ -4317,6 +4319,53 @@ mod tests { assert!(snapshot[0].active); } + #[test] + fn evpn_withdrawal_reports_prior_existence() { + use super::super::evpn::EvpnNlri; + + let mut raw = vec![5, 34]; + raw.extend_from_slice(&[0, 0, 0, 1, 0, 0, 0, 1]); + raw.extend_from_slice(&[0; 14]); + raw.extend_from_slice(&[24, 10, 0, 0, 0, 0, 0, 0, 0, 0, 0, 10]); + let route = RotondaRoute::L2VpnEvpn( + Box::new(EvpnNlri::parse(&raw).unwrap()), + RotondaPaMap::empty_path_attributes(), + ); + + for retain in [false, true] { + let rib = test_rib(); + let insert = |status, ingress| { + rib.insert(&route, status, 0, ingress, retain, false) + .unwrap() + }; + let unknown = insert(RouteStatus::Withdrawn, 1); + assert_eq!(unknown.mui_count, 0); + assert!(!unknown.prefix_new && !unknown.mui_new); + + let announced = insert(RouteStatus::Active, 1); + assert_eq!(announced.mui_count, 1); + assert!(announced.prefix_new && announced.mui_new); + + // The same NLRI from another ingress is not a matching record. + assert_eq!(insert(RouteStatus::Withdrawn, 2).mui_count, 0); + let withdrawn = insert(RouteStatus::Withdrawn, 1); + assert_eq!(withdrawn.mui_count, 1); + assert!(!withdrawn.prefix_new && !withdrawn.mui_new); + assert!(rib.evpn_records_matching(|r| r.active).is_empty()); + + // Retained withdrawn records still count as previously known. + assert_eq!( + insert(RouteStatus::Withdrawn, 1).mui_count, + usize::from(retain) + ); + let reannounced = insert(RouteStatus::Active, 1); + assert_eq!(reannounced.mui_count, 1); + assert_eq!(reannounced.prefix_new, !retain); + assert_eq!(reannounced.mui_new, !retain); + assert_eq!(rib.evpn_records_matching(|r| r.active).len(), 1); + } + } + #[test] fn evpn_shared_record_sizes() { use super::super::evpn::EvpnRecord; From 3455121d81aca0ce02e8da0caad9d4ab964e2ae5 Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:48:20 -0500 Subject: [PATCH 08/10] Document Netom engineering principles and cBMP proposal --- AGENTS.md | 91 +++++++++++++++++ ideas/cBMP.md | 272 ++++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 363 insertions(+) create mode 100644 AGENTS.md create mode 100644 ideas/cBMP.md diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..b158b04 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,91 @@ +# Netom engineering principles + +Netom monitors large BGP/BMP routing tables and many sessions. Correctness +includes protocol semantics, lifecycle and metrics, bounded resource use, and +predictable ingestion under concurrent queries. Evaluate changes at realistic +route/path/ingress counts and at 10x scale, beyond small test fixtures. + +## Keep common objects compact + +- Measure `size_of` before/after changes to core route, payload, queue, and + message types, including enclosing types, alignment, and padding. Justify + increases and report their aggregate cost in the PR. +- A large enum variant enlarges every instance. Keep uncommon large data behind + appropriate indirection; account for allocations, reference counts, and copies. +- On 64-bit targets, the boxed EVPN layout preserves `RotondaRoute` at 64 bytes + and `Payload` at 96 bytes. Inline EVPN raised them to 240/272 bytes: + +176 bytes per object, or roughly 176 MB per million objects. +- Queue/buffer accounting must include owned heap allocations and capacities, + not just stack size. Preserve lightweight layout and accounting regression tests. + +## Filter before copying; preserve snapshot consistency + +- Never deep-clone the entire RIB to satisfy a filtered query or solve ownership. + Apply cheap family, active/withdrawn, and identity filters before copying; + exclude withdrawn records unless requested. +- Prefer references or filtered shared immutable snapshots. Captured records/IDs + must remain valid through consumption; use copy-on-write where needed to keep + snapshots coherent across updates, withdrawals, and removal. +- Snapshot predicates under a writer mutex must be inexpensive. Decode, sort, + transform, serialize, and perform network I/O after releasing the lock. +- Account for snapshot lifetime, retained old versions, and concurrent readers. + Add indexes for frequent expensive queries only with a justified memory cost. + +## Bound work and contention + +- Analyze cost using actual cardinalities: routes, paths, ingresses, updates, + peers, and concurrent queries. Include allocations, copies, and lock duration. +- Avoid linear membership searches inside table scans. For batch ingress cleanup, + build a membership set outside the lock: expected O(I + R), rather than O(I × R). + Empty batches return immediately without locking or scanning. +- Keep hot lock sections short; avoid full-table deep copies, expensive decoding, + serialization, and unnecessary allocations while holding them. Preserve an + explicit consistency model when moving work outside locks. + +## Guard the complete query operation + +- Bound concurrent expensive RIB queries. Hold the resource permit through + extraction, transformation, allocation, and serialization, including work in + `spawn_blocking`; releasing it before JSON encoding defeats the limit. +- Move synchronous CPU-heavy work off Tokio workers. `spawn_blocking` alone + does not bound concurrency or peak memory. +- Serialize typed responses directly. Avoid `json!`/JSON value trees that create + another large representation solely for encoding. +- Assess peak memory across retained state, snapshots, response rows, encoded + bodies, and concurrent requests; final response size alone is insufficient. + +## Preserve identity, lifecycle, and metric semantics + +- Route identity must include the correct NLRI, ingress/session provenance, and + ADD-PATH identity. Equal NLRIs can have independent paths; withdrawing one must + leave the others active. +- Derive metrics from prior existence/state and the event, not just resulting + `active`. A known-route withdrawal must not become an unannounced withdrawal; + a different ingress/path is not a matching prior record. +- Define retained-withdrawn behavior explicitly, including repeated withdrawals, + re-announcements, and “new” counters. Test both retention modes. +- Every insertion needs cleanup for withdrawal, family withdrawal, peer/session + teardown, and ingress removal. Marking withdrawn is not physical removal; + expired sessions must not leak one retained slot per prefix/path. +- Test known/unknown withdrawals, updates and re-announcements, independent + ADD-PATH withdrawals, family/peer cleanup, and snapshot stability during mutation. + +## Make optional family support consistent + +- EVPN requires explicit runtime opt-in (`enable_evpn`, default false). Gate + negotiation (including ADD-PATH), ingestion, and APIs consistently; account + for disabled-family drops without retaining their routing state. +- Configuration affecting negotiated sessions and retained state must not change + on reload without a safe transition. EVPN enablement currently requires restart. + Tests should inject settings rather than race on process-global configuration. +- Keep EVPN semantics distinct: RD is not tenant identity; route-target selection + is not VRF import-policy simulation; a numeric label/VNI does not establish + L2 versus L3 role. Preserve Type 2/Type 5 data needed for symmetric IRB analysis. + +## Validate the cost as well as the behavior + +For significant RIB/API changes, report relevant before/after measurements: +bytes per object and owned heap, allocations/copies, query latency, lock hold +time, ingestion throughput, and peak serialization memory. Show how costs grow +with table size and concurrency. Follow the whole operation so a later stage +does not recreate costs eliminated earlier. diff --git a/ideas/cBMP.md b/ideas/cBMP.md new file mode 100644 index 0000000..4b4b545 --- /dev/null +++ b/ideas/cBMP.md @@ -0,0 +1,272 @@ +# cBMP: a canonical, compressed BMP corpus + +Proposal, based on the current Netom source. No implementation or compression +measurements are implied. + +The idea is to normalize the **structure of BMP observations across all incoming +streams**, then use measured frequency and co-occurrence to arrange their encoded +representation for similarity. Compress the result with zstd using a shared, +versioned dictionary. Sorting BGP path attributes is one small part of this: +repeated peer headers, message shapes, UPDATE packaging, next hops, NLRI and +control-message structures also offer opportunities. + +The central hypothesis is that **compression can replace application-level +deduplication for storage and output**. Serialize every observation, including +repeated headers and bodies, in a consistent representation; let zstd encode the +repetition. The initial design needs no peer/body content tables, content-addressed +object store, reference counting or deduplication lookup on the write path. +Normalization exposes redundancy, statistical ordering brings similar bytes +within reach of the compressor, and the shared dictionary supplies learned +patterns across independently compressed chunks. + +The proposed output is a versioned corpus/container format, provisionally cBMP, +with a decoder that can reconstruct BMP streams. It is not a new encoding to send +directly to an ordinary BMP TCP collector. BMP's existing message boundaries and +per-peer headers remain the interoperability contract +([RFC 7854](https://www.rfc-editor.org/rfc/rfc7854.html)). + +## What is already in place + +| Area | Current implementation | Implication for cBMP | +| --- | --- | --- | +| Input framing | [`bmp_tcp_in/io.rs`](../src/units/bmp_tcp_in/io.rs), `bmp_read` and `BmpStream::next`, read length-bounded messages into `Bytes` and validate them with routecore. Tracing can modify the version byte before validation. | A natural capture boundary exists, but exact capture must precede mutation and parser rejection. | +| Input processing | [`router_handler.rs`](../src/units/bmp_tcp_in/router_handler.rs), `process_msg`, applies post-policy suppression, filters and the BMP state machine. | Capture before this path to include every received observation, independently of routing policy. | +| Raw UPDATE forwarding | [`payload.rs`](../src/payload.rs), `Update::RouteMonitoringRaw`, carries the per-peer header and BGP UPDATE. [`machine.rs`](../src/units/bmp_tcp_in/state_machine/machine.rs) constructs it after successful parsing and can correct the ASN-width A flag. The handler emits it before parsed route payloads when `forward_raw_updates` is enabled. | Useful parsing/forwarding precedent, but not a complete or byte-exact BMP archive: common headers and other message types are absent, and the header can change. | +| Statistics forwarding | `Update::PeerStats` carries the statistics count and TLVs; BMP output rebuilds headers. | Some non-route bodies already survive forwarding, but original envelopes still need capture. | +| Attribute representation | `RotondaPaMap` holds an `Arc<[u8]>` containing RPKI and parse metadata followed by raw path attributes. `PathAttributeInterner` uses sharded hash buckets, weak references and byte equality to share identical buffers. | Existing memory deduplication is byte-based, not semantic normalization, durable object storage or a zstd dictionary. Keep internal metadata separate from canonical wire attributes. | +| BMP reconstruction | [`bmp_builder.rs`](../src/units/bmp_tcp_out/bmp_builder.rs) builds control messages, synthetic OPENs, EORs and Route Monitoring. `DumpAggregator` groups by peer, family, ADD-PATH presence and attribute-blob hash, checks equality, and flushes under byte/message limits. | Reuse protocol encoding knowledge and bounded aggregation patterns. Hash-map iteration and fullest-group eviction are not a canonical statistical ordering. | +| Output fastpath | [`bmp-tcp-out`](../docs/bmp-tcp-out.md) preserves live BGP UPDATE bytes while synthesizing BMP identity headers; initial dumps rebuild from the RIB. | Neither path is a corpus of original messages. A RIB snapshot cannot recover event history or original message packaging. | +| History storage | [`clickhouse/event.rs`](../src/targets/clickhouse/event.rs) records route observations, sequence/epoch identities, raw attributes and derived columns. [`spool.rs`](../src/targets/clickhouse/spool.rs) writes immutable, checksummed, row-aligned LZ4 frames. [ClickHouse documentation](../docs/clickhouse.md) describes ZSTD(1) columns and uncompressed RowBinary HTTP output. | Reuse lifecycle and durability ideas, but this is route history rather than full BMP capture. Database column compression is not a shared BMP dictionary. | + +[`Cargo.toml`](../Cargo.toml) has LZ4 and other compression dependencies, but no +direct zstd dependency. There is currently no cBMP codec, statistical layout +profile, shared zstd dictionary lifecycle, or full-message corpus target. + +## Fidelity and canonical identity + +Start with a lossless observation corpus: retain every framed message, including +duplicates, source identity, connection epoch, per-connection sequence number, +receive time and original BMP timestamp. Preserve original message boundaries. +Compression must not collapse observations from different routers or times. +Do not equate equal route bodies with equal policy views or peer sessions. + +Define three separate artifacts: + +1. **Canonical content:** a deterministic structural representation under a + named normalization version. Identical supported structures produce identical + bytes irrespective of incidental attribute ordering. +2. **Observation envelope and reconstruction data:** source/session identity, + timing, boundaries, ordering and any original encodings needed to reverse a + transformation. These preserve differences that canonical content factors out. +3. **Physical layout:** a statistics-driven arrangement of content and envelopes + within a bounded segment, followed by zstd compression. Its profile can change + without changing canonical content identity. + +This distinction resolves a tension: a globally adaptive frequency ordering is +not a timeless canonical form. Freeze the normalization rules and layout profile +by version; use stable content bytes for identity and the profile for placement. +If a future format makes statistical order part of canonical serialization, its +identity must also include that profile version. + +For an initial lossless implementation, use reversible transforms and opaque +fallbacks. Exact reconstruction requires original ordering/encoding residuals +where normalization changes bytes; account for their cost. A later explicitly +semantic export could omit those residuals and emit normalized BMP, but must not +claim byte-for-byte recovery. Preserve unknown attributes, TLVs, unsupported +families, duplicate fields and malformed-but-framed bodies as opaque bytes when +their transformation is not demonstrably safe. Invalid lengths and truncated TCP +messages require capture-error records; they are not valid complete BMP messages. + +## Normalize structures, not just attribute lists + +Represent messages using a common envelope and typed bodies. Serialize repeated +per-peer fields consistently for each observation and let compression exploit +their repetition. Keep volatile timestamps and sequence numbers in separate lanes; +delta-encode only with explicit bases and reversible signed deltas. Clock changes +and zero timestamps must remain representable. + +For Route Monitoring, split the nested UPDATE into withdrawals, shared path +attributes, family/next-hop context, announcements and explicit EOR form. Separate +NLRI from MP_REACH/MP_UNREACH containers so otherwise identical attributes are +not made different by the list of prefixes carried in the UPDATE. Retain the +mapping back to each original message and original container layout. + +Normalize known field encodings conservatively. A canonical attribute ordering +and reversible ordering of supported community collections are candidates. +Keep AS sequence order, AS segment kinds, ASN width, AS4 information, next-hop +structure, ADD-PATH IDs and AFI/SAFI distinctions. Do not infer that every list +is a set, remove duplicate fields, or merge withdrawals and announcements. +Unknown or ambiguous encodings take the opaque path. Original OPEN capabilities +and session context must accompany data that depends on them for decoding. + +Use common skeletons for Peer Up/Down, Initiation/Termination, Statistics and +Route Mirroring, with variable fields/TLV bodies factored out only where safe. +The first codec can leave these bodies opaque while capturing all of them. +Synthetic OPENs from the current restreamer are not substitutes for captured +OPENs. Likewise, decoder output must not inherit the builder's documented +multicast-to-unicast family collapse. + +This also makes packaging differences compressible: one UPDATE carrying 100 +prefixes and 100 UPDATEs carrying one prefix each expose the same canonical +attribute bytes to the compressor. Their envelopes and boundary records remain +different. Feed repeated bytes directly to the encoder without building a durable +attribute-reference table or allocating an extra copy for every occurrence. + +## Statistics-driven ordering + +Collect bounded samples across all configured BMP inputs. Measure message/body +shape frequencies, attribute combinations, repeated values, shared byte prefixes, +field cardinality and co-occurrence. Balance sampling across exporters, time, +initial dumps and live updates so one large feed does not define the whole model. +These are corpus statistics, distinct from incoming BMP Statistics Reports. + +Use those measurements to generate a frozen layout profile: + +- Cluster messages/content by structural shape, family and next-hop form, then + by common attribute structure and values. Compare candidate keys using actual + compressed bytes; frequency alone does not establish useful adjacency. +- Place common stable fields together and separate high-cardinality envelope + fields from reusable bodies. Compare record-oriented and field-lane layouts. +- Order structural groups using measured frequency and similarity. Resolve ties + with canonical bytes; retain every occurrence without assigning content IDs. +- Consider prefix locality inside a compatible group, with an inverse permutation + wherever the original order matters. Never sort AS sequences for compression. + +For example, alternating observations `router-A/body-X`, `router-B/body-Y`, +`router-C/body-X` can place the two serialized body-X occurrences together while +retaining all three observations and their replay positions. Zstd can encode the +second occurrence as a match when it is within its matching window, or use +dictionary matches for patterns present there. The application writes both +occurrences and maintains no body-X deduplication entry. + +Only reorder physical storage within bounded segments. Preserve logical order +using `(source, connection epoch, sequence)` and positional correspondence between +envelope and body lanes, with an inverse permutation where needed. +There is no inherent global causal order across independent TCP connections; +record a collector merge ordinal if reproducing the collector's interleaving is +required. Replay restores order before emitting Peer Up, routes, EOR, Peer Down +and termination. A compressed segment should include the identity/session context +needed to interpret it; starting a new BMP session halfway through history may +still require earlier state or a separate snapshot. + +Profiles should update at segment boundaries, not after every observation. +Deterministic tie-breaking and recorded profile parameters allow reproducible +normalization/layout for a fixed input and segmentation policy. Do not promise +identical compressed bytes across zstd versions or settings. + +## One shared zstd dictionary across incoming streams + +Train on the **normalized serialized representation**, after choosing its layout, +using representative samples pooled across inputs. Share the resulting immutable +dictionary generation among compression workers. Each worker keeps its own +compression context; a shared dictionary does not imply a single mutable +compression stream or a global lock. + +Zstd dictionaries require the corresponding dictionary during decompression and +are particularly useful for small inputs +([zstd documentation](https://github.com/facebook/zstd/blob/dev/programs/zstd.1.md)). +They do not provide unbounded cross-frame history or global deduplication of every +previous message. Choose chunk layout and compression-window settings together: +repetitions outside the available history only benefit from the shared dictionary +when it contains useful matching patterns. The goal is good total compression +without maintaining an application-level index of everything seen before. + +Keep manual content deduplication out of the initial implementation. It can be an +optional benchmark comparison if compression leaves substantial repeated content +uncompressed; introduce it only if measured gains justify the extra tables, +references and lifecycle management. The existing RIB `PathAttributeInterner` +addresses a separate problem: sharing live, uncompressed in-memory route state. +Compressing the corpus does not automatically replace that interner, and this +proposal does not require changing it. + +Use independent zstd frames for bounded chunks inside immutable segments. Flush +on a configurable byte or latency threshold. Compare pooled multi-source chunks +with per-source chunks using the same dictionary: pooling may improve locality, +but adds buffering and couples latency. Large chunks may reduce the incremental +benefit of a dictionary, so measure rather than assume it. + +Each segment manifest should identify the format/normalization version, layout +profile digest, dictionary digest and zstd dictionary ID, codec parameters, +source/epoch sequence ranges, frame sizes and checksums. Store dictionary and +profile artifacts durably before publishing dependent segments; retain them while +any segment references them. A zstd dictionary ID alone is not a content-integrity +check. Exports must bundle the artifacts or specify a durable resolution contract. +Missing/wrong dictionaries and corrupt frames must produce explicit errors. + +Bootstrap with dictionary-free zstd until enough samples exist. Train replacement +generations off the ingestion path and adopt only at segment boundaries after +held-out evaluation. Keep old generations readable; do not rewrite the entire +corpus on rotation. Training controls and sampling memory must be bounded. + +## High-level implementation changes + +1. **Introduce a full BMP capture envelope.** Extend the framing boundary in + `bmp_tcp_in/io.rs` to retain original bytes before tracing mutation and protocol + parsing. Deliver complete framed messages, connection lifecycle and capture + errors to an optional capture channel before filtering/state-machine changes. + Assign persistent source identity, fresh connection epochs and sequences here; + transient/rebound `IngressId` values alone are insufficient. Reuse cheap `Bytes` + clones. Capture must work regardless of `forward_raw_updates`. +2. **Add a shared collector target.** A proposed `cbmp-out` target subscribes to + capture from multiple BMP input units. Prefer a dedicated capture subscription + over adding full raw traffic to every existing route consumer. Wire its config + and lifecycle through `src/targets/mod.rs` and the manager. Keep the capture + envelope/codec independent of the RIB and ClickHouse event model. +3. **Implement a versioned codec library.** Suggested modules cover capture + records, conservative normalization, inverse reconstruction, layout profiles, + dictionary artifacts and segment framing. Begin with opaque lossless bodies, + then support Route Monitoring structure. Reuse routecore parsing with recorded + session context and extract suitable builder helpers; do not route exact replay + through synthetic-header output routines. + Serialize normalized occurrences directly into bounded compression buffers; + do not add a persistent content-deduplication layer. +4. **Add bounded processing and persistence.** Run normalization, sorting, + training and compression away from async socket tasks. Bound queued bytes, + worker memory, sorting windows and disk backlog. Adapt the ClickHouse spool's + checksum/recovery patterns, not its RowBinary schema or LZ4 format. Publish + segments atomically after required durability steps and recover partial writes. + Define queue-full behavior explicitly: backpressure for lossless operation, or + an explicitly configured loss mode with durable gap records and counters. +5. **Provide decoding and output tools.** First implement offline inspect, verify + and decode commands for the corpus. Emit reconstructed per-source BMP streams + in logical order. Later add framed corpus transport with dictionary/profile + delivery and resumable segment identities; keep it separate from standard BMP + TCP output. Reuse the same immutable segment bytes for storage and transfer. +6. **Expose operational evidence.** Report input/canonical/compressed bytes, + dictionary and index overhead, opaque fallback rate, queue size, segment age, + compression/decompression CPU, dictionary generation, gaps and replay failures. + +## Validation and rollout + +Build an offline prototype before coupling the codec to live ingestion. Evaluate +the following on the same captured corpus, segment boundaries and zstd settings: + +| Variant | Question answered | +| --- | --- | +| Raw BMP + zstd, no dictionary | What does ordinary chunk compression already achieve? | +| Raw BMP + shared dictionary | What does dictionary training alone contribute? | +| Attribute-order normalization + zstd | How much does the narrow approach buy? | +| Structural normalization, original event order + zstd | What does factoring message structure contribute? | +| Structural normalization + statistical layout + zstd | What does physical ordering contribute? | +| Full design + shared dictionary | Does the shared model improve the normalized corpus further? | +| Optional explicit deduplication + compression | Is there enough additional benefit to justify application-level tables and references? | + +Train on an earlier interval and test on later intervals and held-out exporters. +Include initial dumps, churn, withdrawals, IPv4/IPv6, multiple policy views, +ADD-PATH, unknown attributes, unsupported families and control messages. Compare +total retained bytes including envelopes, inverse permutations, residuals, indexes, +profiles and amortized dictionaries. Measure ingest/decode throughput, peak memory, +latency, random segment reads and recovery time. Do not use the existing synthetic +ClickHouse compression results as a prediction for this format. + +Correctness gates include byte-exact capture/decode round trips, deterministic and +idempotent canonicalization for supported structures, opaque fallback round trips, +preservation of event multiplicity and per-stream order, and dictionary rotation, +restart, truncation and checksum failures. Verify that reconstruction cannot +confuse peer epochs or ADD-PATH parsing context. Property/fuzz tests belong at the +parser/codec boundary once implementation begins. + +Roll out in stages: offline baseline and codec; optional live lossless capture; +conservative structural normalization; statistical layout and shared dictionary; +then corpus transport and optional semantic BMP export. Advance the more complex +transforms only when their measured savings cover metadata, CPU and latency costs. From 7f2690925c40105ce70476e965436a94fc84f0d1 Mon Sep 17 00:00:00 2001 From: Jeroen van Bemmel Date: Fri, 2 Oct 2026 07:51:19 -0500 Subject: [PATCH 09/10] Document routing scale constraints and cBMP proposal --- AGENTS.md | 57 +++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 57 insertions(+) diff --git a/AGENTS.md b/AGENTS.md index b158b04..dfb1dc9 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -89,3 +89,60 @@ bytes per object and owned heap, allocations/copies, query latency, lock hold time, ingestion throughput, and peak serialization memory. Show how costs grow with table size and concurrency. Follow the whole operation so a later stage does not recreate costs eliminated earlier. + +## Implementation references + +### Storage and buffer accounting (`src/payload.rs`) + +- Keep `RotondaRoute` and `Payload` at their tested 64-bit sizes: 64 and 96 bytes. + EVPN NLRI is boxed because storing it inline added 176 bytes to every route + and payload, including non-EVPN traffic. Measure and justify layout changes. +- When adding owned data, update `Update::shallow_bytes()` accounting for both + `Single` and `Bulk`. Despite its name, it includes EVPN heap storage through + `evpn_heap_bytes()`; counting only `size_of` underestimates buffer usage. +- Extend `route_and_payload_sizes` and + `evpn_buffer_accounting_includes_owned_heap` when changing these representations. + +### RIB reads and HTTP responses (`src/units/rib_unit/`) + +- Use `Rib::evpn_records_matching()` to capture only matching `Arc`s. + Its predicate runs under the writer mutex: put cheap filters there, then decode + attributes, sort, and build response rows outside it. Do not restore a + full-table `Vec` clone. +- Records shared with snapshots must remain unchanged by later RIB mutations. + Follow the existing replacement/`Arc::make_mut` pattern; preserve the snapshot + checks in `evpn_tenants_paths_and_peer_lifecycle`. +- Follow `http_ng.rs::search_evpn`: move the `DumpGuard` permit into + `spawn_blocking` and keep it through `serde_json::to_vec`. Serialize typed rows + directly; returning rows for later encoding or building `json!` trees restores + unguarded CPU work and extra response-sized allocations. + +### Withdrawals and teardown (`src/units/rib_unit/rib.rs`) + +- EVPN insertion results feed downstream metrics. Preserve the current contract: + `mui_count = usize::from(active || existed)` and + `prefix_new`/`mui_new = !existed && active`. Capture `existed` before mutation; + using only resulting `active` misclassifies known withdrawals as unannounced. + See `evpn_withdrawal_reports_prior_existence` for both retention modes. +- Use `withdraw_for_ingresses` for withdrawal semantics and + `remove_for_ingresses` for physical session cleanup. Synthesized BMP peers can + receive fresh ingress IDs each session; marking old records withdrawn leaves + retained slots behind. Include new stores in family-scoped and session cleanup. +- In `remove_for_ingresses`, keep the empty-input return before locking and build + ingress membership outside the EVPN mutex. Do not replace the set lookup with + `ids.contains()` per record: that changes cleanup from O(R + I) to O(R × I). +- ADD-PATH uses child ingresses. Follow existing parent/path provenance when + querying or cleaning up; session cleanup must include its children. + +### EVPN configuration (`src/config.rs`) + +- `enable_evpn` defaults to false and is process-wide after startup. Keep gating + in `CombinedConfig::protocols`/`addpath`, `convert_nlri`, and `search_evpn` aligned. + Disabled ingestion uses unsupported-NLRI accounting. +- Changing enablement requires restart because sessions and retained state are + already established. Keep reload rejection; conversion tests should pass an + explicit setting through `convert_nlri_with_evpn`, not mutate the global setting. + +For changes to these paths, extend the named regression tests and report relevant +layout, allocation, or lock/query cost changes at realistic table sizes. Include +snapshot retention and serialization when estimating concurrent-query memory. From c3d19dc6839cff6f3231861e6035b88d0e8e1769 Mon Sep 17 00:00:00 2001 From: J vanBemmel Date: Fri, 2 Oct 2026 08:03:52 -0500 Subject: [PATCH 10/10] Delete AGENTS.md --- AGENTS.md | 148 ------------------------------------------------------ 1 file changed, 148 deletions(-) delete mode 100644 AGENTS.md diff --git a/AGENTS.md b/AGENTS.md deleted file mode 100644 index dfb1dc9..0000000 --- a/AGENTS.md +++ /dev/null @@ -1,148 +0,0 @@ -# Netom engineering principles - -Netom monitors large BGP/BMP routing tables and many sessions. Correctness -includes protocol semantics, lifecycle and metrics, bounded resource use, and -predictable ingestion under concurrent queries. Evaluate changes at realistic -route/path/ingress counts and at 10x scale, beyond small test fixtures. - -## Keep common objects compact - -- Measure `size_of` before/after changes to core route, payload, queue, and - message types, including enclosing types, alignment, and padding. Justify - increases and report their aggregate cost in the PR. -- A large enum variant enlarges every instance. Keep uncommon large data behind - appropriate indirection; account for allocations, reference counts, and copies. -- On 64-bit targets, the boxed EVPN layout preserves `RotondaRoute` at 64 bytes - and `Payload` at 96 bytes. Inline EVPN raised them to 240/272 bytes: - +176 bytes per object, or roughly 176 MB per million objects. -- Queue/buffer accounting must include owned heap allocations and capacities, - not just stack size. Preserve lightweight layout and accounting regression tests. - -## Filter before copying; preserve snapshot consistency - -- Never deep-clone the entire RIB to satisfy a filtered query or solve ownership. - Apply cheap family, active/withdrawn, and identity filters before copying; - exclude withdrawn records unless requested. -- Prefer references or filtered shared immutable snapshots. Captured records/IDs - must remain valid through consumption; use copy-on-write where needed to keep - snapshots coherent across updates, withdrawals, and removal. -- Snapshot predicates under a writer mutex must be inexpensive. Decode, sort, - transform, serialize, and perform network I/O after releasing the lock. -- Account for snapshot lifetime, retained old versions, and concurrent readers. - Add indexes for frequent expensive queries only with a justified memory cost. - -## Bound work and contention - -- Analyze cost using actual cardinalities: routes, paths, ingresses, updates, - peers, and concurrent queries. Include allocations, copies, and lock duration. -- Avoid linear membership searches inside table scans. For batch ingress cleanup, - build a membership set outside the lock: expected O(I + R), rather than O(I × R). - Empty batches return immediately without locking or scanning. -- Keep hot lock sections short; avoid full-table deep copies, expensive decoding, - serialization, and unnecessary allocations while holding them. Preserve an - explicit consistency model when moving work outside locks. - -## Guard the complete query operation - -- Bound concurrent expensive RIB queries. Hold the resource permit through - extraction, transformation, allocation, and serialization, including work in - `spawn_blocking`; releasing it before JSON encoding defeats the limit. -- Move synchronous CPU-heavy work off Tokio workers. `spawn_blocking` alone - does not bound concurrency or peak memory. -- Serialize typed responses directly. Avoid `json!`/JSON value trees that create - another large representation solely for encoding. -- Assess peak memory across retained state, snapshots, response rows, encoded - bodies, and concurrent requests; final response size alone is insufficient. - -## Preserve identity, lifecycle, and metric semantics - -- Route identity must include the correct NLRI, ingress/session provenance, and - ADD-PATH identity. Equal NLRIs can have independent paths; withdrawing one must - leave the others active. -- Derive metrics from prior existence/state and the event, not just resulting - `active`. A known-route withdrawal must not become an unannounced withdrawal; - a different ingress/path is not a matching prior record. -- Define retained-withdrawn behavior explicitly, including repeated withdrawals, - re-announcements, and “new” counters. Test both retention modes. -- Every insertion needs cleanup for withdrawal, family withdrawal, peer/session - teardown, and ingress removal. Marking withdrawn is not physical removal; - expired sessions must not leak one retained slot per prefix/path. -- Test known/unknown withdrawals, updates and re-announcements, independent - ADD-PATH withdrawals, family/peer cleanup, and snapshot stability during mutation. - -## Make optional family support consistent - -- EVPN requires explicit runtime opt-in (`enable_evpn`, default false). Gate - negotiation (including ADD-PATH), ingestion, and APIs consistently; account - for disabled-family drops without retaining their routing state. -- Configuration affecting negotiated sessions and retained state must not change - on reload without a safe transition. EVPN enablement currently requires restart. - Tests should inject settings rather than race on process-global configuration. -- Keep EVPN semantics distinct: RD is not tenant identity; route-target selection - is not VRF import-policy simulation; a numeric label/VNI does not establish - L2 versus L3 role. Preserve Type 2/Type 5 data needed for symmetric IRB analysis. - -## Validate the cost as well as the behavior - -For significant RIB/API changes, report relevant before/after measurements: -bytes per object and owned heap, allocations/copies, query latency, lock hold -time, ingestion throughput, and peak serialization memory. Show how costs grow -with table size and concurrency. Follow the whole operation so a later stage -does not recreate costs eliminated earlier. - -## Implementation references - -### Storage and buffer accounting (`src/payload.rs`) - -- Keep `RotondaRoute` and `Payload` at their tested 64-bit sizes: 64 and 96 bytes. - EVPN NLRI is boxed because storing it inline added 176 bytes to every route - and payload, including non-EVPN traffic. Measure and justify layout changes. -- When adding owned data, update `Update::shallow_bytes()` accounting for both - `Single` and `Bulk`. Despite its name, it includes EVPN heap storage through - `evpn_heap_bytes()`; counting only `size_of` underestimates buffer usage. -- Extend `route_and_payload_sizes` and - `evpn_buffer_accounting_includes_owned_heap` when changing these representations. - -### RIB reads and HTTP responses (`src/units/rib_unit/`) - -- Use `Rib::evpn_records_matching()` to capture only matching `Arc`s. - Its predicate runs under the writer mutex: put cheap filters there, then decode - attributes, sort, and build response rows outside it. Do not restore a - full-table `Vec` clone. -- Records shared with snapshots must remain unchanged by later RIB mutations. - Follow the existing replacement/`Arc::make_mut` pattern; preserve the snapshot - checks in `evpn_tenants_paths_and_peer_lifecycle`. -- Follow `http_ng.rs::search_evpn`: move the `DumpGuard` permit into - `spawn_blocking` and keep it through `serde_json::to_vec`. Serialize typed rows - directly; returning rows for later encoding or building `json!` trees restores - unguarded CPU work and extra response-sized allocations. - -### Withdrawals and teardown (`src/units/rib_unit/rib.rs`) - -- EVPN insertion results feed downstream metrics. Preserve the current contract: - `mui_count = usize::from(active || existed)` and - `prefix_new`/`mui_new = !existed && active`. Capture `existed` before mutation; - using only resulting `active` misclassifies known withdrawals as unannounced. - See `evpn_withdrawal_reports_prior_existence` for both retention modes. -- Use `withdraw_for_ingresses` for withdrawal semantics and - `remove_for_ingresses` for physical session cleanup. Synthesized BMP peers can - receive fresh ingress IDs each session; marking old records withdrawn leaves - retained slots behind. Include new stores in family-scoped and session cleanup. -- In `remove_for_ingresses`, keep the empty-input return before locking and build - ingress membership outside the EVPN mutex. Do not replace the set lookup with - `ids.contains()` per record: that changes cleanup from O(R + I) to O(R × I). -- ADD-PATH uses child ingresses. Follow existing parent/path provenance when - querying or cleaning up; session cleanup must include its children. - -### EVPN configuration (`src/config.rs`) - -- `enable_evpn` defaults to false and is process-wide after startup. Keep gating - in `CombinedConfig::protocols`/`addpath`, `convert_nlri`, and `search_evpn` aligned. - Disabled ingestion uses unsupported-NLRI accounting. -- Changing enablement requires restart because sessions and retained state are - already established. Keep reload rejection; conversion tests should pass an - explicit setting through `convert_nlri_with_evpn`, not mutate the global setting. - -For changes to these paths, extend the named regression tests and report relevant -layout, allocation, or lock/query cost changes at realistic table sizes. Include -snapshot retention and serialization when estimating concurrent-query memory.