diff --git a/daemon/src/event/grpc.rs b/daemon/src/event/grpc.rs index 73202d2d..1e038647 100644 --- a/daemon/src/event/grpc.rs +++ b/daemon/src/event/grpc.rs @@ -667,12 +667,14 @@ impl From for PeerGroup { } } +type PathUuidMap = FnvHashMap)>; + pub(super) struct GrpcService { init: Arc, active_conn_tx: mpsc::UnboundedSender, pub(super) global: GlobalHandle, pub(super) tables: TableHandle, - path_uuid_map: tokio::sync::Mutex)>>, + path_uuid_map: tokio::sync::Mutex, } /// Validate and convert `api::EbgpMultihop` to the internal `Option`. @@ -770,6 +772,70 @@ impl GrpcService { } } + fn withdraw_local_paths( + &self, + uuid_map: &mut PathUuidMap, + family: Family, + nets: &[packet::PathNlri], + ) { + let timestamp = crate::proto::unix_secs(); + for net in nets { + self.tables + .remove_route(table::Source::local(), family, net.clone(), None, timestamp); + } + let removed: FnvHashSet<_> = nets.iter().collect(); + // Forget every handle for a removed path, so a stale UUID cannot + // delete a later announcement with the same NLRI and identifier. + uuid_map.retain(|_, (mapped_family, mapped_nets)| { + if *mapped_family == family { + mapped_nets.retain(|net| !removed.contains(net)); + } + !mapped_nets.is_empty() + }); + } + + fn local_paths_to_delete( + &self, + uuid_map: &PathUuidMap, + family: Option, + vrf: Option<&table::Vrf>, + ) -> FnvHashMap> { + let mut paths: FnvHashMap> = FnvHashMap::default(); + let in_vrf = |net: &packet::PathNlri| { + vrf.is_none_or(|vrf| match &net.nlri { + packet::Nlri::VpnV4(n) => n.rd == vrf.rd && n.labels.labels() == [vrf.label], + packet::Nlri::VpnV6(n) => n.rd == vrf.rd && n.labels.labels() == [vrf.label], + _ => false, + }) + }; + for shard in &self.tables.shards { + let shard = shard.lock().unwrap(); + for f in shard + .rtable + .families() + .filter(|f| family.is_none_or(|family| *f == family)) + { + // iter_reach includes policy-filtered routes and their original + // path identifiers, including paths injected without a UUID. + for reach in shard.rtable.iter_reach(f).filter(|r| r.source.is_local()) { + if in_vrf(&reach.net) { + paths.entry(f).or_default().push(reach.net); + } + } + } + } + // Also invalidate handles for local paths replaced by kernel updates. + for (f, nets) in uuid_map.values() { + if family.is_none_or(|family| *f == family) { + paths + .entry(*f) + .or_default() + .extend(nets.iter().filter(|net| in_vrf(net)).cloned()); + } + } + paths + } + async fn is_available(&self, need_active: bool) -> Result<(), Error> { let global = &self.global.read().await; if need_active && global.asn == 0 { @@ -795,14 +861,38 @@ impl GrpcService { Some(family) => convert::family_from_api(&family), None => Family::IPV4, }; - let net = convert::net_from_api(path.nlri.ok_or(Error::EmptyArgument)?, family) - .map_err(|_| tonic::Status::new(tonic::Code::InvalidArgument, "prefix is invalid"))?; + // GoBGP gives the binary fields precedence when both forms are supplied. + let net = if path.nlri_binary.is_empty() { + convert::net_from_api(path.nlri.ok_or(Error::EmptyArgument)?, family) + .map_err(|_| tonic::Status::invalid_argument("prefix is invalid"))? + } else { + packet::Nlri::decode_from_bytes(family, &path.nlri_binary) + .map_err(|_| tonic::Status::invalid_argument("invalid binary NLRI"))? + }; + let attributes = if path.pattrs_binary.is_empty() { + path.pattrs + .into_iter() + .map(|a| { + convert::attr_from_api(a) + .map_err(|_| tonic::Status::invalid_argument("invalid attribute")) + }) + .collect::, _>>()? + } else { + path.pattrs_binary + .iter() + .map(|a| { + packet::Attribute::decode_from_bytes(a) + .map_err(|_| tonic::Status::invalid_argument("invalid binary attribute")) + }) + .collect::, _>>()? + }; let mut attr = Vec::new(); let mut nexthop = None; - for a in path.pattrs { - let a = convert::attr_from_api(a).map_err(|_| { - tonic::Status::new(tonic::Code::InvalidArgument, "invalid attribute") - })?; + let mut seen = FnvHashSet::default(); + for a in attributes { + if !seen.insert(a.code()) { + return Err(tonic::Status::invalid_argument("duplicate attribute")); + } match a.code() { bgp::Attribute::MP_REACH => { // MP_REACH binary: [AFI:2][SAFI:1][NH_LEN:1][nexthop:NH_LEN][reserved:1][NLRI...] @@ -833,6 +923,9 @@ impl GrpcService { } bgp::Attribute::NEXTHOP => { nexthop = a.binary().and_then(|b| bgp::Nexthop::from_bytes(b)); + if nexthop.is_none() { + return Err(tonic::Status::invalid_argument("invalid nexthop")); + } } // RR attributes are added on reflection and must not be set by operators. // MP_UNREACH has no meaning in an add_path request. @@ -1981,8 +2074,9 @@ impl GoBgpService for GrpcService { let inner = request.into_inner(); let table_type = api::TableType::try_from(inner.table_type).unwrap_or(api::TableType::Global); - let (mut family, nets, attrs, nexthop) = - self.local_path(inner.path.ok_or(Error::EmptyArgument)?)?; + let path = inner.path.ok_or(Error::EmptyArgument)?; + let is_withdraw = path.is_withdraw; + let (mut family, nets, attrs, nexthop) = self.local_path(path)?; let mut insert_nets = nets.clone(); let mut insert_attrs = attrs; if table_type == api::TableType::Vrf { @@ -2008,6 +2102,12 @@ impl GoBgpService for GrpcService { let map_nets = insert_nets.clone(); let timestamp = crate::proto::unix_secs(); let source = table::Source::local(); + // Serialize route mutations and UUID bookkeeping with DeletePath. + let mut uuid_map = self.path_uuid_map.lock().await; + if is_withdraw { + self.withdraw_local_paths(&mut uuid_map, family, &insert_nets); + return Ok(tonic::Response::new(api::AddPathResponse::default())); + } if let Some(attrs) = insert_attrs { for net in insert_nets { self.tables.insert_route( @@ -2022,10 +2122,7 @@ impl GoBgpService for GrpcService { } } let id = uuid::Uuid::new_v4(); - self.path_uuid_map - .lock() - .await - .insert(id, (family, map_nets)); + uuid_map.insert(id, (family, map_nets)); Ok(tonic::Response::new(api::AddPathResponse { uuid: id.as_bytes().to_vec(), })) @@ -2036,25 +2133,73 @@ impl GoBgpService for GrpcService { ) -> Result, tonic::Status> { let inner = request.into_inner(); if inner.uuid.is_empty() { - return Err(tonic::Status::new( - tonic::Code::InvalidArgument, - "uuid is required", - )); + let table_type = api::TableType::try_from(inner.table_type) + .map_err(|_| tonic::Status::invalid_argument("invalid table type"))?; + if !matches!( + table_type, + api::TableType::Unspecified | api::TableType::Global | api::TableType::Vrf + ) { + return Err(tonic::Status::invalid_argument( + "DeletePath only supports global and VRF tables", + )); + } + let vrf = if table_type == api::TableType::Vrf { + if inner.vrf_id.is_empty() { + return Err(tonic::Status::invalid_argument( + "vrf_id is required for VRF table type", + )); + } + Some( + self.tables + .list_vrfs(Some(&inner.vrf_id)) + .into_iter() + .next() + .ok_or_else(|| { + tonic::Status::not_found(format!("VRF '{}' not found", inner.vrf_id)) + })?, + ) + } else { + None + }; + let paths = if let Some(path) = inner.path { + let (mut family, mut nets, attrs, _) = self.local_path(path)?; + if let Some(vrf) = &vrf { + (family, nets, _) = vrf_export_path(family, nets, attrs, vrf)?; + } + Some((family, nets)) + } else { + None + }; + let mut uuid_map = self.path_uuid_map.lock().await; + if let Some((family, nets)) = paths { + self.withdraw_local_paths(&mut uuid_map, family, &nets); + } else { + let mut family = inner.family.as_ref().map(convert::family_from_api); + if vrf.is_some() { + family = match family { + Some(Family::IPV4) => Some(Family::IPV4_VPN), + Some(Family::IPV6) => Some(Family::IPV6_VPN), + None => None, + _ => { + return Err(tonic::Status::invalid_argument( + "VRF DeletePath only supports IPv4/IPv6 families", + )); + } + }; + } + for (family, nets) in self.local_paths_to_delete(&uuid_map, family, vrf.as_ref()) { + self.withdraw_local_paths(&mut uuid_map, family, &nets); + } + } + return Ok(tonic::Response::new(api::DeletePathResponse {})); } let id = uuid::Uuid::from_slice(&inner.uuid) .map_err(|_| tonic::Status::new(tonic::Code::InvalidArgument, "invalid uuid"))?; - let (family, nets) = self - .path_uuid_map - .lock() - .await + let mut uuid_map = self.path_uuid_map.lock().await; + let (family, nets) = uuid_map .remove(&id) .ok_or_else(|| tonic::Status::new(tonic::Code::NotFound, "uuid not found"))?; - let timestamp = crate::proto::unix_secs(); - let source = table::Source::local(); - for net in nets { - self.tables - .remove_route(source.clone(), family, net, None, timestamp); - } + self.withdraw_local_paths(&mut uuid_map, family, &nets); Ok(tonic::Response::new(api::DeletePathResponse {})) } type ListPathStream = Pin< diff --git a/daemon/src/event/mod.rs b/daemon/src/event/mod.rs index 014667c1..0983042c 100644 --- a/daemon/src/event/mod.rs +++ b/daemon/src/event/mod.rs @@ -9278,6 +9278,371 @@ mod tests { } } + async fn delete_gobgp_path(svc: &GrpcService, uuid: Vec) -> Result<(), tonic::Status> { + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + uuid, + ..Default::default() + })) + .await + .map(|_| ()) + } + + fn gobgp_binary_path(family: Family) -> api::Path { + let ipv6 = family.afi() == 2; + let flowspec = family.safi() == 133; + let address = if ipv6 { + "2001:db8::42" + .parse::() + .unwrap() + .octets() + .to_vec() + } else { + vec![198, 51, 100, 42] + }; + let mut prefix = vec![if ipv6 { 128 } else { 32 }]; + if flowspec && ipv6 { + prefix.push(0); // IPv6 FlowSpec prefix offset + } + prefix.extend(address); + let nlri = if flowspec { + let mut body = vec![1]; // destination prefix + body.extend(prefix); + body.extend([3, 0x81, 17, 5, 0x81, 53]); // UDP, destination port 53 + let mut nlri = vec![body.len() as u8]; + nlri.extend(body); + nlri + } else { + prefix + }; + let mut attrs = vec![vec![0x40, 1, 1, 0]]; // ORIGIN IGP + if family == Family::IPV4 { + attrs.push(vec![0x40, 3, 4, 192, 0, 2, 1]); + } else { + let nh = if flowspec { + vec![] + } else { + "2001:db8::1".parse::().unwrap().octets().to_vec() + }; + let mut mp = vec![0, family.afi() as u8, family.safi(), nh.len() as u8]; + mp.extend(nh); + mp.push(0); // reserved + mp.extend(&nlri); + let mut attr = vec![0x80, 14, mp.len() as u8]; + attr.extend(mp); + attrs.push(attr); + } + if flowspec { + attrs.push(vec![0xc0, 16, 8, 0x80, 6, 0, 0, 0, 0, 0, 0]); // discard + } + attrs.push(vec![0xc0, 8, 4, 0xfd, 0xe8, 0x03, 0x85]); // 65000:901 + api::Path { + family: Some(convert::family_to_api(family)), + nlri_binary: nlri, + pattrs_binary: attrs, + ..Default::default() + } + } + + async fn add_gobgp_path( + svc: &GrpcService, + path: api::Path, + ) -> Result { + svc.add_path(tonic::Request::new(api::AddPathRequest { + path: Some(path), + ..Default::default() + })) + .await + .map(tonic::Response::into_inner) + } + + #[tokio::test] + async fn gobgp_compat_binary_announcements_preserve_nlri_attributes_and_nexthop() { + for family in [ + Family::IPV4, + Family::IPV6, + Family::IPV4_FLOWSPEC, + Family::IPV6_FLOWSPEC, + ] { + let svc = make_grpc_service(); + let path = gobgp_binary_path(family); + let expected_nlri = path.nlri_binary.clone(); + let uuid = add_gobgp_path(&svc, path).await.unwrap().uuid; + assert_eq!(uuid.len(), 16); + let paths = svc.tables.collect_loc_rib_paths(family); + assert_eq!(paths.len(), 1); + assert_eq!(paths[0].net.encode_to_bytes(), expected_nlri); + let attrs = &paths[0].new_best().unwrap().attr; + assert!(attrs.iter().any(|a| a.code() == bgp::Attribute::COMMUNITY + && a.binary().unwrap() == &[0xfd, 0xe8, 0x03, 0x85])); + if family.safi() == 133 { + assert!(paths[0].new_best().unwrap().nexthop.is_none()); + assert!( + attrs + .iter() + .any(|a| a.code() == bgp::Attribute::EXTENDED_COMMUNITY + && a.binary().unwrap() == &[0x80, 6, 0, 0, 0, 0, 0, 0]) + ); + } else { + let expected = if family == Family::IPV4 { + "192.0.2.1" + } else { + "2001:db8::1" + }; + assert_eq!( + paths[0] + .new_best() + .unwrap() + .nexthop + .unwrap() + .addr() + .to_string(), + expected + ); + } + } + } + + #[tokio::test] + async fn gobgp_compat_binary_fields_take_precedence_over_structured_fields() { + let svc = make_grpc_service(); + let mut path = gobgp_binary_path(Family::IPV4); + let structured = ipv4_path("10.0.0.0", 24, "10.0.0.1"); + path.nlri = structured.nlri; + path.pattrs = structured.pattrs; + add_gobgp_path(&svc, path).await.unwrap(); + let paths = svc.tables.collect_loc_rib_paths(Family::IPV4); + assert_eq!(paths[0].net.to_string(), "198.51.100.42/32"); + assert_eq!( + paths[0] + .new_best() + .unwrap() + .nexthop + .unwrap() + .addr() + .to_string(), + "192.0.2.1" + ); + } + + #[tokio::test] + async fn gobgp_compat_binary_extended_length_attributes_are_accepted() { + let svc = make_grpc_service(); + let mut path = gobgp_binary_path(Family::IPV4); + let communities = [0xfd, 0xe8, 0x03, 0x85].repeat(70); + let mut attr = vec![0xd0, 8, 1, 24]; // extended length: 280 bytes + attr.extend(&communities); + *path.pattrs_binary.last_mut().unwrap() = attr; + add_gobgp_path(&svc, path).await.unwrap(); + let paths = svc.tables.collect_loc_rib_paths(Family::IPV4); + let attrs = &paths[0].new_best().unwrap().attr; + assert_eq!( + attrs + .iter() + .find(|a| a.code() == 8) + .unwrap() + .binary() + .unwrap(), + &communities + ); + } + + #[tokio::test] + async fn gobgp_compat_invalid_binary_attribute_flags_do_not_install_routes() { + for (family, index, flags) in [ + (Family::IPV4, 0, 0xe0), // ORIGIN: Optional must be clear + (Family::IPV4, 0, 0x00), // ORIGIN: Transitive must be set + (Family::IPV6, 1, 0xc0), // MP_REACH: Transitive must be clear + (Family::IPV6, 1, 0x00), // MP_REACH: Optional must be set + (Family::IPV4, 2, 0x80), // COMMUNITY: Transitive must be set + (Family::IPV4, 2, 0x40), // COMMUNITY: Optional must be set + ] { + let svc = make_grpc_service(); + let mut path = gobgp_binary_path(family); + path.pattrs_binary[index][0] = flags; + assert_eq!( + add_gobgp_path(&svc, path).await.unwrap_err().code(), + tonic::Code::InvalidArgument, + "family {family:?}, attribute {index}, flags {flags:#x}" + ); + assert!(svc.tables.collect_loc_rib_paths(family).is_empty()); + } + } + + #[tokio::test] + async fn gobgp_compat_binary_partial_community_is_accepted() { + let svc = make_grpc_service(); + let mut path = gobgp_binary_path(Family::IPV4); + path.pattrs_binary[2][0] = 0xe0; // Optional, Transitive, Partial + add_gobgp_path(&svc, path).await.unwrap(); + let paths = svc.tables.collect_loc_rib_paths(Family::IPV4); + let community = paths[0] + .new_best() + .unwrap() + .attr + .iter() + .find(|a| a.code() == bgp::Attribute::COMMUNITY) + .unwrap(); + assert_eq!(community.flags(), 0xe0); + } + + #[tokio::test] + async fn gobgp_compat_malformed_binary_requests_do_not_install_routes() { + let valid = gobgp_binary_path(Family::IPV4); + let mut bad_paths = vec![]; + for nlri in [vec![33, 1, 2, 3, 4, 5], vec![32, 1], vec![0, 0]] { + bad_paths.push(api::Path { + nlri_binary: nlri, + ..valid.clone() + }); + } + for attr in [ + vec![0x40], + vec![0x40, 1, 2, 0], + vec![0x40, 1, 1, 3], + vec![0xc0, 8, 3, 1, 2, 3], + vec![0x40, 3, 1, 1], + vec![0x80, 14, 1, 0], + ] { + bad_paths.push(api::Path { + pattrs_binary: vec![attr], + ..valid.clone() + }); + } + let mut duplicate = valid.clone(); + duplicate + .pattrs_binary + .push(duplicate.pattrs_binary[0].clone()); + bad_paths.push(duplicate); + for path in bad_paths { + let svc = make_grpc_service(); + assert_eq!( + add_gobgp_path(&svc, path).await.unwrap_err().code(), + tonic::Code::InvalidArgument + ); + assert!(svc.tables.collect_loc_rib_paths(Family::IPV4).is_empty()); + } + } + + #[tokio::test] + async fn gobgp_compat_legacy_withdraw_and_uuid_delete_work_for_unicast_and_flowspec() { + for family in [ + Family::IPV4, + Family::IPV6, + Family::IPV4_FLOWSPEC, + Family::IPV6_FLOWSPEC, + ] { + let svc = make_grpc_service(); + let path = gobgp_binary_path(family); + let old_uuid = add_gobgp_path(&svc, path.clone()).await.unwrap().uuid; + let withdrawal = api::Path { + is_withdraw: true, + ..path.clone() + }; + assert!( + add_gobgp_path(&svc, withdrawal.clone()) + .await + .unwrap() + .uuid + .is_empty() + ); + assert!(svc.tables.collect_loc_rib_paths(family).is_empty()); + // Withdrawing an absent path is idempotent. + assert!( + add_gobgp_path(&svc, withdrawal) + .await + .unwrap() + .uuid + .is_empty() + ); + let new_uuid = add_gobgp_path(&svc, path).await.unwrap().uuid; + assert_eq!( + delete_gobgp_path(&svc, old_uuid).await.unwrap_err().code(), + tonic::Code::NotFound + ); + assert_eq!(svc.tables.collect_loc_rib_paths(family).len(), 1); + delete_gobgp_path(&svc, new_uuid).await.unwrap(); + assert!(svc.tables.collect_loc_rib_paths(family).is_empty()); + } + } + + #[tokio::test] + async fn gobgp_compat_legacy_withdraw_is_scoped_to_local_source_prefix_and_path_id() { + let svc = make_grpc_service(); + let path = gobgp_binary_path(Family::IPV4); + let nlri = packet::Nlri::decode_from_bytes(Family::IPV4, &path.nlri_binary).unwrap(); + let peer = Arc::new(table::Source::new( + "192.0.2.2".parse().unwrap(), + "192.0.2.1".parse().unwrap(), + 65002, + 65001, + "192.0.2.2".parse().unwrap(), + table::PeerRole::Ebgp, + )); + svc.tables.insert_route( + peer.clone(), + Family::IPV4, + packet::PathNlri::new(nlri), + Some(bgp::Nexthop::V4("192.0.2.2".parse().unwrap())), + Arc::new(vec![packet::Attribute::empty_as_path()]), + None, + 0, + ); + add_gobgp_path(&svc, path.clone()).await.unwrap(); + let other_id_uuid = add_gobgp_path( + &svc, + api::Path { + identifier: 7, + ..path.clone() + }, + ) + .await + .unwrap() + .uuid; + let other_prefix_uuid = add_gobgp_path(&svc, ipv4_path("10.0.0.0", 24, "192.0.2.1")) + .await + .unwrap() + .uuid; + add_gobgp_path( + &svc, + api::Path { + is_withdraw: true, + ..path + }, + ) + .await + .unwrap(); + let paths = svc.tables.collect_loc_rib_paths(Family::IPV4); + assert_eq!(paths.len(), 2); + let target = paths + .iter() + .find(|p| p.net.to_string() == "198.51.100.42/32") + .unwrap(); + assert_eq!(target.current_paths.len(), 2); // remote + local path ID 7 + assert!( + target + .current_paths + .iter() + .any(|p| Arc::ptr_eq(&p.source, &peer)) + ); + delete_gobgp_path(&svc, other_id_uuid).await.unwrap(); + delete_gobgp_path(&svc, other_prefix_uuid).await.unwrap(); + let remaining = svc.tables.collect_loc_rib_paths(Family::IPV4); + assert_eq!(remaining.len(), 1); + assert_eq!(remaining[0].current_paths.len(), 1); + assert!(Arc::ptr_eq(&remaining[0].current_paths[0].source, &peer)); + } + + #[tokio::test] + async fn gobgp_compat_structured_legacy_withdraw_needs_only_nlri() { + let svc = make_grpc_service(); + let mut path = ipv4_path("10.0.0.0", 24, "192.0.2.1"); + add_gobgp_path(&svc, path.clone()).await.unwrap(); + path.is_withdraw = true; + path.pattrs.clear(); + add_gobgp_path(&svc, path).await.unwrap(); + assert!(svc.tables.collect_loc_rib_paths(Family::IPV4).is_empty()); + } + #[tokio::test] async fn list_path_adj_in_restores_structured_and_binary_nexthops() { let svc = make_grpc_service(); @@ -9444,14 +9809,263 @@ mod tests { } #[tokio::test] - async fn delete_path_without_uuid_is_rejected() { + async fn delete_path_without_selectors_is_idempotent() { let svc = make_grpc_service(); let req = tonic::Request::new(api::DeletePathRequest { uuid: vec![], ..Default::default() }); - let err = svc.delete_path(req).await.unwrap_err(); - assert_eq!(err.code(), tonic::Code::InvalidArgument); + svc.delete_path(req).await.unwrap(); + assert!(svc.tables.collect_loc_rib_paths(Family::IPV4).is_empty()); + } + + #[tokio::test] + async fn delete_path_by_nlri_is_idempotent_and_invalidates_uuids() { + for family in [ + Family::IPV4, + Family::IPV6, + Family::IPV4_FLOWSPEC, + Family::IPV6_FLOWSPEC, + ] { + let svc = make_grpc_service(); + let path = gobgp_binary_path(family); + let first = add_gobgp_path(&svc, path.clone()).await.unwrap().uuid; + let second = add_gobgp_path(&svc, path.clone()).await.unwrap().uuid; + let request = api::DeletePathRequest { + path: Some(path.clone()), + ..Default::default() + }; + svc.delete_path(tonic::Request::new(request.clone())) + .await + .unwrap(); + assert!(svc.tables.collect_loc_rib_paths(family).is_empty()); + svc.delete_path(tonic::Request::new(request)).await.unwrap(); + let current = add_gobgp_path(&svc, path).await.unwrap().uuid; + for stale in [first, second] { + assert_eq!( + delete_gobgp_path(&svc, stale).await.unwrap_err().code(), + tonic::Code::NotFound + ); + } + assert_eq!(svc.tables.collect_loc_rib_paths(family).len(), 1); + delete_gobgp_path(&svc, current).await.unwrap(); + } + } + + #[tokio::test] + async fn delete_path_by_structured_nlri_needs_no_attributes() { + let svc = make_grpc_service(); + let mut path = ipv4_path("10.2.0.0", 24, "10.0.0.1"); + add_gobgp_path(&svc, path.clone()).await.unwrap(); + path.pattrs.clear(); + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + path: Some(path), + ..Default::default() + })) + .await + .unwrap(); + assert!(svc.tables.collect_loc_rib_paths(Family::IPV4).is_empty()); + } + + #[tokio::test] + async fn delete_path_by_nlri_preserves_other_identifiers_and_prefixes() { + let svc = make_grpc_service(); + let path = ipv4_path_with_id("10.2.0.0", 24, "10.0.0.1", 1); + let removed = add_gobgp_path(&svc, path.clone()).await.unwrap().uuid; + let other_id = add_gobgp_path(&svc, ipv4_path_with_id("10.2.0.0", 24, "10.0.0.2", 2)) + .await + .unwrap() + .uuid; + let other_prefix = add_gobgp_path(&svc, ipv4_path("10.3.0.0", 24, "10.0.0.3")) + .await + .unwrap() + .uuid; + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + path: Some(path), + ..Default::default() + })) + .await + .unwrap(); + let paths = svc.tables.collect_loc_rib_paths(Family::IPV4); + assert_eq!(paths.len(), 2); + assert!(paths.iter().all(|p| p.current_paths.len() == 1)); + assert_eq!( + delete_gobgp_path(&svc, removed).await.unwrap_err().code(), + tonic::Code::NotFound + ); + delete_gobgp_path(&svc, other_id).await.unwrap(); + delete_gobgp_path(&svc, other_prefix).await.unwrap(); + } + + #[tokio::test] + async fn delete_path_by_nlri_preserves_peer_routes() { + let svc = make_grpc_service(); + let path = ipv4_path("10.2.0.0", 24, "10.0.0.1"); + add_gobgp_path(&svc, path.clone()).await.unwrap(); + let family = Family::IPV4; + let installed = svc.tables.collect_loc_rib_paths(family); + let local = installed[0].new_best().unwrap(); + let peer = Arc::new(table::Source::new( + "192.0.2.1".parse().unwrap(), + "192.0.2.2".parse().unwrap(), + 65002, + 65001, + Ipv4Addr::new(192, 0, 2, 1), + PeerRole::Ebgp, + )); + svc.tables.insert_route( + peer.clone(), + family, + packet::PathNlri { + path_id: 0, + nlri: installed[0].net.clone(), + }, + local.nexthop, + local.attr.clone(), + None, + 0, + ); + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + path: Some(path), + ..Default::default() + })) + .await + .unwrap(); + let paths = svc.tables.collect_loc_rib_paths(family); + assert_eq!(paths.len(), 1); + assert_eq!(paths[0].current_paths.len(), 1); + assert!(Arc::ptr_eq(&paths[0].current_paths[0].source, &peer)); + } + + #[tokio::test] + async fn delete_path_by_nlri_preserves_other_vrfs() { + for family in [Family::IPV4, Family::IPV6] { + let svc = make_grpc_service(); + for (name, rd) in [("blue", 1), ("red", 2)] { + svc.add_vrf(tonic::Request::new(make_vrf_req( + name, 65000, rd, 65000, rd, + ))) + .await + .unwrap(); + } + let path = gobgp_binary_path(family); + let global = add_gobgp_path(&svc, path.clone()).await.unwrap().uuid; + let mut handles = vec![]; + for name in ["blue", "red"] { + handles.push( + svc.add_path(tonic::Request::new(api::AddPathRequest { + table_type: api::TableType::Vrf as i32, + vrf_id: name.into(), + path: Some(path.clone()), + })) + .await + .unwrap() + .into_inner() + .uuid, + ); + } + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + table_type: api::TableType::Vrf as i32, + vrf_id: "blue".into(), + path: Some(path), + ..Default::default() + })) + .await + .unwrap(); + let vpn = if family == Family::IPV4 { + Family::IPV4_VPN + } else { + Family::IPV6_VPN + }; + assert_eq!(svc.tables.collect_loc_rib_paths(vpn).len(), 1); + assert_eq!(svc.tables.collect_loc_rib_paths(family).len(), 1); + assert_eq!( + delete_gobgp_path(&svc, handles.remove(0)) + .await + .unwrap_err() + .code(), + tonic::Code::NotFound + ); + delete_gobgp_path(&svc, handles.remove(0)).await.unwrap(); + delete_gobgp_path(&svc, global).await.unwrap(); + } + } + + #[tokio::test] + async fn delete_path_rejects_invalid_selectors_without_removing_routes() { + let svc = make_grpc_service(); + let path = gobgp_binary_path(Family::IPV4); + let uuid = add_gobgp_path(&svc, path.clone()).await.unwrap().uuid; + let mut bad_path = path.clone(); + bad_path.nlri_binary = vec![32, 1]; + let mut requests = vec![api::DeletePathRequest { + path: Some(bad_path), + ..Default::default() + }]; + for table_type in [ + api::TableType::Local as i32, + api::TableType::AdjIn as i32, + api::TableType::AdjOut as i32, + 99, + ] { + requests.push(api::DeletePathRequest { + table_type, + path: Some(path.clone()), + ..Default::default() + }); + } + requests.push(api::DeletePathRequest { + table_type: api::TableType::Vrf as i32, + path: Some(path.clone()), + ..Default::default() + }); + for req in requests { + assert_eq!( + svc.delete_path(tonic::Request::new(req)) + .await + .unwrap_err() + .code(), + tonic::Code::InvalidArgument + ); + } + assert_eq!( + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + table_type: api::TableType::Vrf as i32, + vrf_id: "missing".into(), + path: Some(path), + ..Default::default() + })) + .await + .unwrap_err() + .code(), + tonic::Code::NotFound + ); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV4).len(), 1); + delete_gobgp_path(&svc, uuid).await.unwrap(); + } + + #[tokio::test] + async fn delete_path_uuid_takes_precedence_over_path() { + let svc = make_grpc_service(); + let uuid = add_gobgp_path(&svc, ipv4_path("10.2.0.0", 24, "10.0.0.1")) + .await + .unwrap() + .uuid; + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + uuid, + path: Some(api::Path::default()), + ..Default::default() + })) + .await + .unwrap(); + assert!(svc.tables.collect_loc_rib_paths(Family::IPV4).is_empty()); + assert_eq!( + delete_gobgp_path(&svc, vec![1, 2]) + .await + .unwrap_err() + .code(), + tonic::Code::InvalidArgument + ); } #[tokio::test] @@ -9466,6 +10080,164 @@ mod tests { assert_eq!(err.code(), tonic::Code::NotFound); } + #[tokio::test] + async fn delete_path_all_respects_family_and_preserves_nonlocal_routes() { + let svc = make_grpc_service(); + let path = ipv4_path("10.2.0.0", 24, "10.0.0.1"); + let mut handles = vec![add_gobgp_path(&svc, path.clone()).await.unwrap().uuid]; + handles.push( + add_gobgp_path( + &svc, + api::Path { + identifier: 2, + ..path + }, + ) + .await + .unwrap() + .uuid, + ); + let ipv6 = add_gobgp_path(&svc, gobgp_binary_path(Family::IPV6)) + .await + .unwrap() + .uuid; + let installed = svc.tables.collect_loc_rib_paths(Family::IPV4); + let local = installed[0].new_best().unwrap(); + let peer = Arc::new(table::Source::new( + "192.0.2.1".parse().unwrap(), + "192.0.2.2".parse().unwrap(), + 65002, + 65001, + Ipv4Addr::new(192, 0, 2, 1), + PeerRole::Ebgp, + )); + for source in [peer, table::Source::kernel()] { + svc.tables.insert_route( + source, + Family::IPV4, + packet::PathNlri { + path_id: 0, + nlri: installed[0].net.clone(), + }, + local.nexthop, + local.attr.clone(), + None, + 0, + ); + } + // This policy-filtered local path has no UUID and is absent from Loc-RIB. + svc.tables.shards[0].lock().unwrap().rtable.insert( + table::Source::local(), + Family::IPV4, + packet::Nlri::from_str("10.3.0.0/24").unwrap(), + 7, + local.nexthop, + local.attr.clone(), + None, + true, + false, + None, + 0, + ); + let request = api::DeletePathRequest { + family: Some(convert::family_to_api(Family::IPV4)), + ..Default::default() + }; + svc.delete_path(tonic::Request::new(request.clone())) + .await + .unwrap(); + assert!( + svc.tables.shards[0] + .lock() + .unwrap() + .rtable + .iter_reach(Family::IPV4) + .all(|r| !r.source.is_local()) + ); + assert_eq!( + svc.tables.collect_loc_rib_paths(Family::IPV4)[0] + .current_paths + .len(), + 2 + ); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV6).len(), 1); + for uuid in handles { + assert_eq!( + delete_gobgp_path(&svc, uuid).await.unwrap_err().code(), + tonic::Code::NotFound + ); + } + svc.delete_path(tonic::Request::new(request)).await.unwrap(); + svc.delete_path(tonic::Request::new(api::DeletePathRequest::default())) + .await + .unwrap(); + assert!(svc.tables.collect_loc_rib_paths(Family::IPV6).is_empty()); + assert_eq!( + delete_gobgp_path(&svc, ipv6).await.unwrap_err().code(), + tonic::Code::NotFound + ); + assert_eq!( + svc.tables.collect_loc_rib_paths(Family::IPV4)[0] + .current_paths + .len(), + 2 + ); + } + + #[tokio::test] + async fn delete_path_all_vrf_respects_family_and_label_scope() { + let svc = make_grpc_service(); + // Even with the same RD, each VRF's distinct export label scopes deletion. + for name in ["blue", "red"] { + svc.add_vrf(tonic::Request::new(make_vrf_req(name, 65000, 1, 65000, 1))) + .await + .unwrap(); + for family in [Family::IPV4, Family::IPV6] { + svc.add_path(tonic::Request::new(api::AddPathRequest { + table_type: api::TableType::Vrf as i32, + vrf_id: name.into(), + path: Some(gobgp_binary_path(family)), + })) + .await + .unwrap(); + } + } + let global = add_gobgp_path(&svc, gobgp_binary_path(Family::IPV4)) + .await + .unwrap() + .uuid; + let request = api::DeletePathRequest { + table_type: api::TableType::Vrf as i32, + vrf_id: "blue".into(), + family: Some(convert::family_to_api(Family::IPV4)), + ..Default::default() + }; + svc.delete_path(tonic::Request::new(request.clone())) + .await + .unwrap(); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV4_VPN).len(), 1); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV6_VPN).len(), 2); + let mut invalid = request.clone(); + invalid.family = Some(convert::family_to_api(Family::IPV4_FLOWSPEC)); + assert_eq!( + svc.delete_path(tonic::Request::new(invalid)) + .await + .unwrap_err() + .code(), + tonic::Code::InvalidArgument + ); + svc.delete_path(tonic::Request::new(api::DeletePathRequest { + family: None, + ..request + })) + .await + .unwrap(); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV4_VPN).len(), 1); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV6_VPN).len(), 1); + assert_eq!(svc.tables.collect_loc_rib_paths(Family::IPV4).len(), 1); + delete_gobgp_path(&svc, global).await.unwrap(); + } + fn ipv4_path_with_id(prefix: &str, prefix_len: u32, nexthop: &str, path_id: u32) -> api::Path { api::Path { identifier: path_id, diff --git a/packet/src/bgp.rs b/packet/src/bgp.rs index c619c26f..1a0580a2 100644 --- a/packet/src/bgp.rs +++ b/packet/src/bgp.rs @@ -581,6 +581,16 @@ impl Nlri { buf } + /// Decode exactly one wire-format NLRI, without an ADD-PATH identifier. + pub fn decode_from_bytes(family: Family, bytes: &[u8]) -> Result { + let mut reader = BgpReader::::new(bytes); + let nlri = Self::decode(family, &mut reader, bytes.len(), true)?; + if reader.remaining_len() != 0 { + return Err(Notification::UpdateMalformedAttributeList); + } + Ok(nlri) + } + // Add a new match arm here when introducing a new SAFI. fn decode( family: Family, @@ -1214,6 +1224,36 @@ impl Attribute { self.flags } + /// Decode one complete BGP attribute, using four-octet AS numbers as in + /// GoBGP's binary gRPC representation. Reject truncation and trailing data. + pub fn decode_from_bytes(bytes: &[u8]) -> Result { + let invalid = || Notification::UpdateMalformedAttributeList; + let mut reader = Cursor::new(bytes); + let flags = reader.read_u8().map_err(|_| invalid())?; + let code = reader.read_u8().map_err(|_| invalid())?; + // Match UPDATE parsing: known attributes must have their canonical + // Optional and Transitive bits. Extended Length and Partial may vary. + if let Some(expected) = Self::canonical_flags(code) + && (flags ^ expected) & (Self::FLAG_OPTIONAL | Self::FLAG_TRANSITIVE) != 0 + { + return Err(invalid()); + } + let len = if flags & Self::FLAG_EXTENDED != 0 { + reader.read_u16::().map_err(|_| invalid())? + } else { + reader.read_u8().map_err(|_| invalid())? as u16 + }; + if reader.position() as usize + len as usize != bytes.len() { + return Err(invalid()); + } + let attribute = + Self::decode(code, flags, &mut reader, len, false).map_err(|_| invalid())?; + if reader.position() as usize != bytes.len() { + return Err(invalid()); + } + Ok(attribute) + } + pub fn new_with_value(code: u8, val: u32) -> Option { Some(Attribute { flags: Self::canonical_flags(code)?, diff --git a/table/src/lib.rs b/table/src/lib.rs index 48af1321..3a91adcf 100644 --- a/table/src/lib.rs +++ b/table/src/lib.rs @@ -1434,7 +1434,12 @@ impl Table { // removes a stale path from a previous GR session (different Source Arc // but same peer) when the peer reconnects and sends a WITHDRAW. let Some(i) = dst.entry.iter().position(|e| { - e.path.source.remote_addr == source.remote_addr && e.remote_path_id == remote_id + e.path.source.remote_addr == source.remote_addr + && e.remote_path_id == remote_id + // Canonical local and kernel sources both use 0.0.0.0, but + // an API withdrawal must never remove a kernel-owned path. + && e.path.source.is_local() == source.is_local() + && e.path.source.is_kernel() == source.is_kernel() }) else { return (None, None); }; @@ -2342,6 +2347,38 @@ mod tests { Arc::new(Vec::new()) } + #[test] + fn remove_keeps_local_and_kernel_sources_distinct() { + for (owner, other) in [ + (Source::kernel(), Source::local()), + (Source::local(), Source::kernel()), + ] { + let mut rt = Table::new(0); + let net = nlri(10, 0, 0, 0, 24); + rt.insert( + owner.clone(), + Family::IPV4, + net.clone(), + 0, + nh(), + empty_attrs(), + None, + false, + false, + None, + 0, + ); + assert!( + rt.remove(other, Family::IPV4, net.clone(), 0, None) + .0 + .is_none() + ); + assert_eq!(rt.collect_loc_rib_paths(&Family::IPV4).len(), 1); + assert!(rt.remove(owner, Family::IPV4, net, 0, None).0.is_some()); + assert!(rt.collect_loc_rib_paths(&Family::IPV4).is_empty()); + } + } + fn attrs_with_local_pref(val: u32) -> Arc> { Arc::new(vec![ packet::Attribute::new_with_value(packet::Attribute::LOCAL_PREF, val).unwrap(), diff --git a/tests/e2e/README.md b/tests/e2e/README.md index 7fabfbd0..34698016 100644 --- a/tests/e2e/README.md +++ b/tests/e2e/README.md @@ -1,8 +1,8 @@ # RustyBGP End-to-End Tests -Each subdirectory is a self-contained test scenario. Every test starts its -own Docker Compose topology, runs assertions against live BGP sessions, and -tears everything down on exit. +Each subdirectory is a self-contained test scenario. The topology tests +start Docker Compose, run assertions against live BGP sessions, and tear +everything down on exit. The CLI-only gRPC path test runs locally. ## Prerequisites @@ -22,6 +22,20 @@ On aarch64 (Apple Silicon, ARM servers): rustup target add aarch64-unknown-linux-musl ``` +## CLI-only gRPC path test + +This test runs without Docker or a BGP listener. Build `rustybgpd` and +install the GoBGP v4 CLI, then run from the repository root: + +``` +cargo build -p rustybgpd +python3 tests/e2e/grpc-paths/run-test.py --gobgp /path/to/gobgp +``` + +Use `--rustybgpd /path/to/rustybgpd` to test a different build. The test +starts a temporary loopback gRPC server and covers IPv4, IPv6, FlowSpec, +path identifiers, VRF isolation, and family-scoped `rib del all`. + ## Running all tests locally ``` diff --git a/tests/e2e/grpc-paths/run-test.py b/tests/e2e/grpc-paths/run-test.py new file mode 100644 index 00000000..69798acc --- /dev/null +++ b/tests/e2e/grpc-paths/run-test.py @@ -0,0 +1,147 @@ +#!/usr/bin/env python3 +"""Exercise the real GoBGP CLI against a local RustyBGP gRPC server.""" + +import argparse +import json +import pathlib +import socket +import subprocess +import tempfile +import time + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--rustybgpd", default="target/debug/rustybgpd") + parser.add_argument("--gobgp", default="gobgp") + args = parser.parse_args() + with socket.socket() as listener: + listener.bind(("127.0.0.1", 0)) + port = listener.getsockname()[1] + + def cli(*command): + result = subprocess.run( + [args.gobgp, "-u", "127.0.0.1", "-p", str(port), *command], + capture_output=True, + text=True, + timeout=15, + ) + if result.returncode: + raise RuntimeError(f"{' '.join(command)}: {result.stderr or result.stdout}") + return result.stdout + + def rib(family, *scope): + return json.loads(cli("-j", *scope, "rib", "-a", family)) + + def check(condition, description): + if not condition: + raise AssertionError(description) + print(f"PASS: {description}", flush=True) + + with tempfile.TemporaryDirectory(prefix="rustybgp-grpc-paths-") as directory: + config = pathlib.Path(directory) / "rustybgp.toml" + config.write_text( + '[global.config]\nas = 65001\nrouter-id = "192.0.2.1"\nport = -1\n' + ) + with (pathlib.Path(directory) / "daemon.log").open("w+") as log: + daemon = subprocess.Popen( + [args.rustybgpd, "-f", str(config), "--api-hosts", f"127.0.0.1:{port}"], + stdout=log, + stderr=subprocess.STDOUT, + ) + try: + for _ in range(100): + if daemon.poll() is not None: + raise RuntimeError("rustybgpd exited during startup") + try: + with socket.create_connection(("127.0.0.1", port), timeout=0.1): + pass + break + except OSError: + time.sleep(0.1) + else: + raise RuntimeError("gRPC server did not become ready") + cli("global") + + for family, prefix in [ + ("ipv4", "198.51.100.42/32"), + ("ipv6", "2001:db8::42/128"), + ]: + cli("global", "rib", "add", prefix, "-a", family) + check( + len(rib(family, "global")) == 1, + f"{family} announcement installed", + ) + cli("global", "rib", "del", prefix, "-a", family) + cli("global", "rib", "del", prefix, "-a", family) + check(not rib(family, "global"), f"{family} rib del is idempotent") + + prefix = "198.51.100.42/32" + for identifier in [1, 2]: + cli("global", "rib", "add", prefix, "identifier", str(identifier)) + cli("global", "rib", "del", prefix, "identifier", "1") + destinations = rib("ipv4", "global") + check( + len(destinations) == 1 + and len(next(iter(destinations.values()))) == 1, + "rib del preserves the other path identifier", + ) + cli("global", "rib", "add", "198.51.100.43/32") + cli("global", "rib", "add", "2001:db8::42/128", "-a", "ipv6") + cli("global", "rib", "del", "all", "-a", "ipv4") + check( + not rib("ipv4", "global") and len(rib("ipv6", "global")) == 1, + "rib del all respects the address family", + ) + cli("global", "rib", "del", "all", "-a", "ipv6") + + for family in ["ipv4-flowspec", "ipv6-flowspec"]: + destination = ( + "198.51.100.42/32" + if family == "ipv4-flowspec" + else "2001:db8::42/128" + ) + rule = [ + "match", + "destination", + destination, + "protocol", + "udp", + "destination-port", + "==53", + "then", + "discard", + ] + cli("global", "rib", "add", "-a", family, *rule) + check(len(rib(family, "global")) == 1, f"{family} rule installed") + cli("global", "rib", "del", "-a", family, *rule) + check( + not rib(family, "global"), + f"{family} rib del withdraws the rule", + ) + + for name, rd in [("blue", "65001:1"), ("red", "65001:2")]: + cli("vrf", "add", name, "rd", rd, "rt", "both", rd) + cli("vrf", name, "rib", "add", prefix) + cli("vrf", "blue", "rib", "del", prefix) + check( + not rib("ipv4", "vrf", "blue") + and len(rib("ipv4", "vrf", "red")) == 1, + "VRF rib del preserves the other VRF", + ) + print("GoBGP CLI path deletion checks passed.", flush=True) + except Exception: + log.seek(0) + print(log.read(), flush=True) + raise + finally: + daemon.terminate() + try: + daemon.wait(timeout=5) + except subprocess.TimeoutExpired: + daemon.kill() + daemon.wait() + + +if __name__ == "__main__": + main()