diff --git a/.github/buildomat/jobs/test-ddm-sextet.sh b/.github/buildomat/jobs/test-ddm-sextet.sh new file mode 100755 index 000000000..eb54c8f8b --- /dev/null +++ b/.github/buildomat/jobs/test-ddm-sextet.sh @@ -0,0 +1,18 @@ +#!/bin/bash +#: +#: name = "test-ddm-sextet" +#: variety = "basic" +#: target = "helios-3.0" +#: rust_toolchain = "stable" +#: output_rules = [ +#: "/work/*.log", +#: ] + +source .github/buildomat/test-ddm-common.sh + +# +# sextet tests +# + +banner "sextet" +pfexec cargo test --release -p mg-tests test_external_peer_sextet -- --nocapture diff --git a/Cargo.lock b/Cargo.lock index ad01e50e1..f147d6570 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1452,6 +1452,7 @@ version = "0.1.0" dependencies = [ "anyhow", "clap", + "client-common", "colored", "ddm-admin-client", "ddm-api-types-versions", @@ -3958,6 +3959,8 @@ dependencies = [ "ddm-admin-client", "ddm-api-types-versions", "mg-common", + "oxnet", + "rand 0.10.2", "slog", "slog-async", "slog-envlogger", @@ -9469,7 +9472,7 @@ dependencies = [ [[package]] name = "ztest" version = "0.1.0" -source = "git+https://github.com/oxidecomputer/falcon?branch=main#021420c1e2f1ac9c66759d5df44979940a1ac483" +source = "git+https://github.com/oxidecomputer/falcon?branch=main#c7952be660468c17c546f1fbaa03cf50c8f904cd" dependencies = [ "anyhow", "libnet", diff --git a/ddm-admin-client/src/lib.rs b/ddm-admin-client/src/lib.rs index 77988121c..53d89ab90 100644 --- a/ddm-admin-client/src/lib.rs +++ b/ddm-admin-client/src/lib.rs @@ -23,5 +23,6 @@ progenitor::generate_api!( PeerInfo = ddm_api_types_versions::latest::db::PeerInfo, PeerStatus = ddm_api_types_versions::latest::db::PeerStatus, Duration = std::time::Duration, + ExternalPeers = ddm_api_types_versions::latest::external_peers::ExternalPeers, } ); diff --git a/ddm-api-types/src/external_peers.rs b/ddm-api-types/src/external_peers.rs new file mode 100644 index 000000000..24a75854f --- /dev/null +++ b/ddm-api-types/src/external_peers.rs @@ -0,0 +1,5 @@ +// This Source Code Form is subject to the terms of the Mozilla Public +// License, v. 2.0. If a copy of the MPL was not distributed with this +// file, You can obtain one at https://mozilla.org/MPL/2.0/. + +pub use ddm_api_types_versions::latest::external_peers::*; diff --git a/ddm-api-types/src/lib.rs b/ddm-api-types/src/lib.rs index 2c22d1240..c041b73d6 100644 --- a/ddm-api-types/src/lib.rs +++ b/ddm-api-types/src/lib.rs @@ -19,4 +19,5 @@ pub mod admin; pub mod db; pub mod exchange; +pub mod external_peers; pub mod net; diff --git a/ddm-api-types/versions/src/external_peers/external_peers.rs b/ddm-api-types/versions/src/external_peers/external_peers.rs new file mode 100644 index 000000000..fea719441 --- /dev/null +++ b/ddm-api-types/versions/src/external_peers/external_peers.rs @@ -0,0 +1,12 @@ +// This Source Code Form is subject to the terms of the Mozilla Public +// License, v. 2.0. If a copy of the MPL was not distributed with this +// file, You can obtain one at https://mozilla.org/MPL/2.0/. + +use schemars::JsonSchema; +use serde::{Deserialize, Serialize}; +use std::collections::BTreeSet; + +#[derive(Debug, Clone, Deserialize, Serialize, JsonSchema, PartialEq, Eq)] +pub struct ExternalPeers { + pub address_objects: BTreeSet, +} diff --git a/ddm-api-types/versions/src/external_peers/mod.rs b/ddm-api-types/versions/src/external_peers/mod.rs new file mode 100644 index 000000000..74788e4b1 --- /dev/null +++ b/ddm-api-types/versions/src/external_peers/mod.rs @@ -0,0 +1,5 @@ +// This Source Code Form is subject to the terms of the Mozilla Public +// License, v. 2.0. If a copy of the MPL was not distributed with this +// file, You can obtain one at https://mozilla.org/MPL/2.0/. + +pub mod external_peers; diff --git a/ddm-api-types/versions/src/latest.rs b/ddm-api-types/versions/src/latest.rs index a2c1b47f0..64fc60327 100644 --- a/ddm-api-types/versions/src/latest.rs +++ b/ddm-api-types/versions/src/latest.rs @@ -25,3 +25,7 @@ pub mod exchange { pub mod net { pub use crate::v1::net::TunnelOrigin; } + +pub mod external_peers { + pub use crate::v3::external_peers::ExternalPeers; +} diff --git a/ddm-api-types/versions/src/lib.rs b/ddm-api-types/versions/src/lib.rs index 7d88b9274..34abae1eb 100644 --- a/ddm-api-types/versions/src/lib.rs +++ b/ddm-api-types/versions/src/lib.rs @@ -34,3 +34,5 @@ pub mod latest; pub mod v1; #[path = "peer_durations/mod.rs"] pub mod v2; +#[path = "external_peers/mod.rs"] +pub mod v3; diff --git a/ddm-api/src/lib.rs b/ddm-api/src/lib.rs index 623b8376b..e7c73464f 100644 --- a/ddm-api/src/lib.rs +++ b/ddm-api/src/lib.rs @@ -26,6 +26,7 @@ api_versions!([ // | example for the next person. // v // (next_int, IDENT), + (3, EXTERNAL_PEERS), (2, PEER_DURATIONS), (1, INITIAL), ]); @@ -170,4 +171,23 @@ pub trait DdmAdminApi { async fn disable_stats( ctx: RequestContext, ) -> Result; + + #[endpoint { + method = PUT, + path = "/external_peers", + versions = VERSION_EXTERNAL_PEERS.., + }] + async fn set_external_peers( + ctx: RequestContext, + request: TypedBody, + ) -> Result; + + #[endpoint { + method = GET, + path = "/external_peers", + versions = VERSION_EXTERNAL_PEERS.., + }] + async fn get_external_peers( + ctx: RequestContext, + ) -> Result, HttpError>; } diff --git a/ddm-protocol/src/v3.rs b/ddm-protocol/src/v3.rs index 80cfcb925..ab4aed6e5 100644 --- a/ddm-protocol/src/v3.rs +++ b/ddm-protocol/src/v3.rs @@ -117,6 +117,17 @@ impl UnderlayUpdate { .collect(), } } + pub fn break_loops(&self, hostname: &String) -> Self { + Self { + announce: self + .announce + .iter() + .filter(|x| !x.path.contains(hostname)) + .cloned() + .collect(), + withdraw: self.withdraw.clone(), + } + } } impl From for Update { diff --git a/ddm/Cargo.toml b/ddm/Cargo.toml index 1168489a9..a66460942 100644 --- a/ddm/Cargo.toml +++ b/ddm/Cargo.toml @@ -36,15 +36,15 @@ oximeter-producer.workspace = true oxnet.workspace = true uuid.workspace = true ddm-api.workspace = true +dpd-client.workspace = true # illumos-only deps used by the routing state machine and platform sys layer. # Gated by the `backend` feature so stub builds (e.g. Linux test fixtures # running `ddmd` with `--api-only`) link cleanly. libnet = { workspace = true, optional = true } -dpd-client = { workspace = true, optional = true } opte-ioctl = { workspace = true, optional = true } oxide-vpc = { workspace = true, optional = true } [features] default = ["backend"] -backend = ["dep:libnet", "dep:dpd-client", "dep:opte-ioctl", "dep:oxide-vpc"] +backend = ["dep:libnet", "dep:opte-ioctl", "dep:oxide-vpc"] diff --git a/ddm/src/admin.rs b/ddm/src/admin.rs index 21d373f4d..9adc372ca 100644 --- a/ddm/src/admin.rs +++ b/ddm/src/admin.rs @@ -3,13 +3,21 @@ // file, You can obtain one at https://mozilla.org/MPL/2.0/. use crate::db::Db; -use crate::sm::{AdminEvent, Event, PrefixSet, SmContext}; +use crate::defaults::{ + DISCOVERY_READ_TIMEOUT, EXCHANGE_TCP_PORT, EXCHANGE_TIMEOUT, + EXPIRE_THRESHOLD, IP_ADDR_WAIT, SOLICIT_INTERVAL, +}; +use crate::sm::{ + AdminEvent, Event, InterfaceState, PrefixSet, SessionStats, SmContext, + StateMachine, +}; use camino::Utf8PathBuf; use ddm_api::DdmAdminApi; use ddm_api::ddm_admin_api_mod; use ddm_api_types::admin::{EnableStatsRequest, ExpirePathParams, PrefixMap}; -use ddm_api_types::db::{PeerInfo, TunnelRoute}; +use ddm_api_types::db::{PeerInfo, RouterKind, TunnelRoute}; use ddm_api_types::exchange::PathVector; +use ddm_api_types::external_peers::ExternalPeers; use ddm_api_types::net::TunnelOrigin; use dropshot::ApiDescription; use dropshot::ApiDescriptionBuildErrors; @@ -24,16 +32,18 @@ use dropshot::RequestContext; use dropshot::TypedBody; use mg_common::lock; use oxnet::Ipv6Net; -use slog::{Logger, error, info, o}; +use slog::{Logger, debug, error, info, o}; use slog_error_chain::InlineErrorChain; -use std::collections::{HashMap, HashSet}; -use std::net::{IpAddr, SocketAddr, SocketAddrV4, SocketAddrV6}; +use std::collections::{BTreeSet, HashMap, HashSet}; +use std::net::{IpAddr, Ipv6Addr, SocketAddr, SocketAddrV4, SocketAddrV6}; use std::sync::Arc; use std::sync::Mutex; use std::sync::atomic::{AtomicU64, Ordering}; -use std::sync::mpsc::Sender; +use std::sync::mpsc::{Sender, channel}; +use std::time::Duration; use tokio::spawn; use tokio::task::JoinHandle; +use uuid::Uuid; pub const DDM_STATS_PORT: u16 = 8001; @@ -47,14 +57,53 @@ pub struct RouterStats { #[derive(Clone)] pub struct HandlerContext { - pub event_channels: Vec>, pub db: Db, pub stats: Arc, pub peers: Vec, pub stats_handler: Arc>>>, + pub tunables: Tunables, + pub router_kind: RouterKind, + pub rack_id: Option, + pub sled_id: Option, + pub router_id: String, pub log: Logger, } +impl HandlerContext { + pub fn event_channels(&self) -> impl Iterator> { + self.peers.iter().map(|x| &x.tx) + } +} + +#[derive(Clone)] +pub struct Tunables { + pub solicit_interval: Duration, + pub expire_threshold: Duration, + pub discovery_read_timeout: Duration, + pub ip_addr_wait: Duration, + pub exchange_timeout: Duration, + pub dendrite: bool, + pub dpd_port: u16, + pub dpd_host: String, + pub exchange_tcp_port: u16, +} + +impl Default for Tunables { + fn default() -> Self { + Self { + solicit_interval: SOLICIT_INTERVAL, + expire_threshold: EXPIRE_THRESHOLD, + discovery_read_timeout: DISCOVERY_READ_TIMEOUT, + ip_addr_wait: IP_ADDR_WAIT, + exchange_timeout: EXCHANGE_TIMEOUT, + dpd_port: dpd_client::default_port(), + dpd_host: "localhost".into(), + exchange_tcp_port: EXCHANGE_TCP_PORT, + dendrite: true, + } + } +} + pub fn handler( addr: IpAddr, port: u16, @@ -159,7 +208,7 @@ impl DdmAdminApi for DdmAdminApiImpl { let addr = params.into_inner().addr; let ctx = lock!(ctx.context()); - for e in &ctx.event_channels { + for e in ctx.event_channels() { e.send(Event::Admin(AdminEvent::Expire(addr))) .map_err(|e| { HttpError::for_internal_error(format!( @@ -238,7 +287,7 @@ impl DdmAdminApi for DdmAdminApiImpl { .originate(&prefixes) .map_err(|e| HttpError::for_internal_error(e.to_string()))?; - for e in &ctx.event_channels { + for e in ctx.event_channels() { e.send(Event::Admin(AdminEvent::Announce(PrefixSet::Underlay( prefixes.clone(), )))) @@ -274,7 +323,7 @@ impl DdmAdminApi for DdmAdminApiImpl { .originate_tunnel(&endpoints) .map_err(|e| HttpError::for_internal_error(e.to_string()))?; - for e in &ctx.event_channels { + for e in ctx.event_channels() { e.send(Event::Admin(AdminEvent::Announce(PrefixSet::Tunnel( endpoints.clone(), )))) @@ -308,7 +357,7 @@ impl DdmAdminApi for DdmAdminApiImpl { .withdraw(&prefixes) .map_err(|e| HttpError::for_internal_error(e.to_string()))?; - for e in &ctx.event_channels { + for e in ctx.event_channels() { e.send(Event::Admin(AdminEvent::Withdraw(PrefixSet::Underlay( prefixes.clone(), )))) @@ -344,7 +393,7 @@ impl DdmAdminApi for DdmAdminApiImpl { .withdraw_tunnel(&endpoints) .map_err(|e| HttpError::for_internal_error(e.to_string()))?; - for e in &ctx.event_channels { + for e in ctx.event_channels() { e.send(Event::Admin(AdminEvent::Withdraw(PrefixSet::Tunnel( endpoints.clone(), )))) @@ -374,7 +423,7 @@ impl DdmAdminApi for DdmAdminApiImpl { ) -> Result { let ctx = lock!(ctx.context()); - for e in &ctx.event_channels { + for e in ctx.event_channels() { e.send(Event::Admin(AdminEvent::Sync)).map_err(|e| { HttpError::for_internal_error(format!("admin event send: {e}")) })?; @@ -388,9 +437,12 @@ impl DdmAdminApi for DdmAdminApiImpl { request: TypedBody, ) -> Result { let rq = request.into_inner(); - let ctx = lock!(ctx.context()); + let (jh, log) = { + let ctx = lock!(ctx.context()); + (ctx.stats_handler.clone(), ctx.log.clone()) + }; - let mut jh = lock!(ctx.stats_handler); + let mut jh = lock!(jh); if jh.is_none() { let hostname = hostname::get() .expect("failed to get hostname") @@ -399,12 +451,11 @@ impl DdmAdminApi for DdmAdminApiImpl { *jh = Some( crate::oxstats::start_server( DDM_STATS_PORT, - ctx.peers.clone(), - ctx.stats.clone(), + ctx.context().clone(), hostname, rq.rack_id, rq.sled_id, - ctx.log.clone(), + log, ) .map_err(|e| { HttpError::for_internal_error(format!( @@ -429,6 +480,128 @@ impl DdmAdminApi for DdmAdminApiImpl { Ok(HttpResponseUpdatedNoContent()) } + + async fn set_external_peers( + ctx: RequestContext, + request: TypedBody, + ) -> Result { + let mut ctx = lock!(ctx.context()); + + if ctx.router_kind != RouterKind::Transit { + return Err(HttpError::for_bad_request( + None, + "external peers only supported for transit routers".into(), + )); + } + + let rq = request.into_inner(); + + let current = ctx.db.get_external_peers(); + let to_create = rq.address_objects.difference(¤t); + let to_remove = current.difference(&rq.address_objects); + + info!(ctx.log, "peer change request"; + "requested" => ?rq.address_objects, + "to_create" => ?to_create, + "to_remove" => ?to_remove, + "current" => ?current, + ); + ctx.db.set_external_peers(rq.address_objects.clone()); + + for addr_obj in to_create.into_iter() { + let (tx, rx) = channel(); + + let config = crate::sm::Config { + solicit_interval: ctx.tunables.solicit_interval, + expire_threshold: ctx.tunables.expire_threshold, + discovery_read_timeout: ctx.tunables.discovery_read_timeout, + ip_addr_wait: ctx.tunables.ip_addr_wait, + exchange_timeout: ctx.tunables.exchange_timeout, + exchange_port: ctx.tunables.exchange_tcp_port, + aobj_name: addr_obj.clone(), + if_name: String::default(), // initialized in state machine + if_index: 0, // initialized in state machine + // External peers are only a thing for transit routers. + kind: RouterKind::Transit, + dpd: if ctx.tunables.dendrite { + Some(crate::sm::DpdConfig { + host: ctx.tunables.dpd_host.clone(), + port: ctx.tunables.dpd_port, + }) + } else { + None + }, + addr: Ipv6Addr::UNSPECIFIED, + rack_id: ctx.rack_id, + sled_id: ctx.sled_id, + }; + let sm_ctx = SmContext { + config, + db: ctx.db.clone(), + event_channels: ctx.event_channels().cloned().collect(), + tx: tx.clone(), + log: ctx.log.clone(), + router_id: ctx.router_id.clone(), + rt: Arc::new(tokio::runtime::Handle::current()), + iface: Arc::new(InterfaceState::external()), + stats: Arc::new(SessionStats::default()), + discovery_stop: None, + first_run: true, + }; + let mut sm = StateMachine { + ctx: sm_ctx.clone(), + rx: Some(rx), + }; + + sm.run().unwrap(); + + ctx.peers.push(sm_ctx.clone()); + } + + // Ensure our indices are unique and ordered. + let mut remove_idx = BTreeSet::default(); + + for aobj in to_remove.into_iter() { + for (i, p) in ctx.peers.iter().enumerate() { + if p.iface.external { + if &p.config.aobj_name == aobj { + let _ = p.tx.send(Event::Admin(AdminEvent::Shutdown)); + remove_idx.insert(i); + info!( + ctx.log, + "removing external peeer for address object {aobj}" + ); + } else { + debug!(ctx.log, "{aobj} != {}", p.config.aobj_name); + } + } + } + } + // remove peers back to front so we don't shift the order out from under + // ourselves for the indexes we just gathered. + for i in remove_idx.iter().rev() { + ctx.peers.remove(*i); + } + + Ok(HttpResponseUpdatedNoContent()) + } + + async fn get_external_peers( + ctx: RequestContext, + ) -> Result, HttpError> { + let ctx = lock!(ctx.context()); + + if ctx.router_kind != RouterKind::Transit { + return Err(HttpError::for_bad_request( + None, + "external peers only supported for transit routers".into(), + )); + } + + Ok(HttpResponseOk(ExternalPeers { + address_objects: ctx.db.get_external_peers(), + })) + } } pub fn api_description() diff --git a/ddm/src/db.rs b/ddm/src/db.rs index d72fd7b20..03d498847 100644 --- a/ddm/src/db.rs +++ b/ddm/src/db.rs @@ -4,12 +4,13 @@ use ddm_api_types::db::TunnelRoute; use ddm_api_types::net::TunnelOrigin; +use ddm_protocol::v3::PathVector; use mg_common::lock; use oxnet::{IpNet, Ipv6Net}; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; use slog::{Logger, error}; -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeSet, HashMap, HashSet}; use std::net::Ipv6Addr; use std::sync::{Arc, Mutex}; @@ -47,6 +48,7 @@ pub struct Db { pub struct DbData { pub imported: HashSet, pub imported_tunnel: HashSet, + pub external_peers: BTreeSet, } const _: () = { @@ -261,6 +263,14 @@ impl Db { } result } + + pub fn get_external_peers(&self) -> BTreeSet { + lock!(self.data).external_peers.clone() + } + + pub fn set_external_peers(&mut self, value: BTreeSet) { + lock!(self.data).external_peers = value; + } } #[derive( @@ -273,6 +283,15 @@ pub struct Route { pub path: Vec, } +impl From for PathVector { + fn from(val: Route) -> Self { + Self { + destination: val.destination, + path: val.path, + } + } +} + #[derive(Debug, Clone)] pub enum EffectiveTunnelRouteSet { /// The routes in the contained set are active with priority greater than diff --git a/ddm/src/defaults.rs b/ddm/src/defaults.rs new file mode 100644 index 000000000..ad1f8ba84 --- /dev/null +++ b/ddm/src/defaults.rs @@ -0,0 +1,19 @@ +// This Source Code Form is subject to the terms of the Mozilla Public +// License, v. 2.0. If a copy of the MPL was not distributed with this +// file, You can obtain one at https://mozilla.org/MPL/2.0/. + +use std::time::Duration; + +pub const SOLICIT_INTERVAL: Duration = Duration::from_millis(2000); +pub const EXPIRE_THRESHOLD: Duration = Duration::from_millis(5000); +pub const DISCOVERY_READ_TIMEOUT: Duration = Duration::from_millis(1000); +pub const IP_ADDR_WAIT: Duration = Duration::from_millis(1000); +pub const EXCHANGE_TIMEOUT: Duration = Duration::from_millis(3000); + +pub const EXCHANGE_TCP_PORT: u16 = 0xdddd; + +pub const fn millis_u64(d: Duration) -> u64 { + let x = d.as_millis(); + assert!(x <= u64::MAX as u128); + x as u64 +} diff --git a/ddm/src/discovery/runtime.rs b/ddm/src/discovery/runtime.rs index 2cd8c53fc..3ef839d50 100644 --- a/ddm/src/discovery/runtime.rs +++ b/ddm/src/discovery/runtime.rs @@ -18,12 +18,12 @@ use serde::{Deserialize, Serialize}; use slog::Logger; use socket2::{Domain, Protocol, SockAddr, Socket, Type}; use std::mem::MaybeUninit; -use std::net::{Ipv6Addr, SocketAddrV6}; +use std::net::{Ipv6Addr, Shutdown, SocketAddrV6}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::mpsc::Sender; use std::sync::{Arc, RwLock}; use std::thread::{sleep, spawn}; -use std::time::{Duration, Instant}; +use std::time::Instant; const DDM_MADDR: Ipv6Addr = Ipv6Addr::new(0xff02, 0, 0, 0, 0, 0, 0, 0xdd); const DDM_PORT: u16 = 0xddd; @@ -98,7 +98,7 @@ pub(crate) fn handler( iface: Arc, stats: Arc, log: Logger, -) -> Result<(), DiscoveryError> { +) -> Result, DiscoveryError> { // listening on 2 sockets, solicitations are sent to DDM_MADDR, but // advertisements are sent to the unicast source addresses of a // solicitation. Binding to a link-scoped multicast address is required for @@ -116,16 +116,12 @@ pub(crate) fn handler( mc.bind(&mc_sa)?; mc.join_multicast_v6(&DDM_MADDR, config.if_index)?; mc.set_multicast_loop_v6(false)?; - mc.set_read_timeout(Some(Duration::from_millis( - config.discovery_read_timeout, - )))?; + mc.set_read_timeout(Some(config.discovery_read_timeout))?; let uc_sa: SockAddr = SocketAddrV6::new(config.addr, DDM_PORT, 0, config.if_index).into(); uc.bind(&uc_sa)?; - uc.set_read_timeout(Some(Duration::from_millis( - config.discovery_read_timeout, - )))?; + uc.set_read_timeout(Some(config.discovery_read_timeout))?; let ctx = HandlerContext { mc_socket: Arc::new(mc), @@ -153,9 +149,9 @@ pub(crate) fn handler( stop.clone(), stats.clone(), )?; - expire(ctx, stop, stats.clone())?; + expire(ctx, stop.clone(), stats.clone())?; - Ok(()) + Ok(stop) } fn send_solicitations( @@ -165,13 +161,17 @@ fn send_solicitations( ) { spawn(move || { loop { + if stop.load(Ordering::Relaxed) { + inf!(ctx.log, ctx.config.if_name, "stopping solicitor"); + break; + } if let Err(e) = solicit(&ctx) { err!(ctx.log, ctx.config.if_name, "solicit failed: {}", e); stop.store(true, Ordering::Relaxed); break; } stats.solicitations_sent.fetch_add(1, Ordering::Relaxed); - sleep(Duration::from_millis(ctx.config.solicit_interval)); + sleep(ctx.config.solicit_interval); } }); } @@ -197,7 +197,7 @@ fn expire( }; if let Some(nbr) = &*guard { let dt = Instant::now().duration_since(nbr.last_seen); - if dt.as_millis() > u128::from(ctx.config.expire_threshold) { + if dt > ctx.config.expire_threshold { wrn!( &ctx.log, ctx.config.if_name, @@ -212,9 +212,7 @@ fn expire( ctx.log.clone(), &ctx.config.if_name, ); - } else if dt.as_millis() - > u128::from(ctx.config.solicit_interval) - { + } else if dt > ctx.config.solicit_interval { wrn!( &ctx.log, ctx.config.if_name, @@ -230,17 +228,22 @@ fn expire( // sockets by trying to listen on a unicast address that a socket // waiting to be dropped is already listening on. if stop.load(Ordering::Relaxed) { + inf!( + &ctx.log, + ctx.config.if_name, + "stopping discovery expiration thread", + ); let event = ctx.event.clone(); let log = ctx.log.clone(); let if_name = ctx.config.if_name.clone(); let wait = ctx.config.discovery_read_timeout; drop(ctx); // Ensure read handlers have registered the stop event. - sleep(Duration::from_millis(wait)); + sleep(wait); emit_solicit_fail(event, log, &if_name); break; } - sleep(Duration::from_millis(ctx.config.solicit_interval)); + sleep(ctx.config.solicit_interval); } }); Ok(()) @@ -258,6 +261,18 @@ fn listen( handle_msg(&ctx, msg, &addr, &stats); }; if stop.load(Ordering::Relaxed) { + inf!( + &ctx.log, + ctx.config.if_name, + "stopping discovery handler" + ); + if let Err(e) = s.shutdown(Shutdown::Both) { + wrn!( + &ctx.log, + ctx.config.if_name, + "failed to shut down discovery socket {e:?}", + ); + } break; } } diff --git a/ddm/src/exchange/runtime.rs b/ddm/src/exchange/runtime.rs index 653824e43..9af01e7df 100644 --- a/ddm/src/exchange/runtime.rs +++ b/ddm/src/exchange/runtime.rs @@ -11,7 +11,7 @@ use super::ExchangeError; use crate::db::{Route, effective_route_set}; use crate::discovery::Version; use crate::sm::{Config, Event, PeerEvent, SmContext}; -use crate::{dbg, err, inf, wrn}; +use crate::{err, inf, wrn}; use ddm_api_types::db::{RouterKind, TunnelRoute}; use ddm_protocol::{v2, v3}; use dropshot::ApiDescription; @@ -33,7 +33,7 @@ use slog::{Logger, o}; use std::collections::HashSet; use std::net::{Ipv6Addr, SocketAddrV6}; use std::sync::Arc; -use std::sync::atomic::Ordering; +use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; use tokio::sync::Mutex; use tokio::time::timeout; @@ -44,7 +44,6 @@ const UNIT_EXCHANGE_SERVER: &str = "exchange_server"; pub struct HandlerContext { ctx: SmContext, peer: Ipv6Addr, - log: Logger, } pub(crate) fn announce_underlay( @@ -151,25 +150,24 @@ fn do_pull_common( } pub(crate) fn pull( - ctx: SmContext, + ctx: &SmContext, addr: Ipv6Addr, version: Version, rt: Arc, - log: Logger, + stop: Arc, ) -> Result<(), ExchangeError> { let pr: v3::PullResponse = match version { - Version::V2 => do_pull_v2(&ctx, &addr, &rt)?.into(), - Version::V3 => do_pull(&ctx, &addr, &rt)?, + Version::V2 => do_pull_v2(ctx, &addr, &rt)?.into(), + Version::V3 => do_pull(ctx, &addr, &rt)?, }; + if stop.load(Ordering::Relaxed) { + return Ok(()); + } + let update = v3::Update::announce(pr); - let hctx = HandlerContext { - ctx, - peer: addr, - log: log.clone(), - }; - handle_update(&update, &hctx); + handle_update(&update, ctx, addr); Ok(()) } @@ -244,9 +242,7 @@ fn send_update_common( let resp = client.request(req); rt.block_on(async move { - match timeout(Duration::from_millis(config.exchange_timeout), resp) - .await - { + match timeout(config.exchange_timeout, resp).await { Ok(_) => Ok(()), Err(e) => { err!( @@ -271,7 +267,6 @@ pub fn handler( ) -> Result, String> { let context = Arc::new(Mutex::new(HandlerContext { ctx: ctx.clone(), - log: log.clone(), peer, })); @@ -366,8 +361,9 @@ async fn push_handler_common( update: v3::Update, ) -> Result { let ctx = ctx.context().lock().await.clone(); + tokio::task::spawn_blocking(move || { - handle_update(&update, &ctx); + handle_update(&update, &ctx.ctx, ctx.peer); }) .await .map_err(|e| { @@ -400,7 +396,7 @@ async fn pull_handler_v2( destination: route.destination, path: route.path.clone(), }; - pv.path.push(ctx.ctx.hostname.clone()); + pv.path.push(ctx.ctx.router_id.clone()); underlay.insert(pv); } for route in &ctx.ctx.db.imported_tunnel() { @@ -419,7 +415,7 @@ async fn pull_handler_v2( for prefix in &originated { let pv = v3::PathVector { destination: *prefix, - path: vec![ctx.ctx.hostname.clone()], + path: vec![ctx.ctx.router_id.clone()], }; underlay.insert(pv); } @@ -476,7 +472,7 @@ async fn pull_handler( destination: route.destination, path: route.path.clone(), }; - pv.path.push(ctx.ctx.hostname.clone()); + pv.path.push(ctx.ctx.router_id.clone()); underlay.insert(pv); } for route in &ctx.ctx.db.imported_tunnel() { @@ -495,7 +491,7 @@ async fn pull_handler( for prefix in &originated { let pv = v3::PathVector { destination: *prefix, - path: vec![ctx.ctx.hostname.clone()], + path: vec![ctx.ctx.router_id.clone()], }; underlay.insert(pv); } @@ -529,50 +525,40 @@ async fn pull_handler( })) } -fn handle_update(update: &v3::Update, ctx: &HandlerContext) { - ctx.ctx - .stats - .updates_received - .fetch_add(1, Ordering::Relaxed); +fn handle_update(update: &v3::Update, ctx: &SmContext, peer_addr: Ipv6Addr) { + ctx.stats.updates_received.fetch_add(1, Ordering::Relaxed); if let Some(underlay_update) = &update.underlay { - handle_underlay_update(underlay_update, ctx); + handle_underlay_update(underlay_update, ctx, peer_addr); } if let Some(tunnel_update) = &update.tunnel { - handle_tunnel_update(tunnel_update, ctx); + handle_tunnel_update(tunnel_update, ctx, peer_addr); } // distribute updates - if ctx.ctx.config.kind == RouterKind::Transit { - dbg!( + if ctx.config.kind == RouterKind::Transit + && let Err(e) = ctx + .tx + .send(Event::Peer(PeerEvent::Redistribute(update.clone()))) + { + wrn!( ctx.log, - ctx.ctx.config.if_name, - "redistributing update to {} peers", - ctx.ctx.event_channels.len() + ctx.config.if_name, + "failed to send update to SM: {e:?}" ); - - let underlay = update - .underlay - .as_ref() - .map(|update| update.with_path_element(ctx.ctx.hostname.clone())); - - let push = v3::Update { - underlay, - tunnel: update.tunnel.clone(), - }; - - for ec in &ctx.ctx.event_channels { - ec.send(Event::Peer(PeerEvent::Push(push.clone()))).unwrap(); - } } } -fn handle_tunnel_update(update: &v3::TunnelUpdate, ctx: &HandlerContext) { +fn handle_tunnel_update( + update: &v3::TunnelUpdate, + ctx: &SmContext, + peer_addr: Ipv6Addr, +) { let mut import = HashSet::new(); let mut remove = HashSet::new(); - let db = &ctx.ctx.db; + let db = &ctx.db; let before = effective_route_set(&db.imported_tunnel()); @@ -584,7 +570,7 @@ fn handle_tunnel_update(update: &v3::TunnelUpdate, ctx: &HandlerContext) { vni: x.vni, metric: x.metric, }, - nexthop: ctx.peer, + nexthop: peer_addr, }); } db.import_tunnel(&import); @@ -597,7 +583,7 @@ fn handle_tunnel_update(update: &v3::TunnelUpdate, ctx: &HandlerContext) { vni: x.vni, metric: x.metric, }, - nexthop: ctx.peer, + nexthop: peer_addr, }); } db.delete_import_tunnel(&remove); @@ -607,72 +593,71 @@ fn handle_tunnel_update(update: &v3::TunnelUpdate, ctx: &HandlerContext) { let to_add = after.difference(&before).copied().collect(); let to_del = before.difference(&after).copied().collect(); - if let Err(e) = crate::sys::add_tunnel_routes( - &ctx.log, - &ctx.ctx.config.if_name, - &to_add, - ) { + if let Err(e) = + crate::sys::add_tunnel_routes(&ctx.log, &ctx.config.if_name, &to_add) + { err!( ctx.log, - ctx.ctx.config.if_name, + ctx.config.if_name, "add tunnel routes: {e}: {:#?}", import, ) } - if let Err(e) = crate::sys::remove_tunnel_routes( - &ctx.log, - &ctx.ctx.config.if_name, - &to_del, - ) { + if let Err(e) = + crate::sys::remove_tunnel_routes(&ctx.log, &ctx.config.if_name, &to_del) + { err!( ctx.log, - ctx.ctx.config.if_name, + ctx.config.if_name, "remove tunnel routes: {e}: {:#?}", import, ) } - ctx.ctx - .stats + ctx.stats .imported_underlay_prefixes - .store(ctx.ctx.db.imported_tunnel_count() as u64, Ordering::Relaxed); + .store(ctx.db.imported_tunnel_count() as u64, Ordering::Relaxed); } -fn handle_underlay_update(update: &v3::UnderlayUpdate, ctx: &HandlerContext) { +fn handle_underlay_update( + update: &v3::UnderlayUpdate, + ctx: &SmContext, + peer_addr: Ipv6Addr, +) { let mut import = HashSet::new(); let mut add = Vec::new(); - let db = &ctx.ctx.db; + let db = &ctx.db; for prefix in &update.announce { + // Skip announcements with ourselves in the path e.g. path vector + // loop breaking. + if prefix.path.contains(&ctx.router_id) { + continue; + } import.insert(Route { destination: prefix.destination, - nexthop: ctx.peer, - ifname: ctx.ctx.config.if_name.clone(), + nexthop: peer_addr, + ifname: ctx.config.if_name.clone(), path: prefix.path.clone(), }); let mut r = crate::sys::Route::new( prefix.destination.addr().into(), prefix.destination.width(), - ctx.peer.into(), + peer_addr.into(), ); - r.ifname.clone_from(&ctx.ctx.config.if_name); + r.ifname.clone_from(&ctx.config.if_name); add.push(r); } db.import(&import); - crate::sys::add_underlay_routes( - &ctx.log, - &ctx.ctx.config, - add, - &ctx.ctx.rt, - ); + crate::sys::add_underlay_routes(&ctx.log, &ctx.config, add, &ctx.rt); let mut withdraw = HashSet::new(); for prefix in &update.withdraw { withdraw.insert(Route { destination: prefix.destination, - nexthop: ctx.peer, - ifname: ctx.ctx.config.if_name.clone(), + nexthop: peer_addr, + ifname: ctx.config.if_name.clone(), path: prefix.path.clone(), }); } @@ -694,20 +679,19 @@ fn handle_underlay_update(update: &v3::UnderlayUpdate, ctx: &HandlerContext) { w.destination.width(), w.nexthop.into(), ); - r.ifname.clone_from(&ctx.ctx.config.if_name); + r.ifname.clone_from(&ctx.config.if_name); del.push(r); } } crate::sys::remove_underlay_routes( &ctx.log, - &ctx.ctx.config.if_name, - &ctx.ctx.config.dpd, + &ctx.config.if_name, + &ctx.config.dpd, del, - &ctx.ctx.rt, + &ctx.rt, ); - ctx.ctx - .stats + ctx.stats .imported_underlay_prefixes - .store(ctx.ctx.db.imported_count() as u64, Ordering::Relaxed); + .store(ctx.db.imported_count() as u64, Ordering::Relaxed); } diff --git a/ddm/src/lib.rs b/ddm/src/lib.rs index 6a5e2a68a..19326ec0a 100644 --- a/ddm/src/lib.rs +++ b/ddm/src/lib.rs @@ -4,6 +4,7 @@ pub mod admin; pub mod db; +pub mod defaults; pub mod discovery; pub mod exchange; pub mod oxstats; diff --git a/ddm/src/oxstats.rs b/ddm/src/oxstats.rs index a641d728f..a1a5babfa 100644 --- a/ddm/src/oxstats.rs +++ b/ddm/src/oxstats.rs @@ -2,7 +2,7 @@ // License, v. 2.0. If a copy of the MPL was not distributed with this // file, You can obtain one at https://mozilla.org/MPL/2.0/. -use crate::{admin::RouterStats, sm::SmContext}; +use crate::admin::HandlerContext; use chrono::{DateTime, Utc}; use mg_common::{ lock, @@ -15,7 +15,7 @@ use oximeter::{ }; use oximeter_producer::{ConfigLogging, ConfigLoggingLevel, LogConfig}; use slog::Logger; -use std::sync::atomic::Ordering; +use std::sync::{Mutex, atomic::Ordering}; use std::{net::SocketAddr, sync::Arc, time::Duration}; use tokio::task::JoinHandle; use uuid::Uuid; @@ -46,8 +46,7 @@ pub(crate) struct Stats { hostname: String, rack_id: Uuid, sled_id: Uuid, - peers: Vec, - router_stats: Arc, + ctx: Arc>, } macro_rules! ddm_session_counter { @@ -135,12 +134,14 @@ impl Producer for Stats { // level stats. let mut samples: Vec = Vec::with_capacity(2 + 13); + let ctx = lock!(self.ctx); + samples.push(ddm_router_quantity!( self.hostname.clone().into(), self.rack_id, self.sled_id, OriginatedUnderlayPrefixes, - self.router_stats.originated_underlay_prefixes + ctx.stats.originated_underlay_prefixes )); samples.push(ddm_router_quantity!( @@ -148,10 +149,10 @@ impl Producer for Stats { self.rack_id, self.sled_id, OriginatedTunnelEndpoints, - self.router_stats.originated_tunnel_endpoints + ctx.stats.originated_tunnel_endpoints )); - for peer in &self.peers { + for peer in &ctx.peers { let if_name = lock!(peer.iface.if_name).clone(); samples.push(ddm_session_counter!( self.start_time, @@ -268,8 +269,7 @@ impl Producer for Stats { #[allow(clippy::too_many_arguments)] pub fn start_server( port: u16, - peers: Vec, - router_stats: Arc, + ctx: Arc>, hostname: String, rack_id: Uuid, sled_id: Uuid, @@ -284,11 +284,10 @@ pub fn start_server( let stats_producer = Stats { start_time: chrono::offset::Utc::now(), - peers, + ctx, hostname, rack_id, sled_id, - router_stats, }; registry.register_producer(stats_producer).unwrap(); diff --git a/ddm/src/sm/mod.rs b/ddm/src/sm/mod.rs index 4a52efb55..8d3fdc60f 100644 --- a/ddm/src/sm/mod.rs +++ b/ddm/src/sm/mod.rs @@ -16,16 +16,20 @@ use oxnet::Ipv6Net; use slog::Logger; use std::collections::HashSet; use std::net::Ipv6Addr; -use std::sync::atomic::AtomicU64; +use std::sync::atomic::{AtomicBool, AtomicU64}; use std::sync::mpsc::{Receiver, Sender}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use thiserror::Error; +use uuid::Uuid; + +#[cfg(target_os = "illumos")] +use std::collections::BTreeSet; #[cfg(all(feature = "backend", target_os = "illumos"))] -mod state; +pub(crate) mod state; -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum AdminEvent { /// Announce a set of IPv6 prefixes Announce(PrefixSet), @@ -38,27 +42,50 @@ pub enum AdminEvent { /// Synchronize with active peers by pulling their prefixes. Sync, + + /// A new external peer has been added to the router that can be reached + /// using the provided sender. + NewExternalPeer(Sender), + + /// Shutdown on receipt of this event. + Shutdown, } -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum PrefixSet { Underlay(HashSet), Tunnel(HashSet), } -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum PeerEvent { + /// Upon reception of this event, a state machine is to push the update to + /// its peer. Push(ddm_protocol::v3::Update), + + /// Upon reception of this event, a state machine is to redistribute the + /// update to it's sibling routers through it's event channels. The state + /// machine is responsible for maintaining the path vector and performing + /// loop breaking. + Redistribute(ddm_protocol::v3::Update), } -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum NeighborEvent { + /// An event sent from the discovery subsystem to the state machine letting + /// it know the link local ipv6 address of the peer and it's version. Advertise((Ipv6Addr, Version)), + + /// An event sent from the discovery subsystem to the state machine letting + /// it know that solicitation has failed. SolicitFail, + + /// An event sent from the discovery subsystem to the state machine letting + /// it know that the peer has expired. Expire, } -#[derive(Debug)] +#[derive(Debug, Clone)] pub enum Event { Neighbor(NeighborEvent), Peer(PeerEvent), @@ -124,22 +151,21 @@ pub struct Config { /// Link local Ipv6 address this state machine is associated with pub addr: Ipv6Addr, - /// How long to wait between solicitations (milliseconds). - pub solicit_interval: u64, + /// How long to wait between solicitations. + pub solicit_interval: Duration, /// How often to check for link failure while waiting for discovery messges. - pub discovery_read_timeout: u64, + pub discovery_read_timeout: Duration, /// How long to wait between attempts to get an IP address for a specified /// address object. - pub ip_addr_wait: u64, + pub ip_addr_wait: Duration, /// How long to wait without a solicitation response before expiring a peer - /// (milliseconds). - pub expire_threshold: u64, + pub expire_threshold: Duration, /// How long to wait for a response to exchange messages. - pub exchange_timeout: u64, + pub exchange_timeout: Duration, /// The kind of router this is, server or transit. pub kind: RouterKind, @@ -149,6 +175,12 @@ pub struct Config { /// Dendrite dpd config pub dpd: Option, + + /// Rack ID + pub rack_id: Option, + + /// Sled ID + pub sled_id: Option, } #[derive(Clone)] @@ -184,12 +216,20 @@ pub struct PeerIdentity { pub struct InterfaceState { pub if_index: Mutex, pub if_name: Mutex, + pub external: bool, pub fsm_state: Mutex, pub last_fsm_state_change: Mutex, pub peer_identity: Mutex>, } impl InterfaceState { + pub fn external() -> Self { + Self { + external: true, + ..Default::default() + } + } + pub fn transition(&self, state: FsmState) { *lock!(self.fsm_state) = state; *lock!(self.last_fsm_state_change) = Instant::now(); @@ -215,6 +255,7 @@ impl Default for InterfaceState { Self { if_index: Mutex::new(0), if_name: Mutex::new(String::new()), + external: false, fsm_state: Mutex::new(FsmState::Init), last_fsm_state_change: Mutex::new(Instant::now()), peer_identity: Mutex::new(None), @@ -248,13 +289,42 @@ pub struct SmContext { pub tx: Sender, pub event_channels: Vec>, pub rt: Arc, - pub hostname: String, + pub router_id: String, pub iface: Arc, pub stats: Arc, pub log: Logger, + pub discovery_stop: Option>, + pub first_run: bool, } pub struct StateMachine { pub ctx: SmContext, pub rx: Option>, } + +#[cfg(not(target_os = "illumos"))] +impl StateMachine { + pub fn run(&mut self) -> Result<(), SmError> { + Ok(()) + } +} + +/// Send an event to all channels in the list, removing any channels from the +/// list that are dead. +#[cfg(target_os = "illumos")] +pub(crate) fn send(e: Event, event_channels: &mut Vec>) { + // Ensure our indices are unique and ordered. + let mut dead_channels = BTreeSet::default(); + for (i, c) in event_channels.iter().enumerate() { + if c.send(e.clone()).is_err() { + dead_channels.insert(i); + } + } + // we need to remove in descending order, so we don't remove `i` and then + // try to remove `i+1` later which wlll be a _different_ item than we grabbed + // the index for. Removing from the top down causes no shifting for subsequent + // index removals. + for i in dead_channels.iter().rev() { + event_channels.remove(*i); + } +} diff --git a/ddm/src/sm/state.rs b/ddm/src/sm/state.rs index ecb8ffa68..67e004b00 100644 --- a/ddm/src/sm/state.rs +++ b/ddm/src/sm/state.rs @@ -21,7 +21,7 @@ use std::net::IpAddr; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::mpsc::Receiver; -use std::thread::{sleep, spawn}; +use std::thread::{JoinHandle, sleep, spawn}; use std::time::Duration; use crate::discovery::Version; @@ -34,10 +34,13 @@ impl StateMachine { let mut rx = self.rx.take().unwrap(); let log = self.ctx.log.clone(); spawn(move || { - let mut state: Box = - Box::new(Init::new(ctx.clone(), log.clone())); + let mut state: Option> = + Some(Box::new(Init::new(ctx.clone(), log.clone()))); loop { - (state, rx) = state.run(rx); + (state, rx) = match &mut state { + Some(st) => st.run(rx), + None => break, + } } }); @@ -49,7 +52,7 @@ trait State { fn run( &mut self, event: Receiver, - ) -> (Box, Receiver); + ) -> (Option>, Receiver); } struct Init { @@ -67,10 +70,28 @@ impl State for Init { fn run( &mut self, event: Receiver, - ) -> (Box, Receiver) { + ) -> (Option>, Receiver) { self.ctx.iface.transition(FsmState::Init); self.ctx.iface.clear_peer(); loop { + // Check for shutdown or new peers + while let Ok(e) = event.try_recv() { + match e { + Event::Admin(AdminEvent::Shutdown) => { + return (None, event); + } + Event::Admin(AdminEvent::NewExternalPeer(tx)) => { + self.ctx.event_channels.push(tx); + } + _ => { + wrn!( + self.log, + self.ctx.config.aobj_name, + "event unhandled in init state: {e:?}" + ); + } + } + } let info = match get_ipaddr_info(&self.ctx.config.aobj_name) { Ok(info) => info, Err(e) => { @@ -81,7 +102,7 @@ impl State for Init { &self.ctx.config.aobj_name, e ); - sleep(Duration::from_millis(self.ctx.config.ip_addr_wait)); + sleep(self.ctx.config.ip_addr_wait); continue; } }; @@ -94,7 +115,7 @@ impl State for Init { "specified address {} is not IPv6", &self.ctx.config.aobj_name ); - sleep(Duration::from_millis(self.ctx.config.ip_addr_wait)); + sleep(self.ctx.config.ip_addr_wait); continue; } }; @@ -115,17 +136,31 @@ impl State for Init { // Now that we have an ip address to run discovery on, start the // discovery handler and jump into the solicit state. - discovery::handler( - self.ctx.hostname.clone(), + let discovery_stop = match discovery::handler( + self.ctx.router_id.clone(), self.ctx.config.clone(), self.ctx.tx.clone(), self.ctx.iface.clone(), self.ctx.stats.clone(), self.ctx.log.clone(), - ) - .unwrap(); // TODO unwrap + ) { + Ok(stop) => stop, + Err(e) => { + wrn!( + self.log, + self.ctx.config.if_name, + "failed to start discovery handler: {e}", + ); + sleep(self.ctx.config.solicit_interval); + continue; + } + }; + self.ctx.discovery_stop = Some(discovery_stop); return ( - Box::new(Solicit::new(self.ctx.clone(), self.log.clone())), + Some(Box::new(Solicit::new( + self.ctx.clone(), + self.log.clone(), + ))), event, ); } @@ -147,7 +182,7 @@ impl State for Solicit { fn run( &mut self, event: Receiver, - ) -> (Box, Receiver) { + ) -> (Option>, Receiver) { self.ctx.iface.transition(FsmState::Solicit); loop { let e = match event.recv() { @@ -170,12 +205,12 @@ impl State for Solicit { "transition solicit -> exchange" ); return ( - Box::new(Exchange::new( + Some(Box::new(Exchange::new( self.ctx.clone(), addr, version, self.log.clone(), - )), + ))), event, ); } @@ -187,7 +222,10 @@ impl State for Solicit { "exiting solicit state due to failed solicit", ); return ( - Box::new(Init::new(self.ctx.clone(), self.log.clone())), + Some(Box::new(Init::new( + self.ctx.clone(), + self.log.clone(), + ))), event, ); } @@ -199,6 +237,15 @@ impl State for Solicit { e ); } + Event::Admin(AdminEvent::NewExternalPeer(tx)) => { + self.ctx.event_channels.push(tx); + } + Event::Admin(AdminEvent::Shutdown) => { + if let Some(x) = self.ctx.discovery_stop.as_mut() { + x.store(true, Ordering::Relaxed); + } + return (None, event); + } Event::Admin(e) => { wrn!( self.log, @@ -234,30 +281,30 @@ impl Exchange { } } - fn initial_pull(&self, stop: Arc) { - let ctx = self.ctx.clone(); + fn initial_pull(&mut self, stop: Arc) -> JoinHandle<()> { let peer = self.peer; let version = self.version; let rt = self.ctx.rt.clone(); let log = self.log.clone(); let interval = self.ctx.config.solicit_interval; let if_name = self.ctx.config.if_name.clone(); + let ctx = self.ctx.clone(); spawn(move || { while let Err(e) = crate::exchange::pull( - ctx.clone(), + &ctx, peer, version, rt.clone(), - log.clone(), + stop.clone(), ) { - sleep(Duration::from_millis(interval)); + sleep(interval); wrn!(log, if_name, "exchange pull: {}", e); if stop.load(Ordering::Relaxed) { break; } } - }); + }) } fn wait_for_exchange_server_to_start(&self) { @@ -295,10 +342,21 @@ impl Exchange { &mut self, exchange_thread: &tokio::task::JoinHandle<()>, pull_stop: &AtomicBool, + initial_pull: JoinHandle<()>, ) { + // Stop other threads that may update the rib after we clear for expiry. + pull_stop.store(true, Ordering::Relaxed); + if let Err(e) = initial_pull.join() { + err!( + self.log, + self.ctx.config.if_name, + "failed to join initial pull thread: {e:?}", + ); + } + exchange_thread.abort(); self.ctx.iface.clear_peer(); - let (to_remove, to_remove_tnl) = + let (mut to_remove, to_remove_tnl) = self.ctx.db.remove_nexthop_routes(self.peer); let mut routes: Vec = Vec::new(); for x in &to_remove { @@ -335,6 +393,26 @@ impl Exchange { self.ctx.event_channels.len() ); + // Only send withdraws for expirations that result in a total loss + // of reachability to a destination for the given path. + let imported = self.ctx.db.imported(); + dbg!(self.log, self.ctx.config.if_name, "imported: {imported:#?}"); + dbg!( + self.log, + self.ctx.config.if_name, + "to_remove: {to_remove:#?}" + ); + to_remove.retain(|x| { + !imported + .iter() + .any(|y| y.destination == x.destination && y.path == x.path) + }); + dbg!( + self.log, + self.ctx.config.if_name, + "to_remove (retained): {to_remove:#?}" + ); + let underlay = if to_remove.is_empty() { None } else { @@ -345,7 +423,7 @@ impl Exchange { destination: x.destination, path: { let mut ps = x.path.clone(); - ps.push(self.ctx.hostname.clone()); + ps.push(self.ctx.router_id.clone()); ps }, }) @@ -362,11 +440,11 @@ impl Exchange { }; let push = Update { underlay, tunnel }; - for ec in &self.ctx.event_channels { - ec.send(Event::Peer(PeerEvent::Push(push.clone()))).unwrap(); - } + super::send( + Event::Peer(PeerEvent::Push(push.clone())), + &mut self.ctx.event_channels, + ); } - pull_stop.store(true, Ordering::Relaxed); } } @@ -374,7 +452,7 @@ impl State for Exchange { fn run( &mut self, event: Receiver, - ) -> (Box, Receiver) { + ) -> (Option>, Receiver) { self.ctx.iface.transition(FsmState::Exchange); let exchange_thread = loop { match exchange::handler( @@ -404,7 +482,15 @@ impl State for Exchange { // Do an initial pull, in the event that exchange events are fired while // this pull is taking place, they will be queued and handled in the // loop below. - self.initial_pull(pull_stop.clone()); + let initial_pull = self.initial_pull(pull_stop.clone()); + + if self.ctx.iface.external && self.ctx.first_run { + self.ctx.first_run = false; + crate::sm::send( + Event::Admin(AdminEvent::NewExternalPeer(self.ctx.tx.clone())), + &mut self.ctx.event_channels, + ); + } loop { let e = match event.recv() { @@ -427,7 +513,7 @@ impl State for Exchange { .iter() .map(|x| PathVector { destination: *x, - path: vec![self.ctx.hostname.clone()], + path: vec![self.ctx.router_id.clone()], }) .collect(); if let Err(e) = crate::exchange::announce_underlay( @@ -451,12 +537,16 @@ impl State for Exchange { "expiring peer {} due to failed announce", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } @@ -486,12 +576,16 @@ impl State for Exchange { "expiring peer {} due to failed tunnel announce", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } @@ -503,7 +597,7 @@ impl State for Exchange { .iter() .map(|x| PathVector { destination: *x, - path: vec![self.ctx.hostname.clone()], + path: vec![self.ctx.router_id.clone()], }) .collect(); if let Err(e) = crate::exchange::withdraw_underlay( @@ -527,12 +621,16 @@ impl State for Exchange { "expiring peer {} due to failed withdraw", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } @@ -562,12 +660,16 @@ impl State for Exchange { "expiring peer {} due to failed tunnel withdraw", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } @@ -580,23 +682,28 @@ impl State for Exchange { "administratively expiring peer {}", peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } } Event::Admin(AdminEvent::Sync) => { + let rt = self.ctx.rt.clone(); if let Err(e) = crate::exchange::pull( - self.ctx.clone(), + &self.ctx, self.peer, self.version, - self.ctx.rt.clone(), - self.log.clone(), + rt, + pull_stop.clone(), ) { err!( self.log, @@ -606,6 +713,40 @@ impl State for Exchange { ); } } + Event::Admin(AdminEvent::NewExternalPeer(tx)) => { + if tx + .send(Event::Peer(PeerEvent::Push(Update { + underlay: Some( + UnderlayUpdate { + announce: self + .ctx + .db + .imported() + .into_iter() + .map(Into::into) + .collect(), + withdraw: HashSet::default(), + } + .with_path_element(self.ctx.router_id.clone()), + ), + tunnel: None, + }))) + .is_ok() + { + self.ctx.event_channels.push(tx.clone()); + } + } + Event::Admin(AdminEvent::Shutdown) => { + if let Some(x) = self.ctx.discovery_stop.as_mut() { + x.store(true, Ordering::Relaxed); + } + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); + return (None, event); + } Event::Peer(PeerEvent::Push(update)) => { inf!( self.log, @@ -638,12 +779,16 @@ impl State for Exchange { "expiring peer {} due to failed announce", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } @@ -670,17 +815,60 @@ impl State for Exchange { "expiring peer {} due to failed withdraw", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } } } + Event::Peer(PeerEvent::Redistribute(mut update)) => { + dbg!( + self.log, + self.ctx.config.if_name, + "redistributing update to {} peers", + self.ctx.event_channels.len() + ); + update.underlay = update.underlay.map(|u| { + let sans_loops = u.break_loops(&self.ctx.router_id); + + let announce_loop = + u.announce.difference(&sans_loops.announce); + let withdraw_loop = + u.withdraw.difference(&sans_loops.withdraw); + + for x in announce_loop { + wrn!( + self.log, + self.ctx.config.if_name, + "loop detected: dropping announcement: {:?}", + x, + ) + } + for x in withdraw_loop { + wrn!( + self.log, + self.ctx.config.if_name, + "loop detected: dropping withdraw {:?}", + x, + ) + } + + sans_loops.with_path_element(self.ctx.router_id.clone()) + }); + super::send( + Event::Peer(PeerEvent::Push(update)), + &mut self.ctx.event_channels, + ); + } Event::Neighbor(NeighborEvent::Expire) => { wrn!( self.log, @@ -688,12 +876,16 @@ impl State for Exchange { "expiring peer {} due to discovery event", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Solicit::new( + Some(Box::new(Solicit::new( self.ctx.clone(), self.log.clone(), - )), + ))), event, ); } @@ -704,9 +896,16 @@ impl State for Exchange { "expiring peer {} due to failed solicit", self.peer, ); - self.expire_peer(&exchange_thread, &pull_stop); + self.expire_peer( + &exchange_thread, + &pull_stop, + initial_pull, + ); return ( - Box::new(Init::new(self.ctx.clone(), self.log.clone())), + Some(Box::new(Init::new( + self.ctx.clone(), + self.log.clone(), + ))), event, ); } diff --git a/ddm/src/sys.rs b/ddm/src/sys.rs index c251ab600..48f13b8dc 100644 --- a/ddm/src/sys.rs +++ b/ddm/src/sys.rs @@ -178,23 +178,38 @@ pub fn add_routes_dendrite( // TODO this is gross, use link type properties rather than futzing // around with strings. - let Some(egress_port_num) = ifname - .strip_prefix("tfportrear") - .and_then(|x| x.strip_suffix("_0")) - .map(|x| x.trim()) - .and_then(|x| x.parse::().ok()) - else { - err!(log, ifname, "expected tfportrear"); - continue; - }; - // TODO this assumes ddm only operates on rear ports, which will not be - // true for multi-rack deployments. - let port_name = format!("rear{}", egress_port_num); - let port_id = match types::Rear::try_from(&port_name) { - Ok(rear) => PortId::Rear(rear), - Err(e) => { - err!(log, ifname, "bad port name ({port_name}): {e}"); + let port_id = { + if let Some(egress_port_num) = ifname + .strip_prefix("tfportrear") + .and_then(|x| x.strip_suffix("_0")) + .map(|x| x.trim()) + .and_then(|x| x.parse::().ok()) + { + let port_name = format!("rear{}", egress_port_num); + match types::Rear::try_from(&port_name) { + Ok(rear) => PortId::Rear(rear), + Err(e) => { + err!(log, ifname, "bad port name ({port_name}): {e}"); + continue; + } + } + } else if let Some(egress_port_num) = ifname + .strip_prefix("tfportqsfp") + .and_then(|x| x.strip_suffix("_0")) + .map(|x| x.trim()) + .and_then(|x| x.parse::().ok()) + { + let port_name = format!("qsfp{}", egress_port_num); + match types::Qsfp::try_from(&port_name) { + Ok(qsfp) => PortId::Qsfp(qsfp), + Err(e) => { + err!(log, ifname, "bad port name ({port_name}): {e}"); + continue; + } + } + } else { + err!(log, ifname, "expected tfportrear or tfportqsfp"); continue; } }; diff --git a/ddmadm/Cargo.toml b/ddmadm/Cargo.toml index cde57b1c4..25341da3f 100644 --- a/ddmadm/Cargo.toml +++ b/ddmadm/Cargo.toml @@ -5,6 +5,7 @@ edition = "2024" [dependencies] mg-common = { path = "../mg-common" } +client-common = { path = "../client-common" } ddm-admin-client = { path = "../ddm-admin-client" } ddm-api-types-versions.workspace = true anyhow.workspace = true diff --git a/ddmadm/src/main.rs b/ddmadm/src/main.rs index 21d56fdf3..3cd230a46 100644 --- a/ddmadm/src/main.rs +++ b/ddmadm/src/main.rs @@ -4,9 +4,11 @@ use anyhow::Result; use clap::Parser; +use client_common::println_nopipe; use colored::*; use ddm_admin_client::Client; use ddm_api_types_versions::latest::db::PeerStatus; +use ddm_api_types_versions::latest::external_peers::ExternalPeers; use ddm_api_types_versions::latest::net as types; use mg_common::cli::oxide_cli_style; use mg_common::format_duration_human; @@ -65,6 +67,14 @@ enum SubCommand { /// Sync prefix information from peers. Sync, + + /// Set external peers as a list of address objects. + SetExternalPeers { + addr_obj: Vec, + }, + + // Get external peers + GetExternalPeers, } #[derive(Debug, Parser)] @@ -170,13 +180,16 @@ async fn run() -> Result<()> { for pv in &mut destinations { // show path from perspective of this node, e.g. nearest node // first - pv.path.reverse(); - let strpath = pv.path.join(" "); - writeln!( - &mut tw, - "{}\t{}\t{}", - pv.destination, nexthop, strpath, - )?; + if let Some(p) = pv.path.pop() { + writeln!( + &mut tw, + "{}\t{}\t{}", + pv.destination, nexthop, p, + )?; + } + for p in &pv.path { + writeln!(&mut tw, "\t\t{}", p,)?; + } } } tw.flush()?; @@ -265,6 +278,19 @@ async fn run() -> Result<()> { SubCommand::Sync => { client.sync().await?; } + SubCommand::SetExternalPeers { addr_obj } => { + client + .set_external_peers(&ExternalPeers { + address_objects: addr_obj.iter().cloned().collect(), + }) + .await?; + } + SubCommand::GetExternalPeers => { + let peers = client.get_external_peers().await?.into_inner(); + for p in peers.address_objects { + println_nopipe!("{p}"); + } + } } Ok(()) diff --git a/ddmd/src/main.rs b/ddmd/src/main.rs index 671e5cdd1..5cdb48816 100644 --- a/ddmd/src/main.rs +++ b/ddmd/src/main.rs @@ -4,8 +4,12 @@ use camino::Utf8PathBuf; use clap::Parser; -use ddm::admin::{HandlerContext, RouterStats}; +use ddm::admin::{HandlerContext, RouterStats, Tunables}; use ddm::db::Db; +use ddm::defaults::{ + DISCOVERY_READ_TIMEOUT, EXCHANGE_TCP_PORT, EXCHANGE_TIMEOUT, + EXPIRE_THRESHOLD, IP_ADDR_WAIT, SOLICIT_INTERVAL, millis_u64, +}; #[cfg(all(feature = "backend", target_os = "illumos"))] use ddm::sm::{DpdConfig, InterfaceState, SmContext, StateMachine}; #[cfg(not(all(feature = "backend", target_os = "illumos")))] @@ -13,12 +17,14 @@ use ddm::sm::{DpdConfig, SmContext, StateMachine}; #[cfg(all(feature = "backend", target_os = "illumos"))] use ddm::sys::Route; use ddm_api_types::db::RouterKind; +use mg_common::lock; use signal::handle_signals; use slog::{Drain, Logger, error}; use std::net::{IpAddr, Ipv6Addr}; #[cfg(all(feature = "backend", target_os = "illumos"))] use std::sync::mpsc::channel; use std::sync::{Arc, Mutex}; +use std::time::Duration; use uuid::Uuid; mod signal; @@ -32,26 +38,26 @@ struct Arg { addresses: Vec, /// How long to wait between solicitations (milliseconds). - #[arg(long, default_value_t = 2000)] + #[arg(long, default_value_t = millis_u64(SOLICIT_INTERVAL))] solicit_interval: u64, /// How long to wait without a solicitation response before expiring a peer /// (milliseconds). - #[arg(long, default_value_t = 5000)] + #[arg(long, default_value_t = millis_u64(EXPIRE_THRESHOLD))] expire_threshold: u64, /// How often to check for link failure while waiting for discovery messges /// (milliseconds). - #[arg(long, default_value_t = 1000)] + #[arg(long, default_value_t = millis_u64(DISCOVERY_READ_TIMEOUT))] discovery_read_timeout: u64, /// How long to wait between attempts to get an IP address for a specified /// address object (milliseconds). - #[arg(long, default_value_t = 1000)] + #[arg(long, default_value_t = millis_u64(IP_ADDR_WAIT))] ip_addr_wait: u64, - /// How long to wait for a response to exchange messages. - #[arg(long, default_value_t = 3000)] + /// How long to wait for a response to exchange messages (milliseconds). + #[arg(long, default_value_t = millis_u64(EXCHANGE_TIMEOUT))] pub exchange_timeout: u64, /// Address to listen on for the admin API. @@ -67,7 +73,7 @@ struct Arg { kind: RouterKind, /// The tcp port to listen on for exchange messages. - #[arg(long, default_value_t = 0xdddd)] + #[arg(long, default_value_t = EXCHANGE_TCP_PORT)] exchange_port: u16, /// Whether or not to use Dendrite as the underlying routing and forwarding @@ -107,6 +113,10 @@ struct Arg { #[arg(long)] sled_uuid: Option, + // Explicitly set router id instead of using hostname + #[arg(long)] + router_id: Option, + /// Serve only the admin API. Skips the routing state machine /// (discovery, exchange, route synchronization), allowing test fixtures /// to obtain a real `ddmd` admin endpoint without the kernel-level @@ -160,49 +170,65 @@ async fn run() { .to_string_lossy() .to_string(); - let (sms, event_channels) = - start_state_machines(&arg, &db, &dpd, &hostname, &rt, &log); + let router_id = match &arg.router_id { + Some(id) => id.clone(), + None => hostname.clone(), + }; + + let sms = start_state_machines(&arg, &db, &dpd, &router_id, &rt, &log); termination_handler(db.clone(), dpd.clone(), rt.clone(), log.clone()); let router_stats = Arc::new(RouterStats::default()); let peers: Vec = sms.iter().map(|x| x.ctx.clone()).collect(); - let stats_handler = if arg.with_stats { - if let (Some(rack_uuid), Some(sled_uuid)) = - (arg.rack_uuid, arg.sled_uuid) - { - match ddm::oxstats::start_server( - arg.oximeter_port, - peers.clone(), - router_stats.clone(), - hostname.clone(), - rack_uuid, - sled_uuid, - log.clone(), - ) { - Ok(handler) => Some(handler), - Err(e) => { - error!(log, "failed to start stats server: {e}"); - None - } - } - } else { - None - } - } else { - None - }; - let context = Arc::new(Mutex::new(HandlerContext { - event_channels, db, stats: router_stats, peers, - stats_handler: Arc::new(Mutex::new(stats_handler)), + stats_handler: Arc::new(Mutex::new(None)), + router_kind: arg.kind, + tunables: Tunables { + solicit_interval: Duration::from_millis(arg.solicit_interval), + expire_threshold: Duration::from_millis(arg.expire_threshold), + discovery_read_timeout: Duration::from_millis( + arg.discovery_read_timeout, + ), + ip_addr_wait: Duration::from_millis(arg.ip_addr_wait), + exchange_timeout: Duration::from_millis(arg.exchange_timeout), + dendrite: arg.dendrite, + dpd_port: arg.dpd_port, + dpd_host: arg.dpd_host.clone(), + exchange_tcp_port: arg.exchange_port, + }, log: log.clone(), + rack_id: arg.rack_uuid, + sled_id: arg.sled_uuid, + router_id: router_id.clone(), })); + if arg.with_stats + && let (Some(rack_uuid), Some(sled_uuid)) = + (arg.rack_uuid, arg.sled_uuid) + { + let h = match ddm::oxstats::start_server( + arg.oximeter_port, + context.clone(), + hostname.clone(), + rack_uuid, + sled_uuid, + log.clone(), + ) { + Ok(handler) => Some(handler), + Err(e) => { + error!(log, "failed to start stats server: {e}"); + None + } + }; + let ctx = lock!(context); + *lock!(ctx.stats_handler) = h; + } + if let Err(e) = sig_tx.send(context.clone()).await { error!(log, "send context to signal handler {e}"); } @@ -231,15 +257,12 @@ fn start_state_machines( arg: &Arg, db: &Db, dpd: &Option, - hostname: &str, + router_id: &str, rt: &Arc, log: &Logger, -) -> ( - Vec, - Vec>, -) { +) -> Vec { if arg.api_only { - return (Vec::new(), Vec::new()); + return Vec::new(); } let mut sms = Vec::new(); @@ -249,11 +272,13 @@ fn start_state_machines( let (tx, rx) = channel(); let config = ddm::sm::Config { - solicit_interval: arg.solicit_interval, - expire_threshold: arg.expire_threshold, - discovery_read_timeout: arg.discovery_read_timeout, - ip_addr_wait: arg.ip_addr_wait, - exchange_timeout: arg.exchange_timeout, + solicit_interval: Duration::from_millis(arg.solicit_interval), + expire_threshold: Duration::from_millis(arg.expire_threshold), + discovery_read_timeout: Duration::from_millis( + arg.discovery_read_timeout, + ), + ip_addr_wait: Duration::from_millis(arg.ip_addr_wait), + exchange_timeout: Duration::from_millis(arg.exchange_timeout), exchange_port: arg.exchange_port, aobj_name: name.clone(), if_name: String::new(), @@ -261,6 +286,8 @@ fn start_state_machines( kind: arg.kind, dpd: dpd.clone(), addr: Ipv6Addr::UNSPECIFIED, + sled_id: arg.sled_uuid, + rack_id: arg.rack_uuid, }; let ctx = SmContext { @@ -269,10 +296,12 @@ fn start_state_machines( event_channels: Vec::new(), tx: tx.clone(), log: log.clone(), - hostname: hostname.to_string(), + router_id: router_id.to_string(), rt: rt.clone(), iface: Arc::new(InterfaceState::default()), stats: Arc::new(ddm::sm::SessionStats::default()), + discovery_stop: None, + first_run: true, }; let sm = StateMachine { ctx, rx: Some(rx) }; @@ -296,7 +325,7 @@ fn start_state_machines( sm.run().unwrap(); } - (sms, event_channels) + sms } /// Non-illumos variant: the routing state machine depends on illumos @@ -310,11 +339,8 @@ fn start_state_machines( _hostname: &str, _rt: &Arc, _log: &Logger, -) -> ( - Vec, - Vec>, -) { - (Vec::new(), Vec::new()) +) -> Vec { + Vec::new() } /// Install a Ctrl-C handler that withdraws ddmd's imported routes from the diff --git a/ddmd/src/smf.rs b/ddmd/src/smf.rs index 8bcfb69e8..353a13e86 100644 --- a/ddmd/src/smf.rs +++ b/ddmd/src/smf.rs @@ -60,14 +60,16 @@ fn refresh_stats_server( } }; - let context = lock!(ctx); - let mut handler = lock!(context.stats_handler); + let (handler, log) = { + let ctx = lock!(ctx); + (ctx.stats_handler.clone(), ctx.log.clone()) + }; + let mut handler = lock!(handler); if handler.is_none() { info!(log, "starting stats server on smf refresh"); match ddm::oxstats::start_server( DDM_STATS_PORT, - context.peers.clone(), - context.stats.clone(), + ctx.clone(), hostname, props.rack_uuid, props.sled_uuid, diff --git a/openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json.gitstub b/openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json.gitstub new file mode 100644 index 000000000..b5ab91ce1 --- /dev/null +++ b/openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json.gitstub @@ -0,0 +1 @@ +561931c41eafde867f931229e799d1b30df7a1d4:openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json diff --git a/openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json b/openapi/ddm-admin/ddm-admin-3.0.0-3255d9.json similarity index 91% rename from openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json rename to openapi/ddm-admin/ddm-admin-3.0.0-3255d9.json index e891ca0e6..894267329 100644 --- a/openapi/ddm-admin/ddm-admin-2.0.0-45d40c.json +++ b/openapi/ddm-admin/ddm-admin-3.0.0-3255d9.json @@ -6,7 +6,7 @@ "url": "https://oxide.computer", "email": "api@oxide.computer" }, - "version": "2.0.0" + "version": "3.0.0" }, "paths": { "/disable-stats": { @@ -51,6 +51,53 @@ } } }, + "/external_peers": { + "get": { + "operationId": "get_external_peers", + "responses": { + "200": { + "description": "successful operation", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ExternalPeers" + } + } + } + }, + "4XX": { + "$ref": "#/components/responses/Error" + }, + "5XX": { + "$ref": "#/components/responses/Error" + } + } + }, + "put": { + "operationId": "set_external_peers", + "requestBody": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ExternalPeers" + } + } + }, + "required": true + }, + "responses": { + "204": { + "description": "resource updated" + }, + "4XX": { + "$ref": "#/components/responses/Error" + }, + "5XX": { + "$ref": "#/components/responses/Error" + } + } + } + }, "/originated": { "get": { "operationId": "get_originated", @@ -414,6 +461,21 @@ "request_id" ] }, + "ExternalPeers": { + "type": "object", + "properties": { + "address_objects": { + "type": "array", + "items": { + "type": "string" + }, + "uniqueItems": true + } + }, + "required": [ + "address_objects" + ] + }, "IpNet": { "x-rust-type": { "crate": "oxnet", diff --git a/openapi/ddm-admin/ddm-admin-latest.json b/openapi/ddm-admin/ddm-admin-latest.json index 4eb6e8dbc..49ff483ef 120000 --- a/openapi/ddm-admin/ddm-admin-latest.json +++ b/openapi/ddm-admin/ddm-admin-latest.json @@ -1 +1 @@ -ddm-admin-2.0.0-45d40c.json \ No newline at end of file +ddm-admin-3.0.0-3255d9.json \ No newline at end of file diff --git a/smf/ddm/manifest.xml b/smf/ddm/manifest.xml index 5e5da1cf2..86c3ac819 100644 --- a/smf/ddm/manifest.xml +++ b/smf/ddm/manifest.xml @@ -29,6 +29,7 @@ + diff --git a/smf/ddm_method_script.sh b/smf/ddm_method_script.sh index 706da75f6..a1e32a11b 100755 --- a/smf/ddm_method_script.sh +++ b/smf/ddm_method_script.sh @@ -46,6 +46,12 @@ if [[ "$val" != 'unknown' ]]; then args+=( "$val" ) fi +val=$(svcprop -c -p config/router_id "${SMF_FMRI}") +if [[ "$val" != 'unknown' ]]; then + args+=( '--router-id' ) + args+=( "$val" ) +fi + for x in $(svcprop -c -p config/interfaces "${SMF_FMRI}"); do args+=( '-a' ) args+=( "$x" ) diff --git a/tests/Cargo.toml b/tests/Cargo.toml index 3bc726298..6e08b83b7 100644 --- a/tests/Cargo.toml +++ b/tests/Cargo.toml @@ -18,3 +18,5 @@ slog-async.workspace = true tokio.workspace = true ztest.workspace = true uuid.workspace = true +rand.workspace = true +oxnet.workspace = true diff --git a/tests/conf/softnpu-quartet.toml b/tests/conf/softnpu-quartet-sidecar.quartet.toml similarity index 100% rename from tests/conf/softnpu-quartet.toml rename to tests/conf/softnpu-quartet-sidecar.quartet.toml diff --git a/tests/conf/softnpu-sextet-sidecar_a.sextet.toml b/tests/conf/softnpu-sextet-sidecar_a.sextet.toml new file mode 100644 index 000000000..87bce5690 --- /dev/null +++ b/tests/conf/softnpu-sextet-sidecar_a.sextet.toml @@ -0,0 +1,7 @@ +p4_program = "/opt/libsidecar_lite.so" +ports = [ + { sidecar = "sw0", scrimlet = "sq0", mtu = 1500 }, + { sidecar = "sw1", scrimlet = "sq1", mtu = 1500 }, + { sidecar = "sw2", scrimlet = "sr0", mtu = 1500 }, + { sidecar = "sw3", scrimlet = "sr1", mtu = 1500 }, +] diff --git a/tests/conf/softnpu-sextet-sidecar_b.sextet.toml b/tests/conf/softnpu-sextet-sidecar_b.sextet.toml new file mode 100644 index 000000000..e799a1ee2 --- /dev/null +++ b/tests/conf/softnpu-sextet-sidecar_b.sextet.toml @@ -0,0 +1,7 @@ +p4_program = "/opt/libsidecar_lite.so" +ports = [ + { sidecar = "sw4", scrimlet = "sq2", mtu = 1500 }, + { sidecar = "sw5", scrimlet = "sq3", mtu = 1500 }, + { sidecar = "sw6", scrimlet = "sr2", mtu = 1500 }, + { sidecar = "sw7", scrimlet = "sr3", mtu = 1500 }, +] diff --git a/tests/conf/softnpu-trio.toml b/tests/conf/softnpu-trio-sidecar.trio.toml similarity index 100% rename from tests/conf/softnpu-trio.toml rename to tests/conf/softnpu-trio-sidecar.trio.toml diff --git a/tests/src/ddm.rs b/tests/src/ddm.rs index 988ed973d..cb934b6e2 100644 --- a/tests/src/ddm.rs +++ b/tests/src/ddm.rs @@ -5,10 +5,14 @@ use anyhow::{Result, anyhow}; use client_common::{eprintln_nopipe, println_nopipe}; use ddm_admin_client::Client; +use ddm_api_types_versions::latest::external_peers::ExternalPeers; use ddm_api_types_versions::latest::net::TunnelOrigin; +use oxnet::Ipv6Net; use slog::{Drain, Logger}; +use std::collections::{BTreeMap, BTreeSet}; use std::env; use std::net::Ipv6Addr; +use std::ops::{Deref, DerefMut}; use std::thread::sleep; use std::time::Duration; use zone::Zlogin; @@ -45,6 +49,12 @@ macro_rules! softnpu_dump { }}; } +macro_rules! ip6_net { + ($x:expr) => { + $x.parse().unwrap() + }; +} + const ZONE_BRAND: &str = "omicron1"; struct SoftnpuZone<'a> { @@ -60,7 +70,7 @@ impl<'a> SoftnpuZone<'a> { ifx: &[&'a str], testname: &'a str, ) -> Result { - let softnpu_mount = format!("/tmp/softnpu/{}", testname); + let softnpu_mount = format!("/tmp/softnpu/{testname}/{name}"); std::fs::create_dir_all(&softnpu_mount)?; let fs = &[FsMount::new(&softnpu_mount, "/opt/mnt")]; @@ -91,7 +101,10 @@ impl<'a> SoftnpuZone<'a> { )?; self.zfs.copy_workspace_to_zone( &self.zone.name, - &format!("tests/conf/softnpu-{}.toml", self.testname), + &format!( + "tests/conf/softnpu-{}-{}.toml", + self.testname, self.zone.name + ), "opt/softnpu.toml", )?; self.zone.zexec(&format!( @@ -135,6 +148,7 @@ struct RouterZone<'a> { zone: Zone, transit: bool, testname: String, + port_map: BTreeMap, } impl<'a> RouterZone<'a> { @@ -144,7 +158,7 @@ impl<'a> RouterZone<'a> { mgmt: &'a str, rtr_ifx: &[&'a str], ) -> Result { - Self::new(name, zfs, mgmt, rtr_ifx, false, "") + Self::new(name, zfs, mgmt, rtr_ifx, false, "", "") } fn transit( @@ -153,8 +167,9 @@ impl<'a> RouterZone<'a> { mgmt: &'a str, rtr_ifx: &[&'a str], testname: &str, + softnpu_name: &str, ) -> Result { - Self::new(name, zfs, mgmt, rtr_ifx, true, testname) + Self::new(name, zfs, mgmt, rtr_ifx, true, testname, softnpu_name) } fn new( @@ -164,12 +179,14 @@ impl<'a> RouterZone<'a> { rtr_ifx: &[&'a str], transit: bool, testname: &str, + softnpu_name: &str, ) -> Result { let mut ifx = vec![mgmt]; ifx.extend_from_slice(rtr_ifx); let fs = if transit { - let softnpu_mount = format!("/tmp/softnpu/{}", testname); + let softnpu_mount = + format!("/tmp/softnpu/{testname}/{softnpu_name}"); std::fs::create_dir_all(&softnpu_mount)?; vec![FsMount::new(&softnpu_mount, "/opt/mnt")] } else { @@ -183,23 +200,51 @@ impl<'a> RouterZone<'a> { zone, transit, testname: testname.into(), + port_map: BTreeMap::default(), }) } + fn set_port_map(&mut self, pm: BTreeMap) { + self.port_map = pm; + } + fn stop_router(&self) -> Result { self.zone.zexec("pkill ddmd") } fn start_router(&self, restart_dpd: bool) -> Result<()> { - let addrs = self.ifx[1..] + let mapped_ports = self.ifx[1..] + .iter() + .map(|x| x.to_string()) + .map(|x| self.port_map.get(&x).unwrap_or(&x).clone()) + .collect::>(); + + let rear_ports = mapped_ports .iter() - .map(|x| format!("-a {}/v6", x)) - .collect::>() - .join(" "); + .filter(|&x| x.contains("rear")) + .cloned() + .collect::>(); + let front_ports = mapped_ports + .iter() + .filter(|&x| x.contains("qsfp")) + .cloned() + .collect::>(); + + let addrs = if self.transit { + &rear_ports + } else { + &mapped_ports + } + .iter() + .map(|x| format!("-a {}/v6", x)) + .collect::>() + .join(" "); let ddm = "/opt/ddmd"; + + // Tighter solicit interval and expire threshold are to speed up tests. let extra_args = format!( - "--rack-uuid {} --sled-uuid {}", + "--rack-uuid {} --sled-uuid {} --solicit-interval 200 --expire-threshold 500", uuid::Uuid::new_v4(), uuid::Uuid::new_v4(), ); @@ -216,11 +261,13 @@ impl<'a> RouterZone<'a> { self.zone.zexec( "svccfg -s dendrite setprop config/uds_path = /opt/mnt", )?; - self.zone.zexec( - "svccfg -s dendrite setprop config/port_config = /opt/dpd-ports.toml")?; + self.zone.zexec(&format!( + "svccfg -s dendrite setprop config/front_ports = {}", + front_ports.len(), + ))?; self.zone.zexec(&format!( "svccfg -s dendrite setprop config/rear_ports = {}", - self.ifx.len() - 1 + rear_ports.len(), ))?; self.zone.zexec("svcadm refresh dendrite:default")?; self.zone.zexec("svcadm enable dendrite:default")?; @@ -265,7 +312,18 @@ impl<'a> RouterZone<'a> { ), )?; - for ifx in &self.ifx[1..] { + for (link, vnic) in &self.port_map { + self.zone + .zexec(&format!("dladm create-vnic -t -l {link} {vnic}"))?; + } + + let mapped_ports = self.ifx[1..] + .iter() + .map(|x| x.to_string()) + .map(|x| self.port_map.get(&x).unwrap_or(&x).clone()) + .collect::>(); + + for ifx in &mapped_ports { self.zone.zcmd( &z, &format!("ipadm create-addr -t -T addrconf {}/v6", ifx), @@ -437,6 +495,7 @@ async fn test_trio() -> Result<()> { &mg1.name, &[&tf0_sr0.end_a, &tf1_sr1.end_a], "trio", + "sidecar.trio", )?; println_nopipe!("waiting for zones to come up"); @@ -653,7 +712,7 @@ async fn run_trio_tests( #[tokio::test] async fn test_quartet() -> Result<()> { - // A quartet of servers in a star topology. + // A quartet of routers in a star topology. // // sled1 // ,----------, @@ -726,6 +785,7 @@ async fn test_quartet() -> Result<()> { &mgt1.name, &[&tf0_sr0.end_a, &tf1_sr1.end_a, &tf2_sr2.end_a], "quartet", + "sidecar.quartet", )?; println_nopipe!("waiting for zones to come up"); @@ -802,6 +862,593 @@ async fn run_quartet_tests( Ok(()) } +#[tokio::test(flavor = "multi_thread")] +async fn test_external_peer_sextet() -> Result<()> { + // A sextet of routers in a multi-rack topology. + // + // sled1 + // ,----------, + // ,-----, ,-----, + // ,-| sl0 | | mg2 |-* + // scrimletA sidecarA | '-----' '-----' + // ,-----------, ,-----------------, | '----------' + // | ,-----, ,-----, ,-----, ,-----, | + // | | tr0a|--| sr0 |-| |-| sw2 |-' sled2 + // | '-----' '-----' | | '-----' ,----------, + // ,-----, ,-----, ,-----, |soft | ,-----, ,-----, ,-----, + // *-| mg1 | | tr1a|--| sr1 |-| npu|-| sw3 |---| sl1 | | mg3 |-* + // '-----' '-----' '-----' | | '-----' '-----' '-----' + // | ,-----, ,-----, | | ,-----, '----------' + // | | tq0a|--| sq0 |-| |-| sw0 |----, + // | '-----' '-----' | | '-----' | + // | ,-----, ,-----, | | ,-----, | + // | | tq1a|--| sq1 |-| |-| sw1 |-, | + // | '-----' '-----' '-----' '-----' | | + // '-----------' '-----------------' | | + // | | + // | | + // | | + // scrimletB sidecarB | | + // ,-----------, ,-----------------, | | + // | ,-----, ,-----, ,-----, ,-----, | | + // | | tq0b|--| sq2 |-| |-| sw4 |-' | + // | '-----' '-----' | | '-----' | + // ,-----, ,-----, ,-----, |soft | ,-----, | + // *-| mg4 | | tq1b|--| sq3 |-| npu|-| sw5 |----' sled3 + // '-----' '-----' '-----' | | '-----' ,----------, + // | ,-----, ,-----, | | ,-----, ,-----, ,-----, + // | | tr0b|--| sr2 |-| |-| sw6 |---| sl2 | | mg5 |-* + // | '-----' '-----' | | '-----' '-----' '-----' + // | ,-----, ,-----, | | ,-----, '----------' + // | | tr1b|--| sr3 |-| |-| sw7 |-, + // | '-----' '-----' '-----' '-----' | sled4 + // '-----------' '-----------------' | ,----------, + // | ,-----, ,-----, + // '-| sl3 | | mg6 |-* + // '-----' '-----' + // '----------' + + // Scrimlet A <-> Sidecar A + let tqa_sq_0 = SimnetLink::new("tqa0", "sq0")?; + let tqa_sq_1 = SimnetLink::new("tqa1", "sq1")?; + let tra_sr_0 = SimnetLink::new("tra0", "sr0")?; + let tra_sr_1 = SimnetLink::new("tra1", "sr1")?; + + // Sidecar A <-> Sleds + let sl0_sw2 = SimnetLink::new("sl0", "sw2")?; + let sl1_sw3 = SimnetLink::new("sl1", "sw3")?; + + // Scrimlet B <-> Sidecar B + let tqb_sq_0 = SimnetLink::new("tqb0", "sq2")?; + let tqb_sq_1 = SimnetLink::new("tqb1", "sq3")?; + let trb_sr_0 = SimnetLink::new("trb0", "sr2")?; + let trb_sr_1 = SimnetLink::new("trb1", "sr3")?; + + // Sidecar B <-> Sleds + let sl2_sw6 = SimnetLink::new("sl2", "sw6")?; + let sl3_sw7 = SimnetLink::new("sl3", "sw7")?; + + // Sidecar A <-> Sidecar B + let sw0_sw4 = SimnetLink::new("sw0", "sw5")?; + let sw1_sw5 = SimnetLink::new("sw1", "sw4")?; + + let mgmt0 = Etherstub::new("mgmt0")?; + let mg0 = Vnic::new("mg0", &mgmt0.name)?; + let mgs1 = Vnic::new("mgs1", &mgmt0.name)?; + let mgs2 = Vnic::new("mgs2", &mgmt0.name)?; + let mgs3 = Vnic::new("mgs3", &mgmt0.name)?; + let mgs4 = Vnic::new("mgs4", &mgmt0.name)?; + let mgs5 = Vnic::new("mgs5", &mgmt0.name)?; + let mgs6 = Vnic::new("mgs6", &mgmt0.name)?; + + let _mgip = Ip::new("10.0.0.254/24", &mg0.name, "test")?; + + let zfs = Zfs::new("mgtest")?; + + let sidecar_a = SoftnpuZone::new( + "sidecar_a.sextet", + &zfs, + &[ + &tqa_sq_0.end_b, + &tqa_sq_1.end_b, + &tra_sr_0.end_b, + &tra_sr_1.end_b, + &sl0_sw2.end_b, + &sl1_sw3.end_b, + &sw0_sw4.end_a, + &sw1_sw5.end_a, + ], + "sextet", + )?; + + let sidecar_b = SoftnpuZone::new( + "sidecar_b.sextet", + &zfs, + &[ + &tqb_sq_0.end_b, + &tqb_sq_1.end_b, + &trb_sr_0.end_b, + &trb_sr_1.end_b, + &sl2_sw6.end_b, + &sl3_sw7.end_b, + &sw0_sw4.end_b, + &sw1_sw5.end_b, + ], + "sextet", + )?; + + println_nopipe!("start zone s1"); + let s1 = + RouterZone::server("s1.sextet", &zfs, &mgs2.name, &[&sl0_sw2.end_a])?; + + println_nopipe!("start zone s2"); + let s2 = + RouterZone::server("s2.sextet", &zfs, &mgs3.name, &[&sl1_sw3.end_a])?; + + println_nopipe!("start zone s3"); + let s3 = + RouterZone::server("s3.sextet", &zfs, &mgs5.name, &[&sl2_sw6.end_a])?; + + println_nopipe!("start zone s4"); + let s4 = + RouterZone::server("s4.sextet", &zfs, &mgs6.name, &[&sl3_sw7.end_a])?; + + println_nopipe!("start zone t1"); + let mut t1 = RouterZone::transit( + "t1.sextet", + &zfs, + &mgs1.name, + &[ + &tqa_sq_0.end_a, + &tqa_sq_1.end_a, + &tra_sr_0.end_a, + &tra_sr_1.end_a, + ], + "sextet", + "sidecar_a.sextet", + )?; + t1.set_port_map(BTreeMap::from([ + ("tra0".into(), "tfportrear0_0".into()), + ("tra1".into(), "tfportrear1_0".into()), + ("tqa0".into(), "tfportqsfp0_0".into()), + ("tqa1".into(), "tfportqsfp1_0".into()), + ])); + + println_nopipe!("start zone t2"); + let mut t2 = RouterZone::transit( + "t2.sextet", + &zfs, + &mgs4.name, + &[ + &tqb_sq_0.end_a, + &tqb_sq_1.end_a, + &trb_sr_0.end_a, + &trb_sr_1.end_a, + ], + "sextet", + "sidecar_b.sextet", + )?; + t2.set_port_map(BTreeMap::from([ + ("trb0".into(), "tfportrear0_0".into()), + ("trb1".into(), "tfportrear1_0".into()), + ("tqb0".into(), "tfportqsfp0_0".into()), + ("tqb1".into(), "tfportqsfp1_0".into()), + ])); + + println_nopipe!("waiting for zones to come up"); + sleep(Duration::from_secs(10)); + + sidecar_a.setup()?; + sidecar_b.setup()?; + s1.setup(1)?; + s2.setup(2)?; + s3.setup(3)?; + s4.setup(4)?; + t1.setup(5)?; + t2.setup(6)?; + + run_topo!(run_sextet_tests(&s1, &s2, &s3, &s4, &t1, &t2).await)?; + + Ok(()) +} + +async fn run_sextet_tests( + _zs1: &RouterZone<'_>, + _zs2: &RouterZone<'_>, + _zs3: &RouterZone<'_>, + _zs4: &RouterZone<'_>, + _zt1: &RouterZone<'_>, + _zt2: &RouterZone<'_>, +) -> Result<()> { + let log = init_logger(); + + // A ddm client that dumps out information when it drops. Primarily used for + // debugging test failures when an assert pops. + struct DropDump { + c: Client, + name: String, + } + impl Deref for DropDump { + type Target = Client; + fn deref(&self) -> &Self::Target { + &self.c + } + } + impl DerefMut for DropDump { + fn deref_mut(&mut self) -> &mut Self::Target { + &mut self.c + } + } + impl Drop for DropDump { + fn drop(&mut self) { + // Async just loves to make things difficult, it's taken over all + // the things, but heaven forbid you need to do an async thing in + // the most basic of object lifecycle management traits ... + let rt = tokio::runtime::Handle::current(); + let c = self.c.clone(); + let name = self.name.clone(); + tokio::task::block_in_place(|| { + rt.block_on(async move { + println_nopipe!("{name}:"); + if let Ok(peers) = c.get_peers().await { + println_nopipe!("peers: {peers:#?}"); + } + if let Ok(prefixes) = c.get_prefixes().await { + println_nopipe!("prefixes: {prefixes:#?}"); + } + }); + }); + } + } + + macro_rules! drop_dump { + ($name:ident, $endpoint:expr) => { + let $name = DropDump { + c: Client::new($endpoint, log.clone()), + name: stringify!($name).to_string(), + }; + }; + } + + #[derive(Default)] + struct PeerCounts { + s1: usize, + s2: usize, + s3: usize, + s4: usize, + t1: usize, + t2: usize, + } + impl PeerCounts { + fn server(mut self, c: usize) -> Self { + self.s1 = c; + self.s2 = c; + self.s3 = c; + self.s4 = c; + self + } + fn transit(mut self, c: usize) -> Self { + self.t1 = c; + self.t2 = c; + self + } + } + + struct PeerReachablePrefixes { + s1: BTreeSet, + s2: BTreeSet, + s3: BTreeSet, + s4: BTreeSet, + t1: BTreeSet, + t2: BTreeSet, + } + + drop_dump!(s1, "http://10.0.0.1:8000"); + drop_dump!(s2, "http://10.0.0.2:8000"); + drop_dump!(s3, "http://10.0.0.3:8000"); + drop_dump!(s4, "http://10.0.0.4:8000"); + drop_dump!(t1, "http://10.0.0.5:8000"); + drop_dump!(t2, "http://10.0.0.6:8000"); + + // While this would be better as a simple lambda function, when an assert + // pops within we only see the line number of the asserting statement in the + // lambda and all that's available in RUST_BACKTRACE=1 is a pile of useless + // tokio noise. + macro_rules! assert_peer_count { + ($client:expr, $count:expr) => {{ + println_nopipe!( + "ensure {} has {} peers", + stringify!($client), + $count + ); + wait_for_eq!( + $client.get_peers().await.map(|x| x.len()).ok(), + Some($count) + ); + }}; + } + + macro_rules! assert_peer_counts { + ($c:expr) => {{ + assert_peer_count!(s1, $c.s1); + assert_peer_count!(s2, $c.s2); + assert_peer_count!(s3, $c.s3); + assert_peer_count!(s4, $c.s4); + assert_peer_count!(t1, $c.t1); + assert_peer_count!(t2, $c.t2); + }}; + } + + macro_rules! assert_peer_reach { + ($client:expr, $reach:expr) => {{ + println_nopipe!( + "ensure {} has imported prefixes {:?}", + stringify!($client), + $reach + ); + wait_for_eq!( + $client + .get_prefixes() + .await + .map(|x| x + .values() + .cloned() + .into_iter() + .flat_map(|x| x + .clone() + .into_iter() + .map(|y| y.destination)) + .collect::>()) + .ok(), + Some($reach) + ); + }}; + } + + macro_rules! assert_reach { + ($r:expr) => {{ + assert_peer_reach!(s1, $r.s1.clone()); + assert_peer_reach!(s2, $r.s2.clone()); + assert_peer_reach!(s3, $r.s3.clone()); + assert_peer_reach!(s4, $r.s4.clone()); + assert_peer_reach!(t1, $r.t1.clone()); + assert_peer_reach!(t2, $r.t2.clone()); + }}; + } + + macro_rules! assert_nexthops_are_peers { + ($client:expr) => {{ + let pfx_nexthops = + $client.get_prefixes().await.expect("get prefixes"); + let peers = $client + .get_peers() + .await + .expect("get peers") + .values() + .map(|x| x.addr.to_string()) + .collect::>(); + + for p in pfx_nexthops.keys() { + assert!(peers.contains(p), "nexthop {p} is not a peer"); + } + }}; + } + + macro_rules! assert_all_nexthops_are_peers { + () => {{ + assert_nexthops_are_peers!(s1); + assert_nexthops_are_peers!(s2); + assert_nexthops_are_peers!(s3); + assert_nexthops_are_peers!(s4); + assert_nexthops_are_peers!(t1); + assert_nexthops_are_peers!(t2); + }}; + } + + // + // Initialize announcements for each server peer + // + + let s1_origin: Vec = [ip6_net!("fd00:1::/64")].into(); + let s2_origin: Vec = [ip6_net!("fd00:2::/64")].into(); + let s3_origin: Vec = [ip6_net!("fd00:3::/64")].into(); + let s4_origin: Vec = [ip6_net!("fd00:4::/64")].into(); + let t1_origin: Vec = s1_origin + .clone() + .into_iter() + .chain(s2_origin.clone()) + .collect(); + let t2_origin: Vec = s3_origin + .clone() + .into_iter() + .chain(s4_origin.clone()) + .collect(); + + s1.advertise_prefixes(&s1_origin).await?; + s2.advertise_prefixes(&s2_origin).await?; + s3.advertise_prefixes(&s3_origin).await?; + s4.advertise_prefixes(&s4_origin).await?; + + // + // Starting out we should have just the backplane peers. + // + + assert_peer_counts!(PeerCounts::default().server(1).transit(2)); + + // + // Specifying two external peers should result in two additional peers for + // each transit router and no changes for the number of server router peers. + // + + // The address objects for each qsfp on each switch in the test environment. + const QSFP0: &str = "tfportqsfp0_0/v6"; + const QSFP1: &str = "tfportqsfp1_0/v6"; + + let ext_peers_both = ExternalPeers { + address_objects: [QSFP0, QSFP1].map(String::from).into(), + }; + t1.set_external_peers(&ext_peers_both).await?; + t2.set_external_peers(&ext_peers_both).await?; + assert_peer_counts!(PeerCounts::default().server(1).transit(4)); + + // + // Going down to the first peer should result in three peering sessions + // per transit router. Note in the model above that qsfp0/qsfp1 are cross + // connected between the two transit routers. Here we are connecting + // swA/qsfp0 <-> swB/qsfp1 + // + + let ext_peers_qsfp0 = ExternalPeers { + address_objects: [QSFP0].map(String::from).into(), + }; + let ext_peers_qsfp1 = ExternalPeers { + address_objects: [QSFP1].map(String::from).into(), + }; + t1.set_external_peers(&ext_peers_qsfp0).await?; + t2.set_external_peers(&ext_peers_qsfp1).await?; + assert_peer_counts!(PeerCounts::default().server(1).transit(3)); + + // + // Go back to full peering and then switch to swA/qsfp1 <-> swB/qsfp0 + // + + t1.set_external_peers(&ext_peers_both).await?; + t2.set_external_peers(&ext_peers_both).await?; + assert_peer_counts!(PeerCounts::default().server(1).transit(4)); + t1.set_external_peers(&ext_peers_qsfp1).await?; + t2.set_external_peers(&ext_peers_qsfp0).await?; + assert_peer_counts!(PeerCounts::default().server(1).transit(3)); + + // + // Switch from swA/qsfp1 <-> swB/qsfp0 to swA/qsfp0 <-> swB/qsfp1 + // + + t1.set_external_peers(&ext_peers_qsfp0).await?; + t2.set_external_peers(&ext_peers_qsfp1).await?; + assert_peer_counts!(PeerCounts::default().server(1).transit(3)); + + // + // Go to no external peers + // + + let ext_peers_none = ExternalPeers { + address_objects: BTreeSet::default(), + }; + t1.set_external_peers(&ext_peers_none).await?; + t2.set_external_peers(&ext_peers_none).await?; + assert_peer_counts!(PeerCounts::default().server(1).transit(2)); + + // + // A bit of combinatorial exercise + // + + fn peer_is_set(x: &ExternalPeers, s: &str) -> bool { + x.address_objects.contains(&String::from(s)) + } + + fn expected_external_peerings( + x: &ExternalPeers, + y: &ExternalPeers, + ) -> PeerCounts { + // The count starts at two because each transit router has two backplane + // connections that we each expect to have a server peering session on. + let mut ext_count: usize = 2; + + // A peering is expected when qsfp0 and qsfp1 are configured as an + // external router in either direction. This is a property of the + // testing topology (see diagram in test_external_peer_sextet). + if peer_is_set(x, QSFP0) && peer_is_set(y, QSFP1) { + ext_count += 1; + } + if peer_is_set(x, QSFP1) && peer_is_set(y, QSFP0) { + ext_count += 1; + } + PeerCounts::default().server(1).transit(ext_count) + } + + let expected_reachable_prefixes = + |x: &ExternalPeers, y: &ExternalPeers| -> PeerReachablePrefixes { + let counts = expected_external_peerings(x, y); + // Servers can always see the originated prefixes of other routers + // reachable over a single hop transit router path (e.g. in the same + // rack). Transit routers can always see the originated prefixes of + // their directly connected server routers. + let mut reach = PeerReachablePrefixes { + s1: s2_origin.iter().cloned().collect(), + s2: s1_origin.iter().cloned().collect(), + s3: s4_origin.iter().cloned().collect(), + s4: s3_origin.iter().cloned().collect(), + t1: t1_origin.iter().cloned().collect(), + t2: t2_origin.iter().cloned().collect(), + }; + // If there is any peering between transit routers, each server router + // should see prefixes originated from the router adjacent to their + // transit router. + if counts.t1 > 2 || counts.t2 > 2 { + // Origins from servers connected to t2 propagating to servers + // connected to t1. + reach.s1.extend(&s3_origin); + reach.s1.extend(&s4_origin); + reach.s2.extend(&s3_origin); + reach.s2.extend(&s4_origin); + + // Origins from servers connected to t1 propagating to servers + // connected to t2. + reach.s3.extend(&s1_origin); + reach.s3.extend(&s2_origin); + reach.s4.extend(&s1_origin); + reach.s4.extend(&s2_origin); + + // Origins from transit routers propagate to each other. + reach.t1.extend(&t2_origin); + reach.t2.extend(&t1_origin); + } + reach + }; + + // The choices we have for each switch are none, one or both peers where the + // one case can be either of the peers. + let choices = [ + ext_peers_none, + ext_peers_qsfp0, + ext_peers_qsfp1, + ext_peers_both, + ]; + + // Go through 100 rounds. Ideally we'd have more than this, but peer + // expiration and re-establishment is currently a second or two, so at 100 + // rounds this is already taking over a minute. It'd be nice to have really + // quick peering timer settings for tests so we can rapidly iterate through + // sweeps like this. + const N: usize = 100; + for i in 0..N { + println_nopipe!("{i}/{N}"); + let a: usize = rand::random_range(0..choices.len()); + let b: usize = rand::random_range(0..choices.len()); + + let x = &choices[a]; + let y = &choices[b]; + + t1.set_external_peers(x).await?; + t2.set_external_peers(y).await?; + let counts = expected_external_peerings(x, y); + assert_peer_counts!(counts); + + let xx = t1.get_external_peers().await?.into_inner(); + let yy = t2.get_external_peers().await?.into_inner(); + + assert_eq!(x, &xx, "t1 reports different peers than we set"); + assert_eq!(y, &yy, "t2 reports different peers than we set"); + + let reach = expected_reachable_prefixes(x, y); + assert_reach!(reach); + + assert_all_nexthops_are_peers!(); + } + + Ok(()) +} + async fn prefix_count(c: &Client) -> Result { Ok(c.get_prefixes() .await?