diff --git a/.gitignore b/.gitignore index 2dbd36d349..c30846aa3e 100644 --- a/.gitignore +++ b/.gitignore @@ -1,7 +1,7 @@ **.profraw **/__fuzz__/** +.fuzz-corpus/ # libfuzzer writes one of these per worker into the working directory when -# `just fuzz` is given -j; the corpus itself lives under __fuzz__. fuzz-*.log # qemu-user core dumps from SIGABRT under emulated tests. **/qemu_*.core diff --git a/Cargo.lock b/Cargo.lock index 47b54c27c4..26f6af9125 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1185,6 +1185,7 @@ dependencies = [ "arrayvec", "axum", "axum-server", + "bolero", "dataplane-acl-filter", "dataplane-args", "dataplane-clock", @@ -1196,6 +1197,7 @@ dependencies = [ "dataplane-flow-filter", "dataplane-id", "dataplane-lifecycle", + "dataplane-lpm", "dataplane-mgmt", "dataplane-nat", "dataplane-net", @@ -1224,6 +1226,7 @@ dependencies = [ "tokio", "tracing", "tracing-subscriber", + "tracing-test", ] [[package]] diff --git a/config/src/external/overlay/algebra.rs b/config/src/external/overlay/algebra.rs new file mode 100644 index 0000000000..039b7450db --- /dev/null +++ b/config/src/external/overlay/algebra.rs @@ -0,0 +1,1260 @@ +// SPDX-License-Identifier: Apache-2.0 +// Copyright Open Network Fabric Authors + +use std::collections::{BTreeMap, BTreeSet}; +use std::net::Ipv4Addr; +use std::ops::Bound::Included; + +use bolero::{Driver, ValueGenerator}; +use lpm::prefix::{IpPrefix, Ipv4Prefix, Prefix}; + +use crate::ConfigError; +use crate::external::overlay::Overlay; +use crate::external::overlay::vpc::{Vpc, VpcTable}; +use crate::external::overlay::vpcpeering::{VpcExpose, VpcManifest, VpcPeering, VpcPeeringTable}; + +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash)] +pub struct VpcHandle(pub u8); + +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash)] +pub struct PeeringHandle(pub u8); + +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash)] +pub enum Side { + Left, + Right, +} + +impl Side { + #[must_use] + pub fn other(self) -> Self { + match self { + Side::Left => Side::Right, + Side::Right => Side::Left, + } + } + + fn index(self) -> usize { + match self { + Side::Left => 0, + Side::Right => 1, + } + } +} + +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug, Hash)] +pub enum Flavour { + Forward, + Masquerade, +} + +pub const MAX_EXPOSES: u8 = 4; + +pub const MAX_VPCS: u8 = 6; + +const SLOTS_PER_PEERING: u32 = 2 * MAX_EXPOSES as u32; + +impl VpcHandle { + fn name(self) -> String { + format!("VPC-{:03}", self.0) + } + + fn id(self) -> String { + format!("V{:04}", self.0) + } + + fn vni(self) -> u32 { + 1000 + u32::from(self.0) + } +} + +impl PeeringHandle { + fn name(self) -> String { + format!("PEERING-{:03}", self.0) + } + + fn block(self, side: Side, slot: u8) -> u32 { + u32::from(self.0) * SLOTS_PER_PEERING + + u32::try_from(side.index()).unwrap_or_else(|_| unreachable!()) + * u32::from(MAX_EXPOSES) + + u32::from(slot) + } +} + +fn private_prefix(index: u32) -> Prefix { + prefix_v4(0x0A00_0000 | (index << 8), 24) +} + +fn public_prefix(index: u32) -> Prefix { + prefix_v4(0xAC10_0000 | (index << 8), 24) +} + +fn prefix_v4(bits: u32, len: u8) -> Prefix { + Prefix::from( + Ipv4Prefix::new(Ipv4Addr::from_bits(bits), len) + .unwrap_or_else(|_| unreachable!("a well-formed prefix")), + ) +} + +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +pub struct ExposeSpec { + slot: u8, + flavour: Flavour, +} + +impl ExposeSpec { + #[must_use] + pub fn flavour(self) -> Flavour { + self.flavour + } + + #[must_use] + pub fn private(self, peering: PeeringHandle, side: Side) -> Prefix { + private_prefix(peering.block(side, self.slot)) + } + + #[must_use] + pub fn public(self, peering: PeeringHandle, side: Side) -> Prefix { + match self.flavour { + Flavour::Forward => self.private(peering, side), + Flavour::Masquerade => public_prefix(peering.block(side, self.slot)), + } + } + + fn expose(self, peering: PeeringHandle, side: Side) -> VpcExpose { + let private = self.private(peering, side); + match self.flavour { + Flavour::Forward => VpcExpose::empty().ip(private.into()), + Flavour::Masquerade => VpcExpose::empty() + .make_masquerade(None) + .unwrap_or_else(|_| unreachable!("an empty expose accepts masquerade")) + .ip(private.into()) + .as_range(self.public(peering, side).into()) + .unwrap_or_else(|_| unreachable!("a masquerade expose accepts a public range")), + } + } +} + +#[derive(Clone, PartialEq, Eq, Debug)] +pub struct PeeringSpec { + left: VpcHandle, + right: VpcHandle, + exposes: [Vec; 2], +} + +impl PeeringSpec { + #[must_use] + pub fn vpc(&self, side: Side) -> VpcHandle { + match side { + Side::Left => self.left, + Side::Right => self.right, + } + } + + #[must_use] + pub fn exposes(&self, side: Side) -> &[ExposeSpec] { + &self.exposes[side.index()] + } + + fn exposes_mut(&mut self, side: Side) -> &mut Vec { + &mut self.exposes[side.index()] + } + + fn touches(&self, vpc: VpcHandle) -> bool { + self.left == vpc || self.right == vpc + } + + fn side_of(&self, vpc: VpcHandle) -> Option { + if self.left == vpc { + Some(Side::Left) + } else if self.right == vpc { + Some(Side::Right) + } else { + None + } + } + + fn has_stateful(&self, side: Side) -> bool { + self.exposes(side) + .iter() + .any(|expose| expose.flavour == Flavour::Masquerade) + } + + fn manifest(&self, peering: PeeringHandle, side: Side) -> VpcManifest { + self.exposes(side).iter().fold( + VpcManifest::new(&self.vpc(side).name()), + |manifest, spec| manifest.exposing(spec.expose(peering, side)), + ) + } +} + +#[derive(Clone, PartialEq, Eq, Debug, Default)] +pub struct Draft { + vpcs: BTreeSet, + peerings: BTreeMap, +} + +impl Draft { + #[must_use] + pub fn new() -> Self { + Self::default() + } + + #[must_use] + pub fn vpcs(&self) -> impl ExactSizeIterator + '_ { + self.vpcs.iter().copied() + } + + #[must_use] + pub fn peerings(&self) -> impl ExactSizeIterator { + self.peerings.iter().map(|(handle, spec)| (*handle, spec)) + } + + #[must_use] + pub fn peering_between(&self, left: VpcHandle, right: VpcHandle) -> Option { + self.peerings + .iter() + .find(|(_, spec)| spec.touches(left) && spec.touches(right)) + .map(|(handle, _)| *handle) + } + + #[must_use] + pub fn components(&self) -> Vec> { + let mut unvisited: BTreeSet = self.vpcs.clone(); + let mut components = Vec::new(); + while let Some(&seed) = unvisited.iter().next() { + let mut component = Vec::new(); + let mut frontier = vec![seed]; + unvisited.remove(&seed); + while let Some(vpc) = frontier.pop() { + component.push(vpc); + for spec in self.peerings.values() { + let Some(side) = spec.side_of(vpc) else { + continue; + }; + let peer = spec.vpc(side.other()); + if unvisited.remove(&peer) { + frontier.push(peer); + } + } + } + component.sort_unstable(); + components.push(component); + } + components.sort_unstable(); + components + } + + /// # Errors + /// + /// Returns the [`ConfigError`] from the first vpc, peering or overlay that the draft + /// does not describe validly. + pub fn overlay(&self) -> Result { + let mut vpc_table = VpcTable::new(); + for vpc in self.vpcs() { + vpc_table.add(Vpc::new(&vpc.name(), &vpc.id(), vpc.vni())?)?; + } + + let mut peerings = VpcPeeringTable::new(); + for (handle, spec) in self.peerings() { + peerings.add(VpcPeering::with_default_group( + &handle.name(), + spec.manifest(handle, Side::Left), + spec.manifest(handle, Side::Right), + ))?; + } + + Ok(Overlay::new(vpc_table, peerings)) + } +} + +#[derive(Clone, PartialEq, Eq, Debug, Default)] +pub struct Footprint { + vpcs: BTreeSet, + peerings: BTreeSet, +} + +impl Footprint { + fn of( + vpcs: impl IntoIterator, + peerings: impl IntoIterator, + ) -> Self { + Self { + vpcs: vpcs.into_iter().collect(), + peerings: peerings.into_iter().collect(), + } + } + + #[must_use] + pub fn intersects(&self, other: &Self) -> bool { + self.vpcs.intersection(&other.vpcs).next().is_some() + || self.peerings.intersection(&other.peerings).next().is_some() + } + + #[must_use] + pub fn is_empty(&self) -> bool { + self.vpcs.is_empty() && self.peerings.is_empty() + } +} + +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +pub enum Op { + AddVpc(VpcHandle), + RemoveVpc(VpcHandle), + AddPeering { + handle: PeeringHandle, + left: VpcHandle, + right: VpcHandle, + }, + RemovePeering(PeeringHandle), + AddExpose { + peering: PeeringHandle, + side: Side, + slot: u8, + flavour: Flavour, + }, + RemoveExpose { + peering: PeeringHandle, + side: Side, + slot: u8, + }, + SetFlavour { + peering: PeeringHandle, + side: Side, + slot: u8, + flavour: Flavour, + }, +} + +#[derive(Clone, PartialEq, Eq, Debug)] +pub enum Undo { + RemoveVpc(VpcHandle), + RestoreVpc { + handle: VpcHandle, + peerings: Vec<(PeeringHandle, PeeringSpec)>, + }, + RemovePeering(PeeringHandle), + RestorePeering(PeeringHandle, PeeringSpec), + RemoveExpose { + peering: PeeringHandle, + side: Side, + slot: u8, + }, + RestoreExpose { + peering: PeeringHandle, + side: Side, + index: usize, + spec: ExposeSpec, + }, + SetFlavour { + peering: PeeringHandle, + side: Side, + slot: u8, + flavour: Flavour, + }, +} + +impl Op { + #[must_use] + pub fn reads(&self) -> Footprint { + match self { + Op::AddVpc(_) | Op::RemoveVpc(_) | Op::RemovePeering(_) => Footprint::default(), + Op::AddPeering { left, right, .. } => Footprint::of([*left, *right], []), + Op::AddExpose { peering, .. } + | Op::RemoveExpose { peering, .. } + | Op::SetFlavour { peering, .. } => Footprint::of([], [*peering]), + } + } + + #[must_use] + pub fn writes(&self, draft: &Draft) -> Footprint { + match self { + Op::AddVpc(handle) => Footprint::of([*handle], []), + Op::RemoveVpc(handle) => { + let mut vpcs = BTreeSet::from([*handle]); + let mut peerings = BTreeSet::new(); + for (peering, spec) in draft.peerings() { + if spec.touches(*handle) { + peerings.insert(peering); + vpcs.insert(spec.left); + vpcs.insert(spec.right); + } + } + Footprint { vpcs, peerings } + } + Op::AddPeering { + handle, + left, + right, + } => Footprint::of([*left, *right], [*handle]), + Op::RemovePeering(handle) => { + let mut vpcs = BTreeSet::new(); + if let Some(spec) = draft.peerings.get(handle) { + vpcs.insert(spec.left); + vpcs.insert(spec.right); + } + Footprint { + vpcs, + peerings: BTreeSet::from([*handle]), + } + } + Op::AddExpose { peering, .. } + | Op::RemoveExpose { peering, .. } + | Op::SetFlavour { peering, .. } => Footprint::of([], [*peering]), + } + } + + #[must_use] + pub fn applicable(&self, draft: &Draft) -> bool { + let mut trial = draft.clone(); + self.apply(&mut trial).is_some() + } + + pub fn apply(&self, draft: &mut Draft) -> Option { + match *self { + Op::AddVpc(handle) => { + if !draft.vpcs.insert(handle) { + return None; + } + Some(Undo::RemoveVpc(handle)) + } + + Op::RemoveVpc(handle) => remove_vpc(draft, handle), + + Op::AddPeering { + handle, + left, + right, + } => { + if left == right + || draft.peerings.contains_key(&handle) + || !draft.vpcs.contains(&left) + || !draft.vpcs.contains(&right) + || draft.peering_between(left, right).is_some() + { + return None; + } + let first = |slot| ExposeSpec { + slot, + flavour: Flavour::Forward, + }; + draft.peerings.insert( + handle, + PeeringSpec { + left, + right, + exposes: [vec![first(0)], vec![first(0)]], + }, + ); + Some(Undo::RemovePeering(handle)) + } + + Op::RemovePeering(handle) => { + let spec = draft.peerings.remove(&handle)?; + Some(Undo::RestorePeering(handle, spec)) + } + + Op::AddExpose { + peering, + side, + slot, + flavour, + } => add_expose(draft, peering, side, slot, flavour), + + Op::RemoveExpose { + peering, + side, + slot, + } => { + let spec = draft.peerings.get_mut(&peering)?; + if spec.exposes(side).len() <= 1 { + return None; + } + let index = spec.exposes(side).iter().position(|e| e.slot == slot)?; + let removed = spec.exposes_mut(side).remove(index); + Some(Undo::RestoreExpose { + peering, + side, + index, + spec: removed, + }) + } + + Op::SetFlavour { + peering, + side, + slot, + flavour, + } => set_flavour(draft, peering, side, slot, flavour), + } + } +} + +fn remove_vpc(draft: &mut Draft, handle: VpcHandle) -> Option { + if !draft.vpcs.remove(&handle) { + return None; + } + let doomed: Vec = draft + .peerings + .iter() + .filter(|(_, spec)| spec.touches(handle)) + .map(|(peering, _)| *peering) + .collect(); + let peerings = doomed + .into_iter() + .map(|peering| { + let spec = draft + .peerings + .remove(&peering) + .unwrap_or_else(|| unreachable!("just found it")); + (peering, spec) + }) + .collect(); + Some(Undo::RestoreVpc { handle, peerings }) +} + +fn add_expose( + draft: &mut Draft, + peering: PeeringHandle, + side: Side, + slot: u8, + flavour: Flavour, +) -> Option { + if slot >= MAX_EXPOSES { + return None; + } + let spec = draft.peerings.get(&peering)?; + if spec.exposes(side).iter().any(|e| e.slot == slot) { + return None; + } + if flavour == Flavour::Masquerade && spec.has_stateful(side.other()) { + return None; + } + draft + .peerings + .get_mut(&peering) + .unwrap_or_else(|| unreachable!("just found it")) + .exposes_mut(side) + .push(ExposeSpec { slot, flavour }); + Some(Undo::RemoveExpose { + peering, + side, + slot, + }) +} + +fn set_flavour( + draft: &mut Draft, + peering: PeeringHandle, + side: Side, + slot: u8, + flavour: Flavour, +) -> Option { + let spec = draft.peerings.get(&peering)?; + if flavour == Flavour::Masquerade && spec.has_stateful(side.other()) { + return None; + } + let index = spec.exposes(side).iter().position(|e| e.slot == slot)?; + let exposes = draft + .peerings + .get_mut(&peering) + .unwrap_or_else(|| unreachable!("just found it")) + .exposes_mut(side); + let previous = exposes[index].flavour; + exposes[index].flavour = flavour; + Some(Undo::SetFlavour { + peering, + side, + slot, + flavour: previous, + }) +} + +impl Undo { + /// # Panics + /// + /// Panics if the draft is not the one this undo record was made against. + pub fn apply(&self, draft: &mut Draft) { + match self { + Undo::RemoveVpc(handle) => { + assert!(draft.vpcs.remove(handle), "no vpc {handle:?} to remove"); + } + Undo::RestoreVpc { handle, peerings } => { + assert!(draft.vpcs.insert(*handle), "vpc {handle:?} is still there"); + for (peering, spec) in peerings { + assert!( + draft.peerings.insert(*peering, spec.clone()).is_none(), + "peering {peering:?} is still there" + ); + } + } + Undo::RemovePeering(handle) => { + assert!( + draft.peerings.remove(handle).is_some(), + "no peering {handle:?} to remove" + ); + } + Undo::RestorePeering(handle, spec) => { + assert!( + draft.peerings.insert(*handle, spec.clone()).is_none(), + "peering {handle:?} is still there" + ); + } + Undo::RemoveExpose { + peering, + side, + slot, + } => { + let spec = draft + .peerings + .get_mut(peering) + .unwrap_or_else(|| unreachable!("no peering {peering:?}")); + let index = spec + .exposes(*side) + .iter() + .position(|e| e.slot == *slot) + .unwrap_or_else(|| unreachable!("no expose in slot {slot}")); + spec.exposes_mut(*side).remove(index); + } + Undo::RestoreExpose { + peering, + side, + index, + spec, + } => { + draft + .peerings + .get_mut(peering) + .unwrap_or_else(|| unreachable!("no peering {peering:?}")) + .exposes_mut(*side) + .insert(*index, *spec); + } + Undo::SetFlavour { + peering, + side, + slot, + flavour, + } => { + let spec = draft + .peerings + .get_mut(peering) + .unwrap_or_else(|| unreachable!("no peering {peering:?}")); + let index = spec + .exposes(*side) + .iter() + .position(|e| e.slot == *slot) + .unwrap_or_else(|| unreachable!("no expose in slot {slot}")); + spec.exposes_mut(*side)[index].flavour = *flavour; + } + } + } +} + +pub const MAX_SEQUENCE: u8 = 32; + +#[derive(Clone, Copy, Debug)] +pub struct Sequence { + pub len: u8, +} + +impl Default for Sequence { + fn default() -> Self { + Self { len: 24 } + } +} + +impl Sequence { + #[must_use] + pub fn of(len: u8) -> Self { + Self { + len: len.min(MAX_SEQUENCE), + } + } + + /// # Panics + /// + /// Panics if any operation does not apply to the draft the ones before it produced. + #[must_use] + pub fn fold(ops: &[Op]) -> Draft { + let mut draft = Draft::new(); + for (index, op) in ops.iter().enumerate() { + assert!( + op.apply(&mut draft).is_some(), + "operation {index} ({op:?}) does not apply" + ); + } + draft + } +} + +impl ValueGenerator for Sequence { + type Output = Vec; + + fn generate(&self, driver: &mut D) -> Option> { + let count = driver.gen_u8(Included(&0), Included(&self.len.min(MAX_SEQUENCE)))?; + let target = usize::from(driver.gen_u8(Included(&1), Included(&4))?); + let mut draft = Draft::new(); + let mut next_vpc = 0u8; + let mut next_peering = 0u8; + let mut ops = Vec::with_capacity(usize::from(count)); + + for _ in 0..count { + let Some(op) = draw(driver, &draft, &mut next_vpc, &mut next_peering, target) else { + break; + }; + op.apply(&mut draft)?; + ops.push(op); + } + + Some(ops) + } +} + +const MENU: [(Kind, u8); 7] = [ + (Kind::AddVpc, 4), + (Kind::RemoveVpc, 1), + (Kind::AddPeering, 4), + (Kind::RemovePeering, 1), + (Kind::AddExpose, 2), + (Kind::RemoveExpose, 1), + (Kind::SetFlavour, 2), +]; + +#[derive(Clone, Copy, PartialEq, Eq, Debug)] +enum Kind { + AddVpc, + RemoveVpc, + AddPeering, + RemovePeering, + AddExpose, + RemoveExpose, + SetFlavour, +} + +impl Kind { + fn applicable(self, draft: &Draft, next_vpc: u8, next_peering: u8) -> bool { + let sides = || { + draft + .peerings() + .flat_map(|(_, spec)| [Side::Left, Side::Right].map(move |side| spec.exposes(side))) + }; + match self { + Kind::AddVpc => next_vpc < u8::MAX && draft.vpcs.len() < usize::from(MAX_VPCS), + Kind::RemoveVpc => !draft.vpcs.is_empty(), + Kind::AddPeering => { + let vpcs = draft.vpcs.len(); + next_peering < u8::MAX && vpcs >= 2 && draft.peerings.len() < vpcs * (vpcs - 1) / 2 + } + Kind::RemovePeering | Kind::SetFlavour => !draft.peerings.is_empty(), + Kind::AddExpose => sides().any(|exposes| exposes.len() < usize::from(MAX_EXPOSES)), + Kind::RemoveExpose => sides().any(|exposes| exposes.len() > 1), + } + } +} + +fn draw( + driver: &mut D, + draft: &Draft, + next_vpc: &mut u8, + next_peering: &mut u8, + target: usize, +) -> Option { + let mut menu: Vec = MENU + .iter() + .filter(|(kind, _)| kind.applicable(draft, *next_vpc, *next_peering)) + .flat_map(|(kind, weight)| std::iter::repeat_n(*kind, usize::from(*weight))) + .collect(); + + while !menu.is_empty() { + let kind = pick(driver, &menu)?; + let built = match kind { + Kind::AddVpc => Some(Op::AddVpc(VpcHandle(*next_vpc))), + Kind::RemoveVpc => draw_remove_vpc(driver, draft), + Kind::AddPeering => draw_add_peering(driver, draft, *next_peering, target), + Kind::RemovePeering => draw_remove_peering(driver, draft), + Kind::AddExpose => draw_add_expose(driver, draft), + Kind::RemoveExpose => draw_remove_expose(driver, draft), + Kind::SetFlavour => draw_set_flavour(driver, draft), + }; + + let Some(op) = built else { + menu.retain(|other| *other != kind); + continue; + }; + match op { + Op::AddVpc(_) => *next_vpc = next_vpc.checked_add(1)?, + Op::AddPeering { .. } => *next_peering = next_peering.checked_add(1)?, + _ => {} + } + return Some(op); + } + + None +} + +fn draw_remove_vpc(driver: &mut D, draft: &Draft) -> Option { + let vpcs: Vec = draft.vpcs().collect(); + Some(Op::RemoveVpc(pick(driver, &vpcs)?)) +} + +fn draw_add_peering( + driver: &mut D, + draft: &Draft, + next: u8, + target: usize, +) -> Option { + let vpcs: Vec = draft.vpcs().collect(); + let left = pick(driver, &vpcs)?; + let within = driver.produce::()?; + + let free = |candidates: Vec| -> Vec { + candidates + .into_iter() + .filter(|right| *right != left && draft.peering_between(left, *right).is_none()) + .collect() + }; + + let components = draft.components(); + let component_of = |vpc: VpcHandle| -> &[VpcHandle] { + components + .iter() + .find(|component| component.contains(&vpc)) + .map_or(&[][..], Vec::as_slice) + }; + let component = component_of(left).to_vec(); + + let peered = components + .iter() + .filter(|component| component.len() > 1) + .count(); + let left_alone = component.len() <= 1; + let want_more = peered < target; + + let allowed = |right: VpcHandle| { + let right_alone = component_of(right).len() <= 1; + match (left_alone, right_alone) { + (true, true) => true, + (true, false) | (false, true) => !want_more, + (false, false) => peered > target, + } + }; + + let inside = free(component.clone()); + let outside: Vec = free(vpcs) + .into_iter() + .filter(|right| !component.contains(right)) + .filter(|right| allowed(*right)) + .collect(); + + let candidates = if (within && !inside.is_empty()) || outside.is_empty() { + inside + } else { + outside + }; + + Some(Op::AddPeering { + handle: PeeringHandle(next), + left, + right: pick(driver, &candidates)?, + }) +} + +fn draw_remove_peering(driver: &mut D, draft: &Draft) -> Option { + let peerings: Vec = draft.peerings().map(|(handle, _)| handle).collect(); + Some(Op::RemovePeering(pick(driver, &peerings)?)) +} + +fn draw_add_expose(driver: &mut D, draft: &Draft) -> Option { + let room: Vec<(PeeringHandle, Side)> = draft + .peerings() + .flat_map(|(handle, spec)| { + [Side::Left, Side::Right] + .into_iter() + .filter(move |side| spec.exposes(*side).len() < usize::from(MAX_EXPOSES)) + .map(move |side| (handle, side)) + }) + .collect(); + let (peering, side) = pick(driver, &room)?; + let spec = draft.peerings.get(&peering)?; + let slot = (0..MAX_EXPOSES).find(|slot| spec.exposes(side).iter().all(|e| e.slot != *slot))?; + + Some(Op::AddExpose { + peering, + side, + slot, + flavour: draw_flavour(driver, spec, side)?, + }) +} + +fn draw_remove_expose(driver: &mut D, draft: &Draft) -> Option { + let removable: Vec<(PeeringHandle, Side, u8)> = draft + .peerings() + .flat_map(|(handle, spec)| { + [Side::Left, Side::Right] + .into_iter() + .filter(move |side| spec.exposes(*side).len() > 1) + .flat_map(move |side| { + spec.exposes(side) + .iter() + .map(move |expose| (handle, side, expose.slot)) + }) + }) + .collect(); + let (peering, side, slot) = pick(driver, &removable)?; + + Some(Op::RemoveExpose { + peering, + side, + slot, + }) +} + +fn draw_set_flavour(driver: &mut D, draft: &Draft) -> Option { + let exposes: Vec<(PeeringHandle, Side, u8)> = draft + .peerings() + .flat_map(|(handle, spec)| { + [Side::Left, Side::Right].into_iter().flat_map(move |side| { + spec.exposes(side) + .iter() + .map(move |expose| (handle, side, expose.slot)) + }) + }) + .collect(); + let (peering, side, slot) = pick(driver, &exposes)?; + let spec = draft.peerings.get(&peering)?; + + Some(Op::SetFlavour { + peering, + side, + slot, + flavour: draw_flavour(driver, spec, side)?, + }) +} + +fn draw_flavour(driver: &mut D, spec: &PeeringSpec, side: Side) -> Option { + if driver.produce::()? && !spec.has_stateful(side.other()) { + Some(Flavour::Masquerade) + } else { + Some(Flavour::Forward) + } +} + +fn pick(driver: &mut D, items: &[T]) -> Option { + if items.is_empty() { + return None; + } + let index = usize::from(driver.produce::()?) % items.len(); + items.get(index).copied() +} + +impl Draft { + #[must_use] + pub fn restricted(&self, footprint: &Footprint) -> Self { + Self { + vpcs: self + .vpcs + .iter() + .copied() + .filter(|vpc| !footprint.vpcs.contains(vpc)) + .collect(), + peerings: self + .peerings + .iter() + .filter(|(peering, _)| !footprint.peerings.contains(peering)) + .map(|(peering, spec)| (*peering, spec.clone())) + .collect(), + } + } + + #[must_use] + pub fn references_resolve(&self) -> bool { + self.peerings + .values() + .all(|spec| self.vpcs.contains(&spec.left) && self.vpcs.contains(&spec.right)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use bolero::check; + // Vacuity counters, read after `check!()` returns and so outside any execution: they are + // process-lifetime by construction, and a facade atomic in a `static` does not even compile + // under loom, whose `AtomicUsize::new` is not `const`. + use concurrency::process_global::atomic::{AtomicUsize, Ordering::Relaxed}; + + static DRAWN: [AtomicUsize; 7] = [ + AtomicUsize::new(0), + AtomicUsize::new(0), + AtomicUsize::new(0), + AtomicUsize::new(0), + AtomicUsize::new(0), + AtomicUsize::new(0), + AtomicUsize::new(0), + ]; + + const KINDS: [&str; 7] = [ + "AddVpc", + "RemoveVpc", + "AddPeering", + "RemovePeering", + "AddExpose", + "RemoveExpose", + "SetFlavour", + ]; + + fn record(op: Op) { + let index = match op { + Op::AddVpc(_) => 0, + Op::RemoveVpc(_) => 1, + Op::AddPeering { .. } => 2, + Op::RemovePeering(_) => 3, + Op::AddExpose { .. } => 4, + Op::RemoveExpose { .. } => 5, + Op::SetFlavour { .. } => 6, + }; + DRAWN[index].fetch_add(1, Relaxed); + } + + fn assert_every_kind_drawn() { + for (kind, count) in KINDS.iter().zip(&DRAWN) { + assert!( + count.load(Relaxed) > 0, + "no {kind} was ever drawn, so nothing here tested it" + ); + } + } + + #[test] + fn every_sequence_builds_a_valid_configuration() { + let masquerading = AtomicUsize::new(0); + let forwarding = AtomicUsize::new(0); + + check!() + .with_generator(Sequence::default()) + .for_each(|ops| { + let mut draft = Draft::new(); + for (index, op) in ops.iter().enumerate() { + record(*op); + assert!( + op.apply(&mut draft).is_some(), + "the generator drew {op:?} at {index}, which does not apply" + ); + assert!( + draft.references_resolve(), + "{op:?} at {index} left a peering naming an absent vpc: {draft:?}" + ); + } + + for (_, spec) in draft.peerings() { + for side in [Side::Left, Side::Right] { + for expose in spec.exposes(side) { + match expose.flavour() { + Flavour::Masquerade => &masquerading, + Flavour::Forward => &forwarding, + } + .fetch_add(1, Relaxed); + } + } + } + + let overlay = draft + .overlay() + .unwrap_or_else(|e| panic!("{ops:?} does not assemble: {e}")); + if let Err(e) = overlay.validate() { + panic!("{ops:?} builds a configuration the validator refuses: {e}"); + } + }); + + assert_every_kind_drawn(); + assert!( + masquerading.load(Relaxed) > 0 && forwarding.load(Relaxed) > 0, + "only one flavour of expose was ever built (masquerade {}, forward {}), so the \ + combination rules were not exercised", + masquerading.load(Relaxed), + forwarding.load(Relaxed) + ); + } + + #[test] + fn undo_restores_the_configuration() { + let checked = AtomicUsize::new(0); + + check!() + .with_generator(Sequence::default()) + .for_each(|ops| { + let mut draft = Draft::new(); + for op in ops { + let before = draft.clone(); + let undo = op + .apply(&mut draft) + .unwrap_or_else(|| panic!("{op:?} does not apply")); + let after = draft.clone(); + + let mut reversed = after; + undo.apply(&mut reversed); + assert_eq!( + reversed, before, + "{op:?} then {undo:?} is not the configuration it started from" + ); + checked.fetch_add(1, Relaxed); + } + }); + + assert!( + checked.load(Relaxed) > 0, + "no operation was ever undone: every drawn sequence was empty" + ); + } + + #[test] + fn an_operation_leaves_its_complement_alone() { + check!() + .with_generator(Sequence::default()) + .for_each(|ops| { + let mut draft = Draft::new(); + for op in ops { + let writes = op.writes(&draft); + let before = draft.restricted(&writes); + op.apply(&mut draft) + .unwrap_or_else(|| panic!("{op:?} does not apply")); + assert_eq!( + draft.restricted(&writes), + before, + "{op:?} claims to write {writes:?} and changed something outside it" + ); + } + }); + } + + #[test] + fn independent_operations_commute() { + let swapped = AtomicUsize::new(0); + let considered = AtomicUsize::new(0); + + check!() + .with_generator(Sequence::default()) + .for_each(|ops| { + for index in 0..ops.len().saturating_sub(1) { + let mut draft = Sequence::fold(&ops[..index]); + let (first, second) = (ops[index], ops[index + 1]); + + considered.fetch_add(1, Relaxed); + if conflict(&draft, first, second) { + continue; + } + + let mut straight = draft.clone(); + first.apply(&mut straight); + second.apply(&mut straight); + + assert!( + second.apply(&mut draft).is_some(), + "{second:?} does not conflict with {first:?} but only applies after it" + ); + assert!( + first.apply(&mut draft).is_some(), + "{first:?} does not conflict with {second:?} but only applies before it" + ); + + assert_eq!( + draft, straight, + "{first:?} and {second:?} have disjoint footprints and do not commute" + ); + swapped.fetch_add(1, Relaxed); + } + }); + + assert!( + swapped.load(Relaxed) > 0, + "no adjacent pair was ever found independent out of {} considered, so this asserted \ + nothing", + considered.load(Relaxed) + ); + } + + fn conflict(draft: &Draft, first: Op, second: Op) -> bool { + let (rw1, ww1) = (first.reads(), first.writes(draft)); + let (rw2, ww2) = (second.reads(), second.writes(draft)); + ww1.intersects(&ww2) || ww1.intersects(&rw2) || rw1.intersects(&ww2) + } + + #[test] + fn sequences_produce_disjoint_peering_components() { + let split = AtomicUsize::new(0); + let joined = AtomicUsize::new(0); + + check!() + .with_generator(Sequence::default()) + .for_each(|ops| { + let draft = Sequence::fold(ops); + let peered = draft + .components() + .into_iter() + .filter(|component| { + draft + .peerings() + .any(|(_, spec)| component.contains(&spec.left)) + }) + .count(); + if peered >= 2 { + split.fetch_add(1, Relaxed); + } else { + joined.fetch_add(1, Relaxed); + } + }); + + assert!( + split.load(Relaxed) > 0, + "every sequence produced at most one peered component out of {} drawn, so nothing here \ + will ever be able to state isolation", + joined.load(Relaxed) + ); + } + + #[test] + fn removing_a_vpc_removes_exactly_what_referred_to_it() { + let removed = AtomicUsize::new(0); + let cascaded = AtomicUsize::new(0); + + check!() + .with_generator(Sequence::default()) + .for_each(|ops| { + let mut draft = Draft::new(); + for op in ops { + if let Op::RemoveVpc(handle) = *op { + let doomed: BTreeSet = draft + .peerings() + .filter(|(_, spec)| spec.touches(handle)) + .map(|(peering, _)| peering) + .collect(); + let survivors: BTreeSet = draft + .peerings() + .map(|(peering, _)| peering) + .filter(|peering| !doomed.contains(peering)) + .collect(); + + op.apply(&mut draft); + + let left: BTreeSet = + draft.peerings().map(|(peering, _)| peering).collect(); + assert_eq!( + left, survivors, + "removing {handle:?} should have taken {doomed:?} and nothing else" + ); + removed.fetch_add(1, Relaxed); + cascaded.fetch_add(doomed.len(), Relaxed); + } else { + op.apply(&mut draft); + } + } + }); + + assert!(removed.load(Relaxed) > 0, "no vpc was ever removed"); + assert!( + cascaded.load(Relaxed) > 0, + "every removed vpc was unpeered, so the cascade was never exercised across {} removals", + removed.load(Relaxed) + ); + } +} diff --git a/config/src/external/overlay/mod.rs b/config/src/external/overlay/mod.rs index 44f58a1538..fdf8d85186 100644 --- a/config/src/external/overlay/mod.rs +++ b/config/src/external/overlay/mod.rs @@ -4,6 +4,8 @@ //! Dataplane configuration model: overlay configuration pub mod acl; +#[cfg(any(test, feature = "bolero"))] +pub mod algebra; pub mod tests; pub mod validation_tests; pub mod vpc; diff --git a/config/src/external/overlay/vpcpeering.rs b/config/src/external/overlay/vpcpeering.rs index a83644c567..fc037e5875 100644 --- a/config/src/external/overlay/vpcpeering.rs +++ b/config/src/external/overlay/vpcpeering.rs @@ -1100,9 +1100,13 @@ pub mod contract { use super::{VpcExpose, VpcExposeNatConfig, VpcManifest, VpcPeering, VpcPeeringTable}; use crate::ConfigError; + use crate::external::overlay::acl::{ + Acl, AclAction, AclPattern, AclProtoMatch, AclRule, AclScope, + }; use crate::external::overlay::vpc::{Vpc, VpcTable}; use crate::external::overlay::{Overlay, ValidatedOverlay}; use bolero::{Driver, ValueGenerator}; + use lpm::prefix::PrefixPortsSet; use lpm::prefix::{ IpPrefix, Ipv4Prefix, Ipv6Prefix, L4Protocol, PortRange, Prefix, PrefixWithOptionalPorts, }; @@ -1604,6 +1608,87 @@ pub mod contract { Ok(Overlay::new(vpc_table, peerings)) } + #[must_use] + pub const fn peer_vni(n: u8) -> u32 { + REMOTE_VNI + n as u32 + } + + #[must_use] + pub fn peer_prefix(n: u8) -> Prefix { + let octet = u16::from(n) + 1; + format!("10.{octet}.0.0/16") + .parse() + .unwrap_or_else(|_| unreachable!("a well-formed prefix")) + } + + /// # Errors + /// + /// Returns the [`ConfigError`] from the first vpc or peering that does not build. + /// + /// # Panics + /// + /// Panics unless `peers` is between 1 and 254: zero peers has nowhere to send, and + /// more than 254 exhausts the distinct second octets the fixture addresses peers with. + pub fn overlay_with_peers(local: Prefix, peers: u8) -> Result { + assert!(peers > 0, "a local vpc with no peers has nowhere to send"); + assert!(peers <= 254, "more peers than distinct second octets"); + + let mut vpc_table = VpcTable::new(); + vpc_table.add(Vpc::new("VPC-1", "AAAAA", LOCAL_VNI)?)?; + + let mut peerings = VpcPeeringTable::new(); + for n in 0..peers { + let name = format!("PEER-{n}"); + vpc_table.add(Vpc::new(&name, &format!("P{n:04}"), peer_vni(n))?)?; + peerings.add(VpcPeering::with_default_group( + &format!("VPC-1--{name}"), + VpcManifest::new("VPC-1").exposing(VpcExpose::empty().ip(local.into())), + VpcManifest::new(&name).exposing(VpcExpose::empty().ip(peer_prefix(n).into())), + ))?; + } + + Ok(Overlay::new(vpc_table, peerings)) + } + + #[must_use] + pub fn peering_acl(default: AclAction, allow: AclProtoMatch) -> Acl { + Acl::new( + default, + vec![AclRule { + name: "allow-one-protocol".to_owned(), + from: "VPC-1".to_owned(), + to: "VPC-2".to_owned(), + action: match default { + AclAction::Allow => AclAction::Deny, + AclAction::Deny => AclAction::Allow, + }, + pattern: AclPattern { + src: PrefixPortsSet::new(), + dst: PrefixPortsSet::new(), + src_any_ports: Vec::new(), + dst_any_ports: Vec::new(), + proto: allow, + }, + scope: AclScope::Packet, + log: false, + }], + ) + } + + /// # Errors + /// + /// Returns the [`ConfigError`] from the underlying overlay if the exposes do not build. + pub fn overlay_with_exposes_and_acl( + exposes: Vec, + acl: Option<&Acl>, + ) -> Result { + let mut overlay = overlay_with_exposes(exposes)?; + for peering in overlay.peering_table.values_mut() { + peering.acl = acl.cloned(); + } + Ok(overlay) + } + #[must_use] pub fn nat_config(expose: &VpcExpose) -> Option<&VpcExposeNatConfig> { expose.nat_config() diff --git a/dataplane/Cargo.toml b/dataplane/Cargo.toml index 78f13d5cd6..dc3b553da5 100644 --- a/dataplane/Cargo.toml +++ b/dataplane/Cargo.toml @@ -58,11 +58,18 @@ vpcmap = { workspace = true } [dev-dependencies] clock = { workspace = true, features = ["virtual"] } # internal -net = { workspace = true, features = ["test_buffer"] } +config = { workspace = true, features = ["bolero"] } +lpm = { workspace = true, features = ["testing"] } +net = { workspace = true, features = ["bolero", "builder", "test_buffer", "test_meta"] } routing = { workspace = true, features = ["testing"] } # external +bolero = { workspace = true, default-features = false, features = ["alloc"] } n-vm = { workspace = true } -tokio = { workspace = true } +tokio = { workspace = true, features = ["macros", "rt", "test-util", "time"] } tracing = { workspace = true } tracing-subscriber = { workspace = true } +tracing-test = { workspace = true, features = [] } + +[target.'cfg(not(miri))'.dev-dependencies] +dpdk = { workspace = true, features = ["test"] } diff --git a/dataplane/src/packet_processor/fuzz.rs b/dataplane/src/packet_processor/fuzz.rs new file mode 100644 index 0000000000..98b83ec693 --- /dev/null +++ b/dataplane/src/packet_processor/fuzz.rs @@ -0,0 +1,3474 @@ +// SPDX-License-Identifier: Apache-2.0 +// Copyright Open Network Fabric Authors + +#![cfg(test)] +#![cfg(not(miri))] + +use acl_filter::{AclFilter, AclFilterContext, AclFilterContextWriter}; +use concurrency::sync::{Arc, Mutex}; +use config::external::overlay::acl::Acl; +use config::external::overlay::vpcpeering::VpcExpose; +use config::external::overlay::vpcpeering::contract::{ + LOCAL_VNI, REMOTE_VNI, overlay_with_exposes_and_acl, +}; +use config::external::overlay::{Overlay, ValidatedOverlay}; +use flow_entry::flow_table::{FlowLookup, FlowTable}; +use flow_filter::{FlowFilter, FlowFilterContext, FlowFilterContextWriter}; +use lpm::prefix::Prefix; +use nat::masquerade::{MasqueradeConfig, NatAllocatorWriter}; +use nat::portfw::{PortForwarder, PortFwTableWriter}; +use nat::static_nat::NatTablesWriter; +use nat::static_nat::setup::build_nat_configuration; +use nat::{IcmpErrorHandler, Masquerade, StaticNat}; +use net::buffer::{PacketBufferMut, TestBuffer}; +use net::eth::mac::{Mac, SourceMac}; +use net::interface::InterfaceIndex; +use net::packet::{DoneReason, Packet, VpcDiscriminant}; +use net::vxlan::Vni; +use pipeline::{DynPipeline, NetworkFunction}; +use routing::testing::RouterTables; +use routing::testing::{FibGroup, FwAction, NhopKey, RouteOrigin}; +use routing::{EgressObject, FibEntry, PktInstruction, ResolvedEncapsulation, ResolvedVxlan, Vtep}; +use std::net::IpAddr; + +use super::egress::Egress; +use super::ingress::Ingress; +use super::ipforward::IpForwarder; + +pub(crate) struct Fabric { + pipeline: DynPipeline, + flow_table: Arc, + _flow_filter: FlowFilterContextWriter, + _acl: AclFilterContextWriter, + _static_nat: NatTablesWriter, + _portfw: PortFwTableWriter, + _masquerade: NatAllocatorWriter, + _tables: Option, + translations: Arc>, + next_id: u64, +} + +impl Fabric { + pub(crate) fn build(exposes: &[VpcExpose]) -> Option { + Self::build_with_acl(exposes, None) + } + + pub(crate) fn build_with_acl(exposes: &[VpcExpose], acl: Option<&Acl>) -> Option { + Self::assemble(exposes, acl, None) + } + + pub(crate) fn routed(exposes: &[VpcExpose], acl: Option<&Acl>) -> Option { + Self::assemble( + exposes, + acl, + Some(topology(&[vni(LOCAL_VNI), vni(REMOTE_VNI)])), + ) + } + + pub(crate) fn routed_over(overlay: &Overlay, tables: RouterTables) -> Option { + Some(Self::with_overlay( + &overlay.clone().validate().ok()?, + Some(tables), + )) + } + + pub(crate) fn routed_over_validated(overlay: &ValidatedOverlay, tables: RouterTables) -> Self { + Self::with_overlay(overlay, Some(tables)) + } + + fn assemble( + exposes: &[VpcExpose], + acl: Option<&Acl>, + tables: Option, + ) -> Option { + let overlay = overlay_with_exposes_and_acl(exposes.to_vec(), acl) + .ok()? + .validate() + .ok()?; + Some(Self::with_overlay(&overlay, tables)) + } + + fn with_overlay(overlay: &ValidatedOverlay, tables: Option) -> Self { + let translations = Arc::new(Mutex::new(Translations::declaring(overlay))); + let flow_table = Arc::new(FlowTable::default()); + let mut pipeline = DynPipeline::new(); + + if let Some(tables) = &tables { + pipeline = pipeline.add_stage(Ingress::new("ingress", tables.interfaces())); + pipeline = pipeline.add_stage(IpForwarder::new("ip-forward-1", tables.fibs())); + pipeline = pipeline.add_stage(Checkpoint::new( + "after ip-forward-1", + contract::decapsulated, + )); + } + + pipeline = pipeline.add_stage(IcmpErrorHandler::new(flow_table.clone())); + pipeline = pipeline.add_stage(FlowLookup::new("flow-lookup", flow_table.clone())); + + let flow_filter = FlowFilterContextWriter::new(); + flow_filter.store( + FlowFilterContext::try_from(overlay).expect("a validated overlay lowers to tables"), + ); + pipeline = pipeline.add_stage(FlowFilter::new("flow-filter", flow_filter.get_reader())); + pipeline = pipeline.add_stage(Checkpoint::new("after flow-filter", contract::placed)); + + let acl = AclFilterContextWriter::new(); + acl.store(AclFilterContext::try_from(overlay).expect("a validated overlay lowers to acls")); + pipeline = pipeline.add_stage(AclFilter::new("acl-filter", acl.get_reader())); + + let mut static_nat = NatTablesWriter::new(); + static_nat.update_nat_tables( + build_nat_configuration(overlay.vpc_table()) + .expect("a validated overlay lowers to nat"), + ); + pipeline = pipeline.add_stage(StaticNat::with_reader( + "static-nat", + static_nat.get_reader(), + )); + + let mut portfw = PortFwTableWriter::new(); + portfw + .update_from_vpc_table(overlay.vpc_table()) + .expect("a validated overlay lowers to port forwarding"); + pipeline = pipeline.add_stage(PortForwarder::new( + "port-forwarder", + portfw.reader(), + flow_table.clone(), + )); + + let mut masquerade = NatAllocatorWriter::new(); + masquerade.update_nat_allocator( + MasqueradeConfig::new(overlay.vpc_table()).set_randomize(false), + 1, + &flow_table, + ); + pipeline = pipeline.add_stage(Checkpoint::new( + "before masquerade", + contract::ready_to_translate, + )); + let recording = translations.clone(); + pipeline = pipeline.add_stage(Checkpoint::new( + "before masquerade", + move |_: &str, packet: &Packet| { + recording.lock().before(packet); + }, + )); + pipeline = pipeline.add_stage(Masquerade::new( + "masquerade", + flow_table.clone(), + masquerade.get_reader(), + )); + + let checking = translations.clone(); + pipeline = pipeline.add_stage(Checkpoint::new( + "after masquerade", + move |at: &str, packet: &Packet| { + checking.lock().after(at, packet); + }, + )); + + if let Some(tables) = &tables { + pipeline = pipeline.add_stage(IpForwarder::new("ip-forward-2", tables.fibs())); + pipeline = pipeline.add_stage(Egress::new( + "egress", + tables.interfaces(), + tables.adjacencies(), + )); + pipeline = pipeline.add_stage(Checkpoint::new("after egress", contract::finished)); + } + + Self { + pipeline, + flow_table, + _flow_filter: flow_filter, + _acl: acl, + _static_nat: static_nat, + _portfw: portfw, + _masquerade: masquerade, + _tables: tables, + translations, + next_id: 0, + } + } + + pub(crate) fn send(&mut self, packet: Packet) -> Packet { + let mut out = self.send_batch(vec![packet]); + assert_eq!(out.len(), 1, "the pipeline did not return the packet"); + out.pop().unwrap_or_else(|| unreachable!()) + } + + pub(crate) fn flows(&self) -> Option { + self.flow_table.len() + } + + pub(crate) fn send_batch( + &mut self, + mut packets: Vec>, + ) -> Vec> { + let sent = packets.len(); + for packet in &mut packets { + self.next_id += 1; + packet.meta_mut().test = Some(Box::new(net::packet::TestMeta { id: self.next_id })); + } + self.translations.lock().clear(); + + let out: Vec<_> = self.pipeline.process(packets.into_iter()).collect(); + assert_eq!(out.len(), sent, "the pipeline did not return every packet"); + out + } +} + +pub(crate) fn local() -> VpcDiscriminant { + VpcDiscriminant::VNI(Vni::new_checked(LOCAL_VNI).unwrap_or_else(|_| unreachable!())) +} + +pub(crate) fn remote() -> VpcDiscriminant { + VpcDiscriminant::VNI(Vni::new_checked(REMOTE_VNI).unwrap_or_else(|_| unreachable!())) +} + +pub(crate) fn arrive(packet: &mut Packet, src: VpcDiscriminant) { + packet.meta_mut().set_overlay(true); + packet.meta_mut().src_vpcd = Some(src); + packet.meta_mut().set_keep(true); +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum Verdict { + Forwarded { + dst_vpcd: Option, + src: Option, + dst: Option, + }, + Delivered { + oif: Option, + src: Option, + dst: Option, + }, + Dropped(DoneReason), +} + +pub(crate) fn verdict(packet: &Packet) -> Verdict { + match packet.get_done() { + Some(DoneReason::Delivered) => Verdict::Delivered { + oif: packet.meta().oif, + src: packet.ip_source(), + dst: packet.ip_destination(), + }, + Some(reason) => Verdict::Dropped(reason), + None => Verdict::Forwarded { + dst_vpcd: packet.meta().dst_vpcd, + src: packet.ip_source(), + dst: packet.ip_destination(), + }, + } +} + +pub(crate) const MAX_INPUT_LEN: usize = 65536; + +#[cfg(test)] +fn largest_draw(generator: &G) -> usize { + use bolero::generator::bolero_generator::driver::{Options, bytes::Driver}; + + const SAMPLES: usize = 64; + const UNLIMITED: usize = 1 << 20; + + let options = Options::default().with_max_len(UNLIMITED); + let mut state: u64 = 0x243f_6a88_85a3_08d3; + let mut bytes = vec![0u8; UNLIMITED]; + (0..SAMPLES) + .map(|_| { + for byte in &mut bytes { + state = state + .wrapping_mul(6_364_136_223_846_793_005) + .wrapping_add(1); + #[allow(clippy::cast_possible_truncation)] + { + *byte = (state >> 33) as u8; + } + } + let mut driver = Driver::new(&bytes[..], &options); + assert!( + generator.generate(&mut driver).is_some(), + "the generator gave up on a budget it cannot exhaust" + ); + UNLIMITED - driver.as_slice().len() + }) + .max() + .unwrap_or_else(|| unreachable!("SAMPLES is not zero")) +} + +#[cfg(test)] +fn assert_within_budget(name: &str, generator: &G) { + let largest = largest_draw(generator); + eprintln!("{name}: largest draw {largest} of {MAX_INPUT_LEN} bytes"); + assert!( + largest * 2 <= MAX_INPUT_LEN, + "{name} drew {largest} bytes, over half of the {MAX_INPUT_LEN} budget. Raise \ + MAX_INPUT_LEN and `just fuzz`'s `-l`, or the tail of every batch will be zeros." + ); +} + +#[cfg(test)] +fn assert_covered(covered: bool, what: &str) { + if (cfg!(instrumented) || cfg!(emulated)) && !covered { + // Coverage and emulation are both slow enough per case that bolero's + // budget buys a sample too small for "did anything reach this branch" + // to mean anything: qemu-user gets a couple of orders of magnitude + // fewer cases than a native run, and instrumentation is not far behind. + // Say so and carry on -- the point of those runs is the line counts and + // the target's own behaviour, and failing here loses both. + eprintln!("{what} -- not asserted: too few cases under this build"); + return; + } + assert!( + covered, + "{what}. Check for a `__fuzz__` corpus beside this test before reading further: replaying \ + one can spend the budget on inputs chosen for being unusual. Move it aside and re-run to \ + tell that apart from a real gap." + ); +} + +#[cfg(test)] +#[derive(Default)] +pub(crate) struct Translations { + was: std::collections::HashMap, + from: std::collections::HashMap, + given: std::collections::HashMap, + declared: Vec, +} + +#[cfg(test)] +impl Translations { + fn before(&mut self, packet: &Packet) { + if packet.is_done() || !packet.meta().requires_masquerade() { + return; + } + let (Some(test), Ok(key)) = (packet.meta().test.as_ref(), net::FlowKey::try_from(packet)) + else { + return; + }; + self.was.insert(test.id, key); + if let Some(source) = packet.ip_source() { + self.from.insert(test.id, source); + } + } + + fn after(&mut self, at: &str, packet: &Packet) { + if packet.is_done() { + return; + } + let Some(test) = packet.meta().test.as_ref() else { + return; + }; + let Some(was) = self.was.get(&test.id).copied() else { + return; + }; + let (Some(source), Some(port)) = (packet.ip_source(), packet.transport_src_port()) else { + return; + }; + let now = (source, port.get()); + if let Some(before) = self.given.insert(was, now) { + assert_eq!( + before, now, + "{at}: two packets of one flow in one burst were given different public tuples, \ + so the burst allocated more than once for it" + ); + } + + if self + .from + .get(&test.id) + .is_some_and(|before| *before != source) + { + assert!( + self.declared.iter().any(|p| p.covers_addr(&source)), + "{at}: masquerade rewrote a source to {source}, which no declared public range \ + covers" + ); + } + } + + fn clear(&mut self) { + self.was.clear(); + self.from.clear(); + self.given.clear(); + } + + fn declaring(overlay: &ValidatedOverlay) -> Self { + let mut declared = Vec::new(); + for vpc in overlay.vpc_table().values() { + for peering in vpc.peerings() { + for manifest in [peering.local(), peering.remote()] { + for expose in manifest.valexp() { + declared.extend( + expose + .public_ips() + .into_iter() + .map(lpm::prefix::PrefixWithOptionalPorts::prefix), + ); + } + } + } + } + Self { + declared, + ..Self::default() + } + } +} + +#[cfg(test)] +pub(crate) trait Load { + fn next(&mut self) -> Option>; + + fn observe(&mut self, got: &Packet); + + fn finished(&self) -> bool; + + fn checked(&self) -> bool; + + fn describe(&self) -> String; +} + +#[cfg(test)] +pub(crate) fn drive(fabric: &mut Fabric, load: &mut dyn Load) { + for _ in 0..64 { + if load.finished() { + return; + } + let Some(packet) = load.next() else { + panic!( + "a load is neither finished nor willing to send: {}", + load.describe() + ); + }; + let out = fabric.send(packet); + load.observe(&out); + } + panic!("a load did not finish in 64 steps: {}", load.describe()); +} + +#[cfg(test)] +#[derive(Debug, Clone, Copy)] +pub(crate) struct Pick { + pub(crate) load: u8, + pub(crate) take: u8, +} + +#[cfg(test)] +pub(crate) type Poll = Vec; + +#[cfg(test)] +pub(crate) fn run_schedule( + fabric: &mut Fabric, + loads: &mut [Box], + schedule: &[Poll], +) -> Vec> { + let mut bursts = Vec::new(); + for poll in schedule { + let mut burst = Vec::new(); + let mut origin = Vec::new(); + for pick in poll { + if loads.is_empty() { + break; + } + let which = usize::from(pick.load) % loads.len(); + for _ in 0..pick.take { + let Some(packet) = loads[which].next() else { + break; + }; + burst.push(packet); + origin.push(which); + } + } + if burst.is_empty() { + continue; + } + for (answer, which) in fabric.send_batch(burst).iter().zip(&origin) { + loads[*which].observe(answer); + } + bursts.push(origin); + } + + for load in loads { + drive(fabric, load.as_mut()); + } + bursts +} + +#[cfg(test)] +pub(crate) mod derive { + use super::routed::{Blast, Conversation, Inbound}; + use super::*; + use config::external::overlay::ValidatedOverlay; + use config::external::overlay::vpcpeering::ValidatedExpose; + use lpm::prefix::{Prefix, PrefixWithOptionalPorts}; + + #[derive(Debug, Clone, Copy)] + pub(crate) struct Vary { + pub(crate) host: u8, + pub(crate) port: u8, + pub(crate) sport: u16, + pub(crate) dport: u16, + pub(crate) burst: u8, + pub(crate) blast: bool, + } + + fn host_in(prefix: Prefix, n: u8) -> Option { + let full = if matches!(prefix.as_address(), IpAddr::V4(_)) { + 32 + } else { + 128 + }; + let width = u32::from(full - prefix.length()); + if width == 0 { + return Some(prefix.as_address()); + } + let span = 1u128.checked_shl(width.min(7))?; + let offset = u128::from(n) % span; + match prefix.as_address() { + IpAddr::V4(base) => { + let raw = u32::from(base).checked_add(u32::try_from(offset).ok()?)?; + Some(IpAddr::V4(raw.into())) + } + IpAddr::V6(base) => { + let raw = u128::from(base).checked_add(offset)?; + Some(IpAddr::V6(raw.into())) + } + } + } + + fn port_in(entry: &PrefixWithOptionalPorts, n: u8, fallback: u16) -> u16 { + entry.ports().map_or(fallback.max(1), |range| { + let span = u32::from(range.end()) - u32::from(range.start()) + 1; + let offset = u32::from(n) % span; + u16::try_from(u32::from(range.start()) + offset).unwrap_or(range.start()) + }) + } + + fn peer_of( + peering: &config::external::overlay::vpc::ValidatedPeering, + n: u8, + usable: fn(&ValidatedExpose) -> bool, + ) -> Option { + peering + .remote() + .valexp() + .iter() + .filter(|expose| usable(expose)) + .flat_map(|expose| expose.public_ips().into_iter()) + .find_map(|entry| host_in(entry.prefix(), n)) + } + + pub(crate) fn loads_for(overlay: &ValidatedOverlay, vary: &[Vary]) -> Vec> { + let mut loads: Vec> = Vec::new(); + let mut nth = 0usize; + for vpc in overlay.vpc_table().values() { + for peering in vpc.peerings() { + let path = super::routed::Path::new(vpc.vni(), peering.remote_vni()); + for expose in peering.local().valexp() { + let Some(v) = vary.get(nth % vary.len().max(1)).copied() else { + continue; + }; + nth += 1; + let outward = peer_of(peering, v.host, |expose| { + expose.can_receive_connection() && !expose.has_port_forwarding() + }); + let inward = peer_of(peering, v.host, ValidatedExpose::can_init_connection); + + if expose.has_port_forwarding() { + let (Some(outside), Some(inside_entry)) = ( + expose.public_ips().into_iter().next(), + expose.ips().into_iter().next(), + ) else { + continue; + }; + let (Some(external), Some(internal)) = ( + host_in(outside.prefix(), v.host), + host_in(inside_entry.prefix(), v.host), + ) else { + continue; + }; + let Some(peer) = inward else { + continue; + }; + loads.push(Box::new(Inbound::new( + path, + peer, + external, + port_in(outside, v.port, v.dport), + internal, + port_in(inside_entry, v.port, v.dport), + v.sport, + ))); + } else if expose.has_static_nat() || expose.is_default() { + } else { + let (Some(peer), Some(src)) = ( + outward, + expose + .ips() + .into_iter() + .next() + .and_then(|entry| host_in(entry.prefix(), v.host)), + ) else { + continue; + }; + loads.push(if v.blast { + Box::new(Blast::new(path, src, peer, v.sport, v.dport, v.burst)) + } else { + Box::new(Conversation::new(path, src, peer, v.sport, v.dport)) + as Box + }); + } + } + } + } + loads + } +} + +pub(crate) struct Checkpoint { + at: &'static str, + check: F, +} + +impl Checkpoint { + pub(crate) fn new(at: &'static str, check: F) -> Self { + Self { at, check } + } +} + +impl) + 'static> NetworkFunction + for Checkpoint +{ + fn process<'a, Input: Iterator> + 'a>( + &'a mut self, + input: Input, + ) -> impl Iterator> + 'a { + input.inspect(move |packet| (self.check)(self.at, packet)) + } +} + +mod contract { + use super::*; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + + pub(super) static JUDGED: LazyLock<[AtomicU64; 4]> = + LazyLock::new(|| std::array::from_fn(|_| AtomicU64::new(0))); + + pub(super) const DECAPSULATED: usize = 0; + pub(super) const PLACED: usize = 1; + pub(super) const READY: usize = 2; + pub(super) const FINISHED: usize = 3; + + fn judged(which: usize) { + JUDGED[which].fetch_add(1, Ordering::Relaxed); + } + + pub(super) fn decapsulated(at: &str, packet: &Packet) { + if packet.is_done() || !packet.meta().is_overlay() { + return; + } + judged(DECAPSULATED); + assert!( + packet.meta().src_vpcd.is_some(), + "{at}: overlay traffic with no source vpc discriminant" + ); + assert!( + packet.meta().vrf.is_some(), + "{at}: overlay traffic with no vrf to route it in" + ); + } + + pub(super) fn placed(at: &str, packet: &Packet) { + if packet.is_done() || !packet.meta().is_overlay() { + return; + } + judged(PLACED); + assert!( + packet.meta().dst_vpcd.is_some(), + "{at}: forwarded without a destination vpc: nothing chose where this packet goes" + ); + } + + pub(super) fn ready_to_translate(at: &str, packet: &Packet) { + if packet.is_done() || !packet.meta().requires_masquerade() { + return; + } + judged(READY); + assert!( + packet.meta().src_vpcd.is_some() && packet.meta().dst_vpcd.is_some(), + "{at}: a packet is to be masqueraded without both discriminants, which masquerade \ + itself calls a bug" + ); + } + + pub(super) fn finished(at: &str, packet: &Packet) { + judged(FINISHED); + assert!( + packet.is_done(), + "{at}: a packet left the last stage of the pipeline without a verdict" + ); + } +} + +pub(crate) const UPLINK: u32 = 1; +pub(crate) const LOCAL_VTEP: &str = "5.6.7.8"; +pub(crate) const PEER_VTEP: &str = "1.2.3.4"; +pub(crate) const GATEWAY_MAC: Mac = Mac([0x02, 0, 0, 0, 0, 0xaa]); +pub(crate) const PEER_MAC: Mac = Mac([0x02, 0, 0, 0, 0, 0xbb]); + +const UNDERLAY_VRF: u32 = 0; + +pub(crate) fn uplink() -> InterfaceIndex { + InterfaceIndex::try_new(UPLINK).unwrap_or_else(|_| unreachable!()) +} + +pub(crate) fn topology(vnis: &[Vni]) -> RouterTables { + let mut tables = RouterTables::new(); + + tables.vrf(UNDERLAY_VRF, None); + tables.interface( + uplink(), + "uplink", + SourceMac::new(GATEWAY_MAC).unwrap_or_else(|_| unreachable!()), + ); + tables.attach(uplink(), UNDERLAY_VRF); + tables.route_via( + UNDERLAY_VRF, + Prefix::expect_from((LOCAL_VTEP, 32)), + nhop(&LOCAL_VTEP.parse().unwrap_or_else(|_| unreachable!())), + &FibGroup::with_entry(FibEntry::with_inst(PktInstruction::Local(uplink()))), + ); + + for reachable in vnis { + tables.vrf(reachable.as_u32(), Some(*reachable)); + encapsulate_out_of(&mut tables, reachable.as_u32(), *reachable); + } + + tables.adjacency( + PEER_VTEP.parse().unwrap_or_else(|_| unreachable!()), + uplink(), + PEER_MAC, + ); + + tables +} + +fn encapsulate_out_of(tables: &mut RouterTables, vrfid: u32, out_vni: Vni) { + tables.vtep( + vrfid, + Vtep::with_ip_and_mac( + LOCAL_VTEP.parse().unwrap_or_else(|_| unreachable!()), + GATEWAY_MAC, + ), + ); + let peer: IpAddr = PEER_VTEP.parse().unwrap_or_else(|_| unreachable!()); + let mut out = FibEntry::with_inst(PktInstruction::Encap(ResolvedEncapsulation::Vxlan( + ResolvedVxlan { + vni: out_vni, + remote: peer, + dmac: PEER_MAC, + }, + ))); + out.add(PktInstruction::Egress(EgressObject::new( + Some(uplink()), + Some(peer), + ))); + tables.route_via( + vrfid, + Prefix::root_v4(), + nhop(&peer), + &FibGroup::with_entry(out), + ); +} + +fn nhop(address: &IpAddr) -> NhopKey { + NhopKey::new( + RouteOrigin::default(), + Some(*address), + None, + None, + FwAction::Forward, + ) +} + +fn vni(raw: u32) -> Vni { + Vni::new_checked(raw).unwrap_or_else(|_| unreachable!()) +} + +#[cfg(test)] +mod smoke { + use super::*; + use lpm::prefix::Prefix; + use net::packet::test_utils::build_test_udp_ipv4_packet; + + fn masquerade_expose() -> VpcExpose { + VpcExpose::empty() + .make_masquerade(None) + .unwrap() + .ip("1.1.0.0/16".parse::().unwrap().into()) + .as_range("2.2.0.0/16".parse::().unwrap().into()) + .unwrap() + } + + #[tokio::test] + #[dpdk::with_eal] + async fn every_contract_is_reached() { + use net::packet::test_utils::build_test_udp_ipv4_packet; + + let before: Vec = contract::JUDGED + .iter() + .map(|c| c.load(concurrency::sync::atomic::Ordering::Relaxed)) + .collect(); + + let mut fabric = Fabric::routed(&routed::exposes(), None).expect("a valid configuration"); + let inner = build_test_udp_ipv4_packet("1.1.0.1", "3.3.3.1", 1234, 80); + let out = fabric.send(routed::tunnelled(&inner)); + assert!( + matches!(verdict(&out), Verdict::Delivered { .. }), + "the fixture packet did not reach the wire: {:?}", + verdict(&out) + ); + + for (i, name) in ["decapsulated", "placed", "ready_to_translate", "finished"] + .iter() + .enumerate() + { + let after = contract::JUDGED[i].load(concurrency::sync::atomic::Ordering::Relaxed); + assert!( + after > before[i], + "the `{name}` contract judged no packet of an ordinary delivered flow: its guard \ + excuses everything, so it holds without ever having been evaluated" + ); + } + } + + #[tokio::test] + #[dpdk::with_eal] + async fn the_same_input_twice_gives_the_same_answer() { + let scenario = || { + let mut flows: Vec<_> = (0..12u8) + .flat_map(|host| (0..3u16).map(move |round| (host, round))) + .filter_map(|(host, round)| { + let src: IpAddr = format!("1.1.0.{host}").parse().ok()?; + let dst: IpAddr = "3.3.3.1".parse().ok()?; + round_trip::udp(src, dst, 4000 + round, 80).map(|p| routed::tunnelled(&p)) + }) + .collect(); + flows.extend((0..8u8).filter_map(|h| { + let src: IpAddr = format!("1.1.9.{h}").parse().ok()?; + let dst: IpAddr = "3.3.3.1".parse().ok()?; + round_trip::udp(src, dst, 5000, 80).map(|p| routed::tunnelled(&p)) + })); + flows + }; + + let answers = |fabric: &mut Fabric| { + let mut seen = Vec::new(); + let packets = scenario(); + let (singly, burst) = packets.split_at(packets.len() - 8); + for packet in singly.iter().cloned() { + let out = fabric.send(packet); + seen.push(describe(&out)); + } + for out in fabric.send_batch(burst.to_vec()) { + seen.push(describe(&out)); + } + seen + }; + + let mut once = Fabric::routed(&routed::exposes(), None).expect("a valid configuration"); + let mut again = Fabric::routed(&routed::exposes(), None).expect("a valid configuration"); + let first = answers(&mut once); + let second = answers(&mut again); + + assert!( + first.iter().any(|a| a.contains("Delivered")), + "the scenario delivered nothing, so this compares two pipelines doing nothing" + ); + assert_eq!( + first, second, + "two runs of one scenario disagreed: replay-based diagnosis cannot be trusted, and \ + neither can bolero's shrinking" + ); + } + + fn describe(packet: &Packet) -> String { + let carried = routed::inside(packet); + format!( + "{:?} {:?} {:?} {:?}", + verdict(packet), + carried.as_ref().and_then(Packet::ip_source), + carried.as_ref().and_then(Packet::ip_destination), + carried.as_ref().and_then(Packet::transport_src_port), + ) + } + + #[tokio::test] + #[dpdk::with_eal] + async fn the_harness_builds_a_pipeline_that_translates() { + let mut fabric = Fabric::build(&[masquerade_expose()]).expect("a valid configuration"); + + let mut packet = build_test_udp_ipv4_packet("1.1.0.1", "3.3.3.1", 1234, 80); + arrive(&mut packet, local()); + let out = fabric.send(packet); + + match verdict(&out) { + Verdict::Forwarded { dst_vpcd, src, dst } => { + assert_eq!(dst_vpcd, Some(remote()), "the flow filter chose no peer"); + let IpAddr::V4(src) = src.expect("no source") else { + panic!("came out IPv6") + }; + assert_eq!(src.octets()[..2], [2, 2], "not masqueraded into the range"); + assert_eq!(dst, Some("3.3.3.1".parse().unwrap())); + } + Verdict::Dropped(reason) => panic!("the packet was dropped: {reason:?}"), + Verdict::Delivered { .. } => { + unreachable!("the overlay slice has no egress stage") + } + } + } +} + +#[cfg(test)] +mod shapes { + use super::*; + use bolero::{Driver, ValueGenerator}; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::vpcpeering::contract::MasqueradeExposes; + use lpm::prefix::Prefix; + use net::headers::builder::ChainBase; + use net::headers::{Headers, TryIpv4Mut, TryIpv6Mut}; + use net::ipv4::UnicastIpv4Addr; + use net::ipv6::UnicastIpv6Addr; + use net::parse::DeParse; + use std::net::{Ipv4Addr, Ipv6Addr}; + + const MAX_EXPOSES: u8 = 2; + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub(super) enum Shape { + V4Tcp, + V4Udp, + V4Icmp, + VlanV4Tcp, + V4ExoticProto, + V6Tcp, + V6HopByHopTcp, + V6FragmentUdp, + NoIp, + } + + impl Shape { + pub(super) const ALL: [Shape; 9] = [ + Shape::V4Tcp, + Shape::V4Udp, + Shape::V4Icmp, + Shape::VlanV4Tcp, + Shape::V4ExoticProto, + Shape::V6Tcp, + Shape::V6HopByHopTcp, + Shape::V6FragmentUdp, + Shape::NoIp, + ]; + } + + const STACKS_PER_FABRIC: usize = 16; + + struct AnyStack; + + impl ValueGenerator for AnyStack { + type Output = (Shape, Headers); + + fn generate(&self, driver: &mut D) -> Option<(Shape, Headers)> { + let shape = Shape::ALL[usize::from(driver.produce::()?) % Shape::ALL.len()]; + let headers = match shape { + Shape::NoIp => ChainBase::new().eth(|_| {}).generate(driver), + Shape::V4Tcp => ChainBase::new() + .eth(|_| {}) + .ipv4(|_| {}) + .tcp(|_| {}) + .generate(driver), + Shape::V4Udp => ChainBase::new() + .eth(|_| {}) + .ipv4(|_| {}) + .udp(|_| {}) + .generate(driver), + Shape::V4Icmp => ChainBase::new() + .eth(|_| {}) + .ipv4(|_| {}) + .icmp4(|_| {}) + .generate(driver), + Shape::VlanV4Tcp => ChainBase::new() + .eth(|_| {}) + .vlan(|_| {}) + .ipv4(|_| {}) + .tcp(|_| {}) + .generate(driver), + Shape::V4ExoticProto => ChainBase::new() + .eth(|_| {}) + .ipv4(|ip| { + ip.set_next_header(net::ip::NextHeader::new(132)); + }) + .generate(driver), + Shape::V6Tcp => ChainBase::new() + .eth(|_| {}) + .ipv6(|_| {}) + .tcp(|_| {}) + .generate(driver), + Shape::V6HopByHopTcp => ChainBase::new() + .eth(|_| {}) + .ipv6(|_| {}) + .hop_by_hop(|_| {}) + .tcp(|_| {}) + .generate(driver), + Shape::V6FragmentUdp => ChainBase::new() + .eth(|_| {}) + .ipv6(|_| {}) + .fragment(|_| {}) + .udp(|_| {}) + .generate(driver), + }?; + Some((shape, headers)) + } + } + + pub(super) struct Batch; + + impl ValueGenerator for Batch { + type Output = (Vec, Vec<(Shape, Headers)>); + + fn generate(&self, driver: &mut D) -> Option { + let exposes = MasqueradeExposes(MAX_EXPOSES).generate(driver)?; + let mut stacks = Vec::with_capacity(STACKS_PER_FABRIC); + for _ in 0..STACKS_PER_FABRIC { + stacks.push(AnyStack.generate(driver)?); + } + Some((exposes, stacks)) + } + } + + pub(super) fn aim(headers: &mut Headers, private: Option) { + let Some(private) = private else { return }; + match private.as_address() { + IpAddr::V4(addr) => { + if let Some(ip) = headers.try_ipv4_mut() { + ip.set_source( + UnicastIpv4Addr::new(addr).unwrap_or_else(|_| unreachable!("prefix base")), + ); + ip.set_destination( + "3.3.3.1" + .parse::() + .unwrap_or_else(|_| unreachable!()), + ); + } + } + IpAddr::V6(addr) => { + if let Some(ip) = headers.try_ipv6_mut() { + ip.set_source( + UnicastIpv6Addr::new(addr).unwrap_or_else(|_| unreachable!("prefix base")), + ); + ip.set_destination( + "2001:db8:ffff::1" + .parse::() + .unwrap_or_else(|_| unreachable!()), + ); + } + } + } + } + + pub(super) fn wire(headers: &Headers) -> Option> { + let mut buffer = TestBuffer::new(); + headers.deparse(buffer.as_mut()).ok()?; + Packet::new(buffer).ok() + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("shapes::Batch", &Batch); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn every_shape_leaves_the_pipeline_with_a_verdict() { + static FORWARDED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static DROPPED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static BY_SHAPE: LazyLock<[AtomicU64; Shape::ALL.len()]> = + LazyLock::new(|| std::array::from_fn(|_| AtomicU64::new(0))); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Batch) + .for_each(|(exposes, stacks)| { + let Some(mut fabric) = Fabric::build(exposes) else { + return; + }; + let private = exposes.first().and_then(|e| { + e.ips + .first() + .map(lpm::prefix::PrefixWithOptionalPorts::prefix) + }); + + for (shape, headers) in stacks { + let mut headers = headers.clone(); + aim(&mut headers, private); + let Some(mut packet) = wire(&headers) else { + continue; + }; + BY_SHAPE[*shape as usize].fetch_add(1, Ordering::Relaxed); + + arrive(&mut packet, local()); + let out = fabric.send(packet); + + match verdict(&out) { + Verdict::Forwarded { dst_vpcd, .. } => { + assert_eq!( + dst_vpcd, + Some(remote()), + "forwarded without a destination VPC, on a {shape:?} stack: \ + nothing chose where this packet goes" + ); + FORWARDED.fetch_add(1, Ordering::Relaxed); + } + Verdict::Dropped(_) => { + DROPPED.fetch_add(1, Ordering::Relaxed); + } + Verdict::Delivered { .. } => { + unreachable!("the overlay slice has no egress stage") + } + } + } + }); + + let forwarded = FORWARDED.load(Ordering::Relaxed); + let dropped = DROPPED.load(Ordering::Relaxed); + let by_shape: Vec<_> = Shape::ALL + .iter() + .map(|s| format!("{s:?}={}", BY_SHAPE[*s as usize].load(Ordering::Relaxed))) + .collect(); + eprintln!( + "forwarded={forwarded} dropped={dropped}; {}", + by_shape.join(" ") + ); + + super::assert_covered( + forwarded > 0, + "no packet was ever forwarded: the harness is exercising the drop path only", + ); + super::assert_covered(dropped > 0, "no packet was ever dropped"); + for shape in Shape::ALL { + super::assert_covered( + BY_SHAPE[shape as usize].load(Ordering::Relaxed) > 0, + &format!("no {shape:?} packet ever reached the pipeline"), + ); + } + } +} + +#[cfg(test)] +mod round_trip { + use super::*; + use bolero::{Driver, ValueGenerator}; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::vpcpeering::contract::MasqueradeExposes; + use lpm::prefix::{Prefix, PrefixWithOptionalPorts}; + use net::headers::builder::HeaderStack; + use net::ipv4::UnicastIpv4Addr; + use net::ipv6::UnicastIpv6Addr; + use net::parse::DeParse; + use net::udp::UdpPort; + use std::net::{Ipv4Addr, Ipv6Addr}; + + const MAX_EXPOSES: u8 = 2; + const FLOWS_PER_FABRIC: usize = 8; + + #[derive(Debug, Clone, Copy)] + struct FlowSpec { + prefix: u8, + host: u8, + sport: u16, + dport: u16, + } + + struct Batch; + + impl ValueGenerator for Batch { + type Output = (Vec, Vec); + + fn generate(&self, driver: &mut D) -> Option { + let exposes = MasqueradeExposes(MAX_EXPOSES).generate(driver)?; + let mut flows = Vec::with_capacity(FLOWS_PER_FABRIC); + for _ in 0..FLOWS_PER_FABRIC { + flows.push(FlowSpec { + prefix: driver.produce()?, + host: driver.produce()?, + sport: driver.produce::()?.max(1), + dport: driver.produce::()?.max(1), + }); + } + Some((exposes, flows)) + } + } + + fn private_addresses(exposes: &[VpcExpose]) -> Vec { + exposes + .iter() + .flat_map(|e| e.ips.iter().map(PrefixWithOptionalPorts::prefix)) + .collect() + } + + fn peer(family_of: IpAddr) -> IpAddr { + match family_of { + IpAddr::V4(_) => IpAddr::V4( + "3.3.3.1" + .parse::() + .unwrap_or_else(|_| unreachable!()), + ), + IpAddr::V6(_) => IpAddr::V6( + "2001:db8:ffff::1" + .parse::() + .unwrap_or_else(|_| unreachable!()), + ), + } + } + + pub(super) fn udp( + src: IpAddr, + dst: IpAddr, + sport: u16, + dport: u16, + ) -> Option> { + let sport = UdpPort::new_checked(sport).ok()?; + let dport = UdpPort::new_checked(dport).ok()?; + let headers = match (src, dst) { + (IpAddr::V4(src), IpAddr::V4(dst)) => { + let src = UnicastIpv4Addr::new(src).ok()?; + HeaderStack::new() + .eth(|_| {}) + .ipv4(|ip| { + ip.set_source(src); + ip.set_destination(dst); + ip.set_ttl(64); + }) + .udp(|udp| { + udp.set_source(sport); + udp.set_destination(dport); + }) + .build_headers() + .ok()? + } + (IpAddr::V6(src), IpAddr::V6(dst)) => { + let src = UnicastIpv6Addr::new(src).ok()?; + HeaderStack::new() + .eth(|_| {}) + .ipv6(|ip| { + ip.set_source(src); + ip.set_destination(dst); + ip.set_hop_limit(64); + }) + .udp(|udp| { + udp.set_source(sport); + udp.set_destination(dport); + }) + .build_headers() + .ok()? + } + _ => return None, + }; + let mut buffer = TestBuffer::new(); + headers.deparse(buffer.as_mut()).ok()?; + Packet::new(buffer).ok() + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("round_trip::Batch", &Batch); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_translated_flow_comes_back_to_where_it_started() { + static ROUND_TRIPPED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static NOT_FORWARDED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Batch) + .for_each(|(exposes, flows)| { + let Some(mut fabric) = Fabric::build(exposes) else { + return; + }; + let privates = private_addresses(exposes); + if privates.is_empty() { + return; + } + + for flow in flows { + let prefix = privates[usize::from(flow.prefix) % privates.len()]; + let src = match prefix.as_address() { + IpAddr::V4(a) => { + let mut o = a.octets(); + o[3] = o[3].wrapping_add(flow.host % 8); + IpAddr::V4(Ipv4Addr::from(o)) + } + IpAddr::V6(a) => { + let mut o = a.octets(); + o[15] = o[15].wrapping_add(flow.host % 8); + IpAddr::V6(Ipv6Addr::from(o)) + } + }; + let dst = peer(src); + + let Some(mut request) = udp(src, dst, flow.sport, flow.dport) else { + continue; + }; + arrive(&mut request, local()); + let out = fabric.send(request); + + let Verdict::Forwarded { + src: public_src, + dst: reached, + .. + } = verdict(&out) + else { + NOT_FORWARDED.fetch_add(1, Ordering::Relaxed); + continue; + }; + let (Some(public_src), Some(reached)) = (public_src, reached) else { + continue; + }; + let public_port = out + .transport_src_port() + .unwrap_or_else(|| unreachable!("a udp packet has a source port")) + .get(); + + let Some(mut reply) = udp(reached, public_src, flow.dport, public_port) else { + continue; + }; + arrive(&mut reply, remote()); + let back = fabric.send(reply); + + match verdict(&back) { + Verdict::Forwarded { src: s, dst: d, .. } => { + assert_eq!( + d, + Some(src), + "the reply did not come back to the host that sent the request" + ); + assert_eq!(s, Some(dst), "the reply's source was rewritten"); + assert_eq!( + back.transport_dst_port().map(std::num::NonZero::get), + Some(flow.sport), + "the reply did not get the original source port back" + ); + ROUND_TRIPPED.fetch_add(1, Ordering::Relaxed); + } + Verdict::Dropped(reason) => panic!( + "the reply of a forwarded flow was dropped: {reason:?} \ + (request {src} -> {dst} became {public_src}:{public_port})" + ), + Verdict::Delivered { .. } => { + unreachable!("the overlay slice has no egress stage") + } + } + } + }); + + let round_tripped = ROUND_TRIPPED.load(Ordering::Relaxed); + eprintln!( + "round-tripped={round_tripped} not-forwarded={}", + NOT_FORWARDED.load(Ordering::Relaxed) + ); + super::assert_covered( + round_tripped > 0, + "no flow was ever forwarded, so nothing was ever checked to come back", + ); + } +} + +#[cfg(test)] +mod acl { + use super::*; + use bolero::{Driver, TypeGenerator, ValueGenerator}; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::acl::{AclAction, AclProtoMatch}; + use config::external::overlay::vpcpeering::contract::{MasqueradeExposes, peering_acl}; + use lpm::prefix::{Prefix, PrefixWithOptionalPorts}; + use net::headers::builder::ChainBase; + use net::headers::{Headers, TryIpv4Mut, TryIpv6Mut}; + use net::ip::NextHeader; + use net::ipv4::UnicastIpv4Addr; + use net::ipv6::UnicastIpv6Addr; + use net::parse::DeParse; + use std::net::{Ipv4Addr, Ipv6Addr}; + + const MAX_EXPOSES: u8 = 2; + const PACKETS_PER_FABRIC: usize = 12; + + #[derive(Debug, Clone, Copy, TypeGenerator)] + struct PacketSpec { + proto: Proto, + behind_extension: bool, + host: u8, + } + + #[derive(Debug, Clone, Copy, PartialEq, Eq, TypeGenerator)] + enum Proto { + Tcp, + Udp, + Icmp, + } + + impl Proto { + fn as_match(self) -> AclProtoMatch { + match self { + Proto::Tcp => AclProtoMatch::Tcp, + Proto::Udp => AclProtoMatch::Udp, + Proto::Icmp => AclProtoMatch::Other(NextHeader::ICMP.as_u8()), + } + } + } + + #[derive(Debug, Clone, Copy, TypeGenerator)] + enum RuleProto { + Tcp, + Udp, + Icmp, + Icmp6, + Any, + } + + impl RuleProto { + fn as_match(self) -> AclProtoMatch { + match self { + RuleProto::Tcp => AclProtoMatch::Tcp, + RuleProto::Udp => AclProtoMatch::Udp, + RuleProto::Icmp => AclProtoMatch::Other(NextHeader::ICMP.as_u8()), + RuleProto::Icmp6 => AclProtoMatch::Other(NextHeader::ICMP6.as_u8()), + RuleProto::Any => AclProtoMatch::Any, + } + } + } + + struct Batch; + + impl ValueGenerator for Batch { + type Output = (Vec, bool, RuleProto, Vec<(PacketSpec, Headers)>); + + fn generate(&self, driver: &mut D) -> Option { + let exposes = MasqueradeExposes(MAX_EXPOSES).generate(driver)?; + let v6 = exposes + .first() + .and_then(|e| e.ips.first()) + .is_some_and(|p| p.prefix().as_address().is_ipv6()); + let default_allow = driver.produce()?; + let rule_proto = RuleProto::generate(driver)?; + let mut packets = Vec::with_capacity(PACKETS_PER_FABRIC); + for _ in 0..PACKETS_PER_FABRIC { + let spec = PacketSpec::generate(driver)?; + packets.push((spec, stack(driver, spec, v6)?)); + } + Some((exposes, default_allow, rule_proto, packets)) + } + } + + fn carried(proto: Proto, v6: bool) -> AclProtoMatch { + match (proto, v6) { + (Proto::Icmp, true) => AclProtoMatch::Other(NextHeader::ICMP6.as_u8()), + (p, _) => p.as_match(), + } + } + + fn rule_matches(rule: AclProtoMatch, carried: AclProtoMatch) -> bool { + matches!(rule, AclProtoMatch::Any) || rule == carried + } + + fn stack(driver: &mut D, spec: PacketSpec, v6: bool) -> Option { + if v6 { + let chain = ChainBase::new().eth(|_| {}).ipv6(|_| {}); + if spec.behind_extension { + let chain = chain.hop_by_hop(|_| {}); + match spec.proto { + Proto::Tcp => chain.tcp(|_| {}).generate(driver), + Proto::Udp => chain.udp(|_| {}).generate(driver), + Proto::Icmp => chain.icmp6(|_| {}).generate(driver), + } + } else { + match spec.proto { + Proto::Tcp => chain.tcp(|_| {}).generate(driver), + Proto::Udp => chain.udp(|_| {}).generate(driver), + Proto::Icmp => chain.icmp6(|_| {}).generate(driver), + } + } + } else { + let chain = ChainBase::new().eth(|_| {}).ipv4(|_| {}); + if spec.behind_extension { + let chain = chain.ipv4_auth(|_| {}); + match spec.proto { + Proto::Tcp => chain.tcp(|_| {}).generate(driver), + Proto::Udp => chain.udp(|_| {}).generate(driver), + Proto::Icmp => chain.icmp4(|_| {}).generate(driver), + } + } else { + match spec.proto { + Proto::Tcp => chain.tcp(|_| {}).generate(driver), + Proto::Udp => chain.udp(|_| {}).generate(driver), + Proto::Icmp => chain.icmp4(|_| {}).generate(driver), + } + } + } + } + + fn wire( + headers: &Headers, + spec: PacketSpec, + src: IpAddr, + dst: IpAddr, + ) -> Option> { + let mut headers = headers.clone(); + aim(&mut headers, src, dst, spec.host); + let mut buffer = TestBuffer::new(); + headers.deparse(buffer.as_mut()).ok()?; + Packet::new(buffer).ok() + } + + fn aim(headers: &mut Headers, src: IpAddr, dst: IpAddr, host: u8) { + match (src, dst) { + (IpAddr::V4(src), IpAddr::V4(dst)) => { + if let Some(ip) = headers.try_ipv4_mut() { + let mut o = src.octets(); + o[3] = o[3].wrapping_add(host % 8); + if let Ok(src) = UnicastIpv4Addr::new(Ipv4Addr::from(o)) { + ip.set_source(src); + } + ip.set_destination(dst); + } + } + (IpAddr::V6(src), IpAddr::V6(dst)) => { + if let Some(ip) = headers.try_ipv6_mut() { + let mut o = src.octets(); + o[15] = o[15].wrapping_add(host % 8); + if let Ok(src) = UnicastIpv6Addr::new(Ipv6Addr::from(o)) { + ip.set_source(src); + } + ip.set_destination(dst); + } + } + _ => {} + } + } + + fn peer(family_of: IpAddr) -> IpAddr { + match family_of { + IpAddr::V4(_) => IpAddr::V4( + "3.3.3.1" + .parse::() + .unwrap_or_else(|_| unreachable!()), + ), + IpAddr::V6(_) => IpAddr::V6( + "2001:db8:ffff::1" + .parse::() + .unwrap_or_else(|_| unreachable!()), + ), + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("acl::Batch", &Batch); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn the_acl_verdict_follows_the_protocol_the_packet_carries() { + static DENIED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static PERMITTED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static BEHIND_EXT: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static PERMITTED_OUT: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static DENIED_BY_ACL: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Batch) + .for_each(|(exposes, default_allow, rule_proto, packets)| { + let default = if *default_allow { + AclAction::Allow + } else { + AclAction::Deny + }; + let rule = rule_proto.as_match(); + let Some(mut fabric) = + Fabric::build_with_acl(exposes, Some(&peering_acl(default, rule))) + else { + return; + }; + let Some(private) = exposes + .iter() + .flat_map(|e| e.ips.iter().map(PrefixWithOptionalPorts::prefix)) + .next() + .map(|p: Prefix| p.as_address()) + else { + return; + }; + let v6 = private.is_ipv6(); + let dst = peer(private); + + for (spec, headers) in packets { + let Some(mut packet) = wire(headers, *spec, private, dst) else { + continue; + }; + arrive(&mut packet, local()); + let out = fabric.send(packet); + + let permitted = rule_matches(rule, carried(spec.proto, v6)) != *default_allow; + if spec.behind_extension { + BEHIND_EXT.fetch_add(1, Ordering::Relaxed); + } + + let seen = verdict(&out); + let acl_dropped = seen == Verdict::Dropped(DoneReason::AclDropped); + let forwarded = matches!(seen, Verdict::Forwarded { .. }); + if permitted { + assert!( + !acl_dropped, + "the acl dropped a {:?} packet it permits (rule={rule:?} \ + default={default:?} behind_extension={})", + spec.proto, spec.behind_extension + ); + PERMITTED.fetch_add(1, Ordering::Relaxed); + if forwarded { + PERMITTED_OUT.fetch_add(1, Ordering::Relaxed); + } + } else { + assert!( + !forwarded, + "a {:?} packet the acl denies was forwarded (rule={rule:?} \ + default={default:?} behind_extension={})", + spec.proto, spec.behind_extension + ); + DENIED.fetch_add(1, Ordering::Relaxed); + if acl_dropped { + DENIED_BY_ACL.fetch_add(1, Ordering::Relaxed); + } + } + } + }); + + let (permitted, permitted_out, denied, denied_by_acl, behind) = ( + PERMITTED.load(Ordering::Relaxed), + PERMITTED_OUT.load(Ordering::Relaxed), + DENIED.load(Ordering::Relaxed), + DENIED_BY_ACL.load(Ordering::Relaxed), + BEHIND_EXT.load(Ordering::Relaxed), + ); + eprintln!( + "permitted={permitted} (forwarded {permitted_out}) denied={denied} \ + (by the acl {denied_by_acl}) behind-extension={behind}" + ); + super::assert_covered(permitted > 0, "no packet was ever permitted"); + super::assert_covered(denied > 0, "no packet was ever denied"); + super::assert_covered( + permitted_out > 0, + "no permitted packet was ever forwarded: the permit direction is vacuous", + ); + super::assert_covered( + denied_by_acl > 0, + "no denial ever came from the acl: the deny direction is being satisfied by stages \ + ahead of it, and would hold with the acl removed", + ); + super::assert_covered( + behind > 0, + "no packet was ever sent behind an extension header, which is the shape this exists for", + ); + } +} + +#[cfg(test)] +mod port_forward { + use super::round_trip::udp; + use super::routed::{inside, tunnelled_from}; + use super::*; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::vpcpeering::VpcExpose; + use lpm::prefix::{L4Protocol, PortRange, Prefix, PrefixWithOptionalPorts}; + use net::headers::TryVxlan; + + const INTERNAL_NET: &str = "10.0.5.0/28"; + const EXTERNAL_NET: &str = "172.16.5.0/28"; + const HOSTS: u8 = 16; + const INTERNAL_PORT: u16 = 8000; + const EXTERNAL_PORT: u16 = 9000; + const PORTS: u16 = 8; + + fn expose() -> VpcExpose { + VpcExpose::empty() + .make_port_forwarding(None, Some(L4Protocol::Udp)) + .unwrap_or_else(|_| unreachable!("udp port forwarding is a valid flavour")) + .ip(PrefixWithOptionalPorts::new( + INTERNAL_NET + .parse::() + .unwrap_or_else(|_| unreachable!()), + Some( + PortRange::new(INTERNAL_PORT, INTERNAL_PORT + PORTS - 1) + .unwrap_or_else(|_| unreachable!()), + ), + )) + .as_range(PrefixWithOptionalPorts::new( + EXTERNAL_NET + .parse::() + .unwrap_or_else(|_| unreachable!()), + Some( + PortRange::new(EXTERNAL_PORT, EXTERNAL_PORT + PORTS - 1) + .unwrap_or_else(|_| unreachable!()), + ), + )) + .unwrap_or_else(|_| unreachable!("the two ranges are the same size")) + } + + fn outside() -> IpAddr { + "3.3.3.7".parse().unwrap_or_else(|_| unreachable!()) + } + + #[derive(Debug, Clone, Copy)] + struct Reach { + host: u8, + port: u16, + src_port: u16, + past_the_range: bool, + } + + struct Reaches; + + const REACHES_PER_FABRIC: usize = 8; + + impl bolero::ValueGenerator for Reaches { + type Output = Vec; + + fn generate(&self, driver: &mut D) -> Option> { + (0..REACHES_PER_FABRIC) + .map(|_| { + let choice = driver.produce::()?; + Some(Reach { + host: driver.produce::()? % HOSTS, + port: driver.produce::()? % PORTS, + src_port: driver.produce::()?.max(1), + past_the_range: choice % 4 == 3, + }) + }) + .collect() + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("port_forward::Reaches", &Reaches); + } + + fn answers( + fabric: &mut Fabric, + host: IpAddr, + reach: Reach, + external: IpAddr, + dport: u16, + ) -> bool { + let Some(reply) = udp(host, outside(), INTERNAL_PORT + reach.port, reach.src_port) else { + return false; + }; + let back = fabric.send(tunnelled_from(vni(LOCAL_VNI), &reply)); + assert!( + matches!(verdict(&back), Verdict::Delivered { .. }), + "the service's reply did not get out: {:?}", + verdict(&back) + ); + assert_eq!( + back.try_vxlan().map(net::vxlan::Vxlan::vni), + Some(vni(REMOTE_VNI)), + "the reply went back into the wrong vpc" + ); + let answered = inside(&back).expect("a delivered reply was not tunnelled"); + assert_eq!( + answered.ip_source(), + Some(external), + "the reply was sourced from the internal address, which the outside host never \ + addressed" + ); + assert_eq!( + answered.transport_src_port().map(std::num::NonZero::get), + Some(dport), + "the reply came from the right address on the wrong port" + ); + assert_eq!( + answered.ip_destination(), + Some(outside()), + "the reply did not reach the host that made the request" + ); + true + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_forwarded_port_reaches_the_host_behind_it() { + static FORWARDED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static ANSWERED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static REFUSED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Reaches) + .for_each(|reaches| { + let Some(mut fabric) = Fabric::routed(&[expose()], None) else { + unreachable!("the port-forwarding fixture does not configure") + }; + + for reach in reaches { + let external: IpAddr = format!("172.16.5.{}", reach.host) + .parse() + .unwrap_or_else(|_| unreachable!()); + let dport = if reach.past_the_range { + EXTERNAL_PORT + PORTS + (reach.port % PORTS) + } else { + EXTERNAL_PORT + reach.port + }; + + let Some(inbound) = udp(outside(), external, reach.src_port, dport) else { + continue; + }; + let out = fabric.send(tunnelled_from(vni(REMOTE_VNI), &inbound)); + + if reach.past_the_range { + assert!( + !matches!(verdict(&out), Verdict::Delivered { .. }), + "a packet to {external}:{dport}, past the declared range, was \ + forwarded anyway" + ); + REFUSED.fetch_add(1, Ordering::Relaxed); + continue; + } + + assert!( + matches!(verdict(&out), Verdict::Delivered { .. }), + "a packet to the declared {external}:{dport} was not forwarded: {:?}", + verdict(&out) + ); + let arrived = inside(&out).expect("a forwarded packet was not tunnelled"); + let expected_host: IpAddr = format!("10.0.5.{}", reach.host) + .parse() + .unwrap_or_else(|_| unreachable!()); + assert_eq!( + arrived.ip_destination(), + Some(expected_host), + "{external}:{dport} reached the wrong host" + ); + assert_eq!( + arrived.transport_dst_port().map(std::num::NonZero::get), + Some(INTERNAL_PORT + reach.port), + "{external}:{dport} reached the right host on the wrong port" + ); + FORWARDED.fetch_add(1, Ordering::Relaxed); + + if answers(&mut fabric, expected_host, *reach, external, dport) { + ANSWERED.fetch_add(1, Ordering::Relaxed); + } + } + }); + + let (forwarded, answered, refused) = ( + FORWARDED.load(Ordering::Relaxed), + ANSWERED.load(Ordering::Relaxed), + REFUSED.load(Ordering::Relaxed), + ); + eprintln!("forwarded={forwarded} answered={answered} refused={refused}"); + super::assert_covered(forwarded > 0, "no packet was ever forwarded to the service"); + super::assert_covered(answered > 0, "the service never answered"); + super::assert_covered( + refused > 0, + "no packet was ever aimed past the declared range, so the negative half is vacuous", + ); + } +} + +#[cfg(test)] +mod interleaved { + use super::routed::{Blast, Conversation, Path, exposes}; + use super::*; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use std::ops::Bound::Included; + + const LOADS: usize = 6; + const POLLS: usize = 10; + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + enum Kind { + Conversation, + Blast, + } + + #[derive(Debug, Clone, Copy)] + struct Sender { + kind: Kind, + host: u8, + sport: u16, + dport: u16, + count: u8, + } + + struct Interleaving; + + impl bolero::ValueGenerator for Interleaving { + type Output = (Vec, Vec); + + fn generate(&self, driver: &mut D) -> Option { + let senders = (0..LOADS) + .map(|_| { + Some(Sender { + kind: if driver.produce::()? { + Kind::Conversation + } else { + Kind::Blast + }, + host: driver.produce::()?, + sport: driver.produce::()?.max(1), + dport: driver.produce::()?.max(1), + count: driver.gen_u8(Included(&2), Included(&5))?, + }) + }) + .collect::>>()?; + + let schedule = (0..POLLS) + .map(|_| { + let picks = driver.gen_u8(Included(&1), Included(&3))?; + (0..picks) + .map(|_| { + Some(Pick { + load: driver.produce::()?, + take: driver.gen_u8(Included(&1), Included(&3))?, + }) + }) + .collect::>>() + }) + .collect::>>()?; + + Some((senders, schedule)) + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("interleaved::Interleaving", &Interleaving); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn interleaved_traffic_is_each_satisfied() { + static CHECKED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static ABANDONED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static MIXED_LOADS: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static MIXED_KINDS: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Interleaving) + .for_each(|(senders, schedule)| { + let Some(mut fabric) = Fabric::routed(&exposes(), None) else { + return; + }; + + let dst: IpAddr = "3.3.3.1".parse().unwrap_or_else(|_| unreachable!()); + let mut kinds = Vec::new(); + let mut loads: Vec> = Vec::new(); + for (i, sender) in senders.iter().enumerate() { + let Ok(src) = format!("1.1.{i}.{}", sender.host).parse::() else { + continue; + }; + kinds.push(sender.kind); + loads.push(match sender.kind { + Kind::Conversation => Box::new(Conversation::new( + Path::fixture(), + src, + dst, + sender.sport, + sender.dport, + )), + Kind::Blast => Box::new(Blast::new( + Path::fixture(), + src, + dst, + sender.sport, + sender.dport, + sender.count, + )) as Box, + }); + } + + for burst in run_schedule(&mut fabric, &mut loads, schedule) { + let mut loads_in: Vec = burst.clone(); + loads_in.sort_unstable(); + loads_in.dedup(); + if loads_in.len() > 1 { + MIXED_LOADS.fetch_add(1, Ordering::Relaxed); + } + let mut kinds_in: Vec = burst.iter().map(|i| kinds[*i]).collect(); + kinds_in.sort_unstable_by_key(|k| format!("{k:?}")); + kinds_in.dedup(); + if kinds_in.len() > 1 { + MIXED_KINDS.fetch_add(1, Ordering::Relaxed); + } + } + + for load in &loads { + if load.checked() { + CHECKED.fetch_add(1, Ordering::Relaxed); + } else { + ABANDONED.fetch_add(1, Ordering::Relaxed); + } + } + }); + + let (checked, abandoned, mixed_loads, mixed_kinds) = ( + CHECKED.load(Ordering::Relaxed), + ABANDONED.load(Ordering::Relaxed), + MIXED_LOADS.load(Ordering::Relaxed), + MIXED_KINDS.load(Ordering::Relaxed), + ); + eprintln!( + "checked={checked} abandoned={abandoned} mixed-loads={mixed_loads} \ + mixed-kinds={mixed_kinds}" + ); + super::assert_covered(checked > 0, "no sender ever completed its business"); + super::assert_covered( + mixed_loads > 0, + "no burst ever carried more than one sender's traffic, so nothing was interleaved", + ); + super::assert_covered( + mixed_kinds > 0, + "no burst ever mixed a conversation with a blast, so the two shapes never met", + ); + } +} + +#[cfg(test)] +mod offers { + use super::derive::{Vary, loads_for}; + use super::*; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::vpcpeering::VpcExpose; + use config::external::overlay::vpcpeering::contract::overlay_with_exposes; + use lpm::prefix::{L4Protocol, PortRange, Prefix, PrefixWithOptionalPorts}; + use std::ops::Bound::Included; + + const SENDERS: usize = 6; + const POLLS: usize = 10; + + fn overlay() -> config::external::overlay::ValidatedOverlay { + let masquerade = VpcExpose::empty() + .make_masquerade(None) + .unwrap_or_else(|_| unreachable!()) + .ip("1.1.0.0/16" + .parse::() + .unwrap_or_else(|_| unreachable!()) + .into()) + .as_range( + "2.2.0.0/16" + .parse::() + .unwrap_or_else(|_| unreachable!()) + .into(), + ) + .unwrap_or_else(|_| unreachable!()); + + let forwarded = VpcExpose::empty() + .make_port_forwarding(None, Some(L4Protocol::Udp)) + .unwrap_or_else(|_| unreachable!()) + .ip(PrefixWithOptionalPorts::new( + "10.0.5.0/28" + .parse::() + .unwrap_or_else(|_| unreachable!()), + Some(PortRange::new(8000, 8007).unwrap_or_else(|_| unreachable!())), + )) + .as_range(PrefixWithOptionalPorts::new( + "172.16.5.0/28" + .parse::() + .unwrap_or_else(|_| unreachable!()), + Some(PortRange::new(9000, 9007).unwrap_or_else(|_| unreachable!())), + )) + .unwrap_or_else(|_| unreachable!()); + + overlay_with_exposes(vec![masquerade, forwarded]) + .unwrap_or_else(|e| unreachable!("the fixture does not assemble: {e}")) + .validate() + .unwrap_or_else(|e| unreachable!("the fixture does not validate: {e}")) + } + + struct Offered; + + impl bolero::ValueGenerator for Offered { + type Output = (Vec, Vec); + + fn generate(&self, driver: &mut D) -> Option { + let vary = (0..SENDERS) + .map(|_| { + Some(Vary { + host: driver.produce::()?, + port: driver.produce::()?, + sport: driver.produce::()?.max(1), + dport: driver.produce::()?.max(1), + burst: driver.gen_u8(Included(&2), Included(&5))?, + blast: driver.produce::()?, + }) + }) + .collect::>>()?; + + let schedule = (0..POLLS) + .map(|_| { + let picks = driver.gen_u8(Included(&1), Included(&3))?; + (0..picks) + .map(|_| { + Some(Pick { + load: driver.produce::()?, + take: driver.gen_u8(Included(&1), Included(&3))?, + }) + }) + .collect::>>() + }) + .collect::>>()?; + + Some((vary, schedule)) + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("offers::Offered", &Offered); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_configuration_carries_everything_it_offers() { + static CHECKED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static ABANDONED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static DERIVED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static MIXED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static INBOUND: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static OUTBOUND: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + let overlay = overlay(); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Offered) + .for_each(|(vary, schedule)| { + let mut fabric = Fabric::routed_over_validated( + &overlay, + topology(&[vni(LOCAL_VNI), vni(REMOTE_VNI)]), + ); + + let mut loads = loads_for(&overlay, vary); + DERIVED.fetch_add(loads.len() as u64, Ordering::Relaxed); + for load in &loads { + if load.describe().starts_with("[inbound") { + INBOUND.fetch_add(1, Ordering::Relaxed); + } else { + OUTBOUND.fetch_add(1, Ordering::Relaxed); + } + } + + for burst in run_schedule(&mut fabric, &mut loads, schedule) { + let mut seen = burst.clone(); + seen.sort_unstable(); + seen.dedup(); + if seen.len() > 1 { + MIXED.fetch_add(1, Ordering::Relaxed); + } + } + + for load in &loads { + if load.checked() { + CHECKED.fetch_add(1, Ordering::Relaxed); + } else { + ABANDONED.fetch_add(1, Ordering::Relaxed); + } + } + }); + + let (checked, abandoned, derived, mixed) = ( + CHECKED.load(Ordering::Relaxed), + ABANDONED.load(Ordering::Relaxed), + DERIVED.load(Ordering::Relaxed), + MIXED.load(Ordering::Relaxed), + ); + let (inbound, outbound) = ( + INBOUND.load(Ordering::Relaxed), + OUTBOUND.load(Ordering::Relaxed), + ); + eprintln!( + "checked={checked} abandoned={abandoned} derived={derived} \ + (inbound {inbound}, outbound {outbound}) mixed-bursts={mixed}" + ); + super::assert_covered(derived > 0, "the configuration implied no traffic at all"); + super::assert_covered( + inbound > 0, + "the derivation produced no inbound traffic, so the port-forwarding expose this \ + fixture carries was skipped rather than tested", + ); + super::assert_covered( + outbound > 0, + "the derivation produced no outbound traffic, so the masquerade expose was skipped", + ); + super::assert_covered(checked > 0, "no derived sender ever completed its business"); + super::assert_covered( + mixed > 0, + "no burst ever carried more than one sender's traffic, so nothing was interleaved", + ); + } +} + +#[cfg(test)] +mod generated { + use super::derive::{Vary, loads_for}; + use super::*; + use bolero::ValueGenerator; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::algebra::{Op, Sequence}; + use std::ops::Bound::Included; + + const SENDERS: usize = 6; + const POLLS: usize = 8; + + struct Generated; + + impl ValueGenerator for Generated { + type Output = (Vec, Vec, Vec); + + fn generate(&self, driver: &mut D) -> Option { + let ops = Sequence::default().generate(driver)?; + + let vary = (0..SENDERS) + .map(|_| { + Some(Vary { + host: driver.produce::()?, + port: driver.produce::()?, + sport: driver.produce::()?.max(1), + dport: driver.produce::()?.max(1), + burst: driver.gen_u8(Included(&2), Included(&5))?, + blast: driver.produce::()?, + }) + }) + .collect::>>()?; + + let schedule = (0..POLLS) + .map(|_| { + let picks = driver.gen_u8(Included(&1), Included(&3))?; + (0..picks) + .map(|_| { + Some(Pick { + load: driver.produce::()?, + take: driver.gen_u8(Included(&1), Included(&3))?, + }) + }) + .collect::>>() + }) + .collect::>>()?; + + Some((ops, vary, schedule)) + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("generated::Generated", &Generated); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_generated_configuration_carries_its_own_traffic() { + static CHECKED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static DERIVED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static MIXED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static PEERED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static MULTI: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Generated) + .for_each(|(ops, vary, schedule)| { + let draft = Sequence::fold(ops); + let overlay = draft + .overlay() + .unwrap_or_else(|e| panic!("{ops:?} does not assemble: {e}")); + let validated = overlay + .validate() + .unwrap_or_else(|e| panic!("{ops:?} does not validate: {e}")); + + let vnis: Vec = validated + .vpc_table() + .values() + .map(config::external::overlay::vpc::ValidatedVpc::vni) + .collect(); + if vnis.is_empty() { + return; + } + if validated.vpc_table().peerings().next().is_some() { + PEERED.fetch_add(1, Ordering::Relaxed); + } + if vnis.len() > 2 { + MULTI.fetch_add(1, Ordering::Relaxed); + } + + let mut fabric = Fabric::routed_over_validated(&validated, topology(&vnis)); + + let mut loads = loads_for(&validated, vary); + DERIVED.fetch_add(loads.len() as u64, Ordering::Relaxed); + + for burst in run_schedule(&mut fabric, &mut loads, schedule) { + let mut seen = burst.clone(); + seen.sort_unstable(); + seen.dedup(); + if seen.len() > 1 { + MIXED.fetch_add(1, Ordering::Relaxed); + } + } + + for load in &loads { + assert!( + load.checked(), + "a load derived from the configuration did not complete: {}", + load.describe() + ); + CHECKED.fetch_add(1, Ordering::Relaxed); + } + }); + + let (checked, derived, mixed) = ( + CHECKED.load(Ordering::Relaxed), + DERIVED.load(Ordering::Relaxed), + MIXED.load(Ordering::Relaxed), + ); + let (peered, multi) = ( + PEERED.load(Ordering::Relaxed), + MULTI.load(Ordering::Relaxed), + ); + eprintln!( + "checked={checked} derived={derived} peered-configs={peered} \ + configs-past-two-vpcs={multi} mixed-bursts={mixed}" + ); + super::assert_covered(peered > 0, "no generated configuration ever had a peering"); + super::assert_covered( + multi > 0, + "no generated configuration ever had more than two vpcs, so this reached nothing the \ + two-vpc fixtures do not", + ); + super::assert_covered( + derived > 0, + "no generated configuration ever implied any traffic", + ); + super::assert_covered(checked > 0, "no derived sender ever completed its business"); + super::assert_covered( + mixed > 0, + "no burst ever carried more than one sender's traffic, so nothing was interleaved", + ); + } +} + +#[cfg(test)] +mod burst { + use super::round_trip::udp; + use super::routed::{exposes, inside, tunnelled}; + use super::*; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use net::headers::TryVxlan; + + const BURST: usize = 8; + + #[derive(Debug, Clone, Copy)] + struct Member { + host: u8, + dport: u16, + } + + struct Burst; + + impl bolero::ValueGenerator for Burst { + type Output = Vec; + + fn generate(&self, driver: &mut D) -> Option> { + (0..BURST) + .map(|_| { + Some(Member { + host: driver.produce()?, + dport: driver.produce::()?.max(1), + }) + }) + .collect() + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("burst::Burst", &Burst); + } + + #[derive(Debug, PartialEq, Eq)] + struct Treatment { + verdict: Verdict, + vni: Option, + inner_src: Option, + inner_dst: Option, + inner_sport: Option, + } + + fn treatment(packet: &Packet) -> Treatment { + let carried = inside(packet); + Treatment { + verdict: verdict(packet), + vni: packet.try_vxlan().map(|v| v.vni().as_u32()), + inner_src: carried.as_ref().and_then(Packet::ip_source), + inner_dst: carried.as_ref().and_then(Packet::ip_destination), + inner_sport: carried + .as_ref() + .and_then(Packet::transport_src_port) + .map(std::num::NonZero::get), + } + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_burst_of_one_flow_allocates_once() { + static CHECKED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Burst) + .for_each(|members| { + let m = members[0]; + let src: IpAddr = format!("1.1.0.{}", m.host) + .parse() + .unwrap_or_else(|_| unreachable!()); + let dst: IpAddr = "3.3.3.1".parse().unwrap_or_else(|_| unreachable!()); + let packet = || udp(src, dst, 4000, m.dport).map(|p| tunnelled(&p)); + + let (Some(mut alone), Some(mut burst)) = ( + Fabric::routed(&exposes(), None), + Fabric::routed(&exposes(), None), + ) else { + return; + }; + let Some(one) = packet() else { return }; + let single = treatment(&alone.send(one)); + if !matches!(single.verdict, Verdict::Delivered { .. }) { + return; + } + let cost_of_one = alone.flows(); + + let Some(together) = (0..BURST).map(|_| packet()).collect::>>() + else { + return; + }; + let out = burst.send_batch(together); + + for (i, packet) in out.iter().enumerate() { + let t = treatment(packet); + assert_eq!( + t.inner_sport, single.inner_sport, + "packet {i} of a burst of one flow was given a different public port \ + from the same packet sent alone: the burst allocated more than once" + ); + assert_eq!( + t.inner_src, single.inner_src, + "packet {i} of a burst of one flow left under a different public address" + ); + assert_eq!( + t.verdict, single.verdict, + "packet {i} of a burst of one flow reached a different verdict" + ); + } + assert_eq!( + burst.flows(), + cost_of_one, + "a burst of {BURST} packets of one flow cost more flow-table entries than \ + one packet of it did" + ); + CHECKED.fetch_add(1, Ordering::Relaxed); + }); + + let checked = CHECKED.load(Ordering::Relaxed); + eprintln!("single-flow-bursts={checked}"); + super::assert_covered(checked > 0, "no burst of a single flow was ever delivered"); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_burst_is_treated_the_same_as_one_packet_at_a_time() { + static COMPARED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static DELIVERED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Burst) + .for_each(|members| { + let packets = || { + members + .iter() + .enumerate() + .map(|(i, m)| { + let src: IpAddr = format!("1.1.{i}.{}", m.host) + .parse() + .unwrap_or_else(|_| unreachable!()); + let dst: IpAddr = "3.3.3.1".parse().unwrap_or_else(|_| unreachable!()); + udp(src, dst, 1024 + u16::try_from(i).unwrap_or(0), m.dport) + .map(|p| tunnelled(&p)) + }) + .collect::>>() + }; + let (Some(singly), Some(together)) = (packets(), packets()) else { + return; + }; + + let (Some(mut a), Some(mut b)) = ( + Fabric::routed(&exposes(), None), + Fabric::routed(&exposes(), None), + ) else { + return; + }; + + let one_at_a_time: Vec<_> = + singly.into_iter().map(|p| treatment(&a.send(p))).collect(); + let in_a_burst: Vec<_> = b.send_batch(together).iter().map(treatment).collect(); + + assert_eq!( + one_at_a_time.len(), + in_a_burst.len(), + "a burst did not return as many packets as it was given" + ); + for (i, (alone, batched)) in one_at_a_time.iter().zip(in_a_burst.iter()).enumerate() + { + assert_eq!( + alone, batched, + "packet {i} of the burst was treated differently from the same packet \ + sent on its own" + ); + COMPARED.fetch_add(1, Ordering::Relaxed); + if matches!(alone.verdict, Verdict::Delivered { .. }) { + DELIVERED.fetch_add(1, Ordering::Relaxed); + } + } + }); + + let (compared, delivered) = ( + COMPARED.load(Ordering::Relaxed), + DELIVERED.load(Ordering::Relaxed), + ); + eprintln!("compared={compared} delivered={delivered}"); + super::assert_covered(compared > 0, "no burst was ever compared"); + super::assert_covered( + delivered > 0, + "no packet of any burst ever reached the wire, so the comparison is between drops", + ); + } +} + +#[cfg(test)] +mod destination { + use super::round_trip::udp; + use super::routed::{inside, tunnelled}; + use super::*; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use config::external::overlay::vpcpeering::contract::{overlay_with_peers, peer_vni}; + use lpm::prefix::Prefix; + use net::headers::TryVxlan; + + const PEERS: u8 = 3; + + fn local_prefix() -> Prefix { + "1.1.0.0/16" + .parse() + .unwrap_or_else(|_| unreachable!("a well-formed prefix")) + } + + #[derive(Debug, Clone, Copy)] + struct Aim { + peer: Option, + host: u8, + third: u8, + sport: u16, + dport: u16, + } + + struct Aims; + + const AIMS_PER_FABRIC: usize = 12; + + impl bolero::ValueGenerator for Aims { + type Output = Vec; + + fn generate(&self, driver: &mut D) -> Option> { + (0..AIMS_PER_FABRIC) + .map(|_| { + let choice = driver.produce::()?; + Some(Aim { + peer: (choice % 4 != 3).then_some(choice % PEERS), + host: driver.produce()?, + third: driver.produce()?, + sport: driver.produce::()?.max(1), + dport: driver.produce::()?.max(1), + }) + }) + .collect() + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("destination::Aims", &Aims); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_packet_leaves_for_the_vpc_that_exposes_its_destination() { + static REACHED: LazyLock<[AtomicU64; PEERS as usize]> = + LazyLock::new(|| std::array::from_fn(|_| AtomicU64::new(0))); + static REFUSED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Aims) + .for_each(|aims| { + let vnis: Vec<_> = std::iter::once(vni(LOCAL_VNI)) + .chain((0..PEERS).map(|n| vni(peer_vni(n)))) + .collect(); + let overlay = overlay_with_peers(local_prefix(), PEERS).unwrap_or_else(|e| { + unreachable!("the multi-peer contract does not build: {e}") + }); + let Some(mut fabric) = Fabric::routed_over(&overlay, topology(&vnis)) else { + unreachable!("the multi-peer contract does not validate") + }; + + for aim in aims { + let src: IpAddr = format!("1.1.0.{}", aim.host) + .parse() + .unwrap_or_else(|_| unreachable!()); + let dst: IpAddr = match aim.peer { + Some(n) => format!("10.{}.{}.{}", n + 1, aim.third, aim.host), + None => format!("172.16.{}.{}", aim.third, aim.host), + } + .parse() + .unwrap_or_else(|_| unreachable!()); + + let Some(packet) = udp(src, dst, aim.sport, aim.dport) else { + continue; + }; + let out = fabric.send(tunnelled(&packet)); + let left = matches!(verdict(&out), Verdict::Delivered { .. }); + + if let Some(n) = aim.peer { + assert!( + left, + "a packet to {dst}, which peer {n} exposes, did not leave: {:?}", + verdict(&out) + ); + assert_eq!( + out.try_vxlan().map(net::vxlan::Vxlan::vni), + Some(vni(peer_vni(n))), + "a packet to {dst} left for the wrong vpc" + ); + let carried = inside(&out).expect("a delivered packet was not tunnelled"); + assert_eq!( + carried.ip_destination(), + Some(dst), + "the destination was rewritten on the way out" + ); + REACHED[n as usize].fetch_add(1, Ordering::Relaxed); + } else { + assert!( + !left, + "a packet to {dst}, which no peering covers, was sent to {:?}", + out.try_vxlan().map(net::vxlan::Vxlan::vni) + ); + REFUSED.fetch_add(1, Ordering::Relaxed); + } + } + }); + + let reached: Vec = REACHED.iter().map(|c| c.load(Ordering::Relaxed)).collect(); + let refused = REFUSED.load(Ordering::Relaxed); + eprintln!("reached-per-peer={reached:?} refused={refused}"); + + for (n, count) in reached.iter().enumerate() { + super::assert_covered(*count > 0, &format!("peer {n} was never reached")); + } + super::assert_covered( + refused > 0, + "no packet was ever aimed outside every peering, so the negative half is vacuous", + ); + } +} + +#[cfg(test)] +mod routed { + use super::round_trip::udp; + use super::shapes::{Batch, Shape, aim, wire}; + use super::*; + use super::{Load, drive}; + use concurrency::sync::LazyLock; + use concurrency::sync::atomic::{AtomicU64, Ordering}; + use net::buffer::TestBuffer; + use net::headers::{TryEth, TryHeaders, TryHeadersMut, TryIpv4, TryVxlan}; + use net::ip::dscp::Dscp; + use net::ip::ecn::Ecn; + use net::packet::test_utils::{ + build_test_udp_ipv4_packet, build_test_vxlan_ipv4_packet_carrying_vni, + }; + use net::parse::DeParse; + use net::vlan::Vid; + + #[derive(Debug, Clone, Copy)] + struct Flow { + host: u8, + sport: u16, + dport: u16, + } + + struct Flows; + + const FLOWS_PER_FABRIC: usize = 8; + + impl bolero::ValueGenerator for Flows { + type Output = Vec; + + fn generate(&self, driver: &mut D) -> Option> { + (0..FLOWS_PER_FABRIC) + .map(|_| { + Some(Flow { + host: driver.produce()?, + sport: driver.produce::()?.max(1), + dport: driver.produce::()?.max(1), + }) + }) + .collect() + } + } + + #[test] + fn the_generator_fits_the_input_budget() { + super::assert_within_budget("routed::Flows", &Flows); + } + + pub(super) fn exposes() -> Vec { + vec![ + VpcExpose::empty() + .make_masquerade(None) + .unwrap() + .ip("1.1.0.0/16".parse::().unwrap().into()) + .as_range("2.2.0.0/16".parse::().unwrap().into()) + .unwrap(), + ] + } + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub(crate) struct Path { + pub(crate) from: Vni, + pub(crate) to: Vni, + } + + impl Path { + pub(crate) fn new(from: Vni, to: Vni) -> Self { + Self { from, to } + } + + pub(crate) fn fixture() -> Self { + Self::new(vni(LOCAL_VNI), vni(REMOTE_VNI)) + } + + fn reversed(self) -> Self { + Self::new(self.to, self.from) + } + } + + pub(super) fn tunnelled_from(from: Vni, inner: &Packet) -> Packet { + let bytes = inner + .clone() + .serialize() + .expect("the inner frame serializes"); + let mut packet = build_test_vxlan_ipv4_packet_carrying_vni( + from, + Dscp::default(), + Ecn::default(), + bytes.as_ref(), + ) + .expect("a well-formed tunnelled frame"); + packet + .set_eth_destination(GATEWAY_MAC) + .expect("the frame has an ethernet header"); + packet.meta_mut().iif = Some(uplink()); + packet.meta_mut().set_keep(true); + packet + } + + pub(super) fn tunnelled(inner: &Packet) -> Packet { + tunnelled_from(vni(LOCAL_VNI), inner) + } + + pub(super) fn inside(delivered: &Packet) -> Option> { + let mut copy = delivered.clone(); + matches!(copy.vxlan_decap(), Some(Ok(_))).then_some(copy) + } + + fn inner() -> Packet { + build_test_udp_ipv4_packet("1.1.0.1", "3.3.3.1", 1234, 80) + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_tunnelled_frame_is_decapsulated_translated_and_sent_back_out_tunnelled() { + let mut fabric = Fabric::routed(&exposes(), None).expect("a valid configuration"); + + let out = fabric.send(tunnelled(&inner())); + + match verdict(&out) { + Verdict::Delivered { oif, src, dst } => { + assert_eq!(oif, Some(uplink()), "delivered over the wrong interface"); + assert_eq!( + src, + Some(LOCAL_VTEP.parse().unwrap()), + "the outer source is not this gateway's vtep: it did not get re-encapsulated" + ); + assert_eq!(dst, Some(PEER_VTEP.parse().unwrap())); + assert_eq!( + out.try_vxlan().map(net::vxlan::Vxlan::vni), + Some(vni(REMOTE_VNI)), + "encapsulated towards the wrong vpc" + ); + assert_eq!( + out.try_eth().map(|eth| eth.destination().inner()), + Some(PEER_MAC), + "the frame was not addressed to the resolved next hop" + ); + } + other => panic!("the frame did not leave the gateway: {other:?}"), + } + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_vlan_tag_is_refused_at_decapsulation() { + let tables = topology(&[vni(LOCAL_VNI), vni(REMOTE_VNI)]); + let mut pipeline = DynPipeline::new() + .add_stage(Ingress::new("ingress", tables.interfaces())) + .add_stage(IpForwarder::new("ip-forward-1", tables.fibs())); + + let plain = one(&mut pipeline, tunnelled(&inner())); + assert_eq!( + verdict(&plain), + Verdict::Forwarded { + dst_vpcd: None, + src: Some("1.1.0.1".parse().unwrap()), + dst: Some("3.3.3.1".parse().unwrap()), + }, + "an untagged frame did not survive decapsulation: the fixture is refusing everything" + ); + + let tagged = one(&mut pipeline, tunnelled(&tagged_inner())); + assert_eq!( + verdict(&tagged), + Verdict::Dropped(DoneReason::Unhandled), + "a tagged frame survived decapsulation" + ); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_vlan_tag_inside_the_tunnel_never_leaves_the_gateway() { + let mut fabric = Fabric::routed(&exposes(), None).expect("a valid configuration"); + + let out = fabric.send(tunnelled(&tagged_inner())); + + assert!( + !matches!(verdict(&out), Verdict::Delivered { .. }), + "a tagged frame was sent out onto the wire: {:?}", + verdict(&out) + ); + } + + fn tagged_inner() -> Packet { + let mut inner = inner(); + inner + .headers_mut() + .push_vlan(Vid::new(42).expect("a valid vlan id")) + .expect("the inner frame has an ethernet header"); + + let bytes = inner.serialize().expect("a tagged frame serializes"); + let reparsed = Packet::new(TestBuffer::from_raw_data(bytes.as_ref())) + .expect("a tagged frame parses back"); + assert!( + !reparsed.headers().vlan().is_empty(), + "the fixture did not actually tag the frame" + ); + reparsed + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_tagged_shape_never_reaches_the_wire() { + static TAGGED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Batch) + .for_each(|(exposes, stacks)| { + let Some(mut fabric) = Fabric::routed(exposes, None) else { + return; + }; + let private = exposes.first().and_then(|e| { + e.ips + .first() + .map(lpm::prefix::PrefixWithOptionalPorts::prefix) + }); + + for (shape, headers) in stacks { + let mut headers = headers.clone(); + aim(&mut headers, private); + let Some(frame) = wire(&headers) else { + continue; + }; + let tagged = *shape == Shape::VlanV4Tcp; + if tagged { + TAGGED.fetch_add(1, Ordering::Relaxed); + } + + let out = fabric.send(tunnelled(&frame)); + assert!( + !(tagged && matches!(verdict(&out), Verdict::Delivered { .. })), + "a tagged frame was sent out onto the wire" + ); + } + }); + + let tagged = TAGGED.load(Ordering::Relaxed); + eprintln!("tagged={tagged}"); + + let mut control = Fabric::routed(&exposes(), None).expect("a valid configuration"); + assert!( + matches!( + verdict(&control.send(tunnelled(&inner()))), + Verdict::Delivered { .. } + ), + "the untagged control did not reach the wire, so no delivery was observable here" + ); + super::assert_covered(tagged > 0, "no tagged shape was ever generated"); + } + + #[tokio::test] + #[dpdk::with_eal] + async fn a_tunnelled_flow_comes_back_through_the_tunnel() { + static ROUND_TRIPPED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + static ABANDONED: LazyLock = LazyLock::new(|| AtomicU64::new(0)); + + bolero::check!() + .with_max_len(MAX_INPUT_LEN) + .with_generator(Flows) + .for_each(|flows| { + let Some(mut fabric) = Fabric::routed(&exposes(), None) else { + return; + }; + + for flow in flows { + let src: IpAddr = format!("1.1.0.{}", flow.host) + .parse() + .unwrap_or_else(|_| unreachable!()); + let dst: IpAddr = "3.3.3.1".parse().unwrap_or_else(|_| unreachable!()); + let mut load = + Conversation::new(Path::fixture(), src, dst, flow.sport, flow.dport); + + drive(&mut fabric, &mut load); + + if load.checked() { + ROUND_TRIPPED.fetch_add(1, Ordering::Relaxed); + } else { + ABANDONED.fetch_add(1, Ordering::Relaxed); + } + } + }); + + let round_tripped = ROUND_TRIPPED.load(Ordering::Relaxed); + eprintln!( + "tunnelled-round-trips={round_tripped} abandoned={}", + ABANDONED.load(Ordering::Relaxed) + ); + super::assert_covered( + round_tripped > 0, + "no flow ever reached the wire, so no reply was ever checked to come back", + ); + } + + pub(super) struct Conversation { + path: Path, + src: IpAddr, + dst: IpAddr, + sport: u16, + dport: u16, + sent: Option>, + state: State, + log: Vec, + } + + #[derive(Debug, PartialEq, Eq)] + enum State { + Opening, + AwaitingRequest, + Replying { public: (IpAddr, u16) }, + AwaitingReply, + Closed, + Abandoned, + } + + impl Conversation { + pub(super) fn new(path: Path, src: IpAddr, dst: IpAddr, sport: u16, dport: u16) -> Self { + Self { + path, + src, + dst, + sport, + dport, + sent: None, + state: State::Opening, + log: Vec::new(), + } + } + + fn note(&mut self, what: &str) { + self.log.push(what.to_owned()); + } + + fn judge_request(&mut self, got: &Packet) { + if !matches!(verdict(got), Verdict::Delivered { .. }) { + self.note(&format!("request not delivered: {:?}", verdict(got))); + self.state = State::Abandoned; + return; + } + assert_eq!( + got.try_vxlan().map(net::vxlan::Vxlan::vni), + Some(self.path.to), + "the request left tunnelled towards the wrong vpc. {}", + self.describe() + ); + + let carried = inside(got).expect("a delivered request was not tunnelled"); + let sent = self.sent.as_ref().unwrap_or_else(|| unreachable!()); + assert_eq!( + payload_of(&carried), + payload_of(sent), + "the tenant payload did not survive decap, nat and encap. {}", + self.describe() + ); + assert_eq!( + ttl_of(&carried).map(|t| t + 1), + ttl_of(sent), + "the tenant packet was not charged exactly one hop. {}", + self.describe() + ); + + let (Some(public_src), Some(port)) = + (carried.ip_source(), carried.transport_src_port()) + else { + self.note("the delivered request had no source tuple to answer"); + self.state = State::Abandoned; + return; + }; + self.note(&format!("request left as {public_src}:{}", port.get())); + self.state = State::Replying { + public: (public_src, port.get()), + }; + } + + fn judge_reply(&mut self, got: &Packet) { + assert!( + matches!(verdict(got), Verdict::Delivered { .. }), + "the reply of a delivered flow did not reach the wire: {:?}. {}", + verdict(got), + self.describe() + ); + assert_eq!( + got.try_vxlan().map(net::vxlan::Vxlan::vni), + Some(self.path.from), + "the reply went back into the wrong vpc. {}", + self.describe() + ); + + let returned = inside(got).expect("a delivered reply was not tunnelled"); + assert_eq!( + returned.ip_destination(), + Some(self.src), + "the reply did not come back to the host that sent the request. {}", + self.describe() + ); + assert_eq!( + returned.ip_source(), + Some(self.dst), + "the reply's source was rewritten. {}", + self.describe() + ); + assert_eq!( + returned.transport_dst_port().map(std::num::NonZero::get), + Some(self.sport), + "the reply did not get the original source port back. {}", + self.describe() + ); + self.note("reply came back"); + self.state = State::Closed; + } + } + + impl Load for Conversation { + fn next(&mut self) -> Option> { + match self.state { + State::Opening => { + let request = udp(self.src, self.dst, self.sport, self.dport)?; + self.sent = Some(request.clone()); + self.state = State::AwaitingRequest; + Some(tunnelled_from(self.path.from, &request)) + } + State::Replying { public: (ip, port) } => { + let reply = udp(self.dst, ip, self.dport, port)?; + self.state = State::AwaitingReply; + Some(tunnelled_from(self.path.reversed().from, &reply)) + } + State::AwaitingRequest + | State::AwaitingReply + | State::Closed + | State::Abandoned => None, + } + } + + fn observe(&mut self, got: &Packet) { + match self.state { + State::AwaitingRequest => self.judge_request(got), + State::AwaitingReply => self.judge_reply(got), + ref state => panic!( + "a load was given an answer it was not waiting for (state {state:?}). {}", + self.describe() + ), + } + } + + fn finished(&self) -> bool { + matches!(self.state, State::Closed | State::Abandoned) + } + + fn checked(&self) -> bool { + self.state == State::Closed + } + + fn describe(&self) -> String { + format!( + "[conversation {}:{} -> {}:{} | {}]", + self.src, + self.sport, + self.dst, + self.dport, + if self.log.is_empty() { + "nothing yet".to_owned() + } else { + self.log.join("; ") + } + ) + } + } + + pub(super) struct Blast { + src: IpAddr, + dst: IpAddr, + sport: u16, + dport: u16, + path: Path, + to_send: u8, + in_flight: u8, + given: Option<(IpAddr, u16)>, + delivered: u8, + log: Vec, + } + + impl Blast { + pub(super) fn new( + path: Path, + src: IpAddr, + dst: IpAddr, + sport: u16, + dport: u16, + count: u8, + ) -> Self { + Self { + src, + dst, + sport, + dport, + path, + to_send: count.max(2), + in_flight: 0, + given: None, + delivered: 0, + log: Vec::new(), + } + } + } + + impl Load for Blast { + fn next(&mut self) -> Option> { + if self.to_send == 0 { + return None; + } + let packet = udp(self.src, self.dst, self.sport, self.dport)?; + self.to_send -= 1; + self.in_flight += 1; + Some(tunnelled_from(self.path.from, &packet)) + } + + fn observe(&mut self, got: &Packet) { + assert!( + self.in_flight > 0, + "a blast was given an answer it was not waiting for. {}", + self.describe() + ); + self.in_flight -= 1; + + if !matches!(verdict(got), Verdict::Delivered { .. }) { + self.log.push(format!("not delivered: {:?}", verdict(got))); + return; + } + let Some(carried) = inside(got) else { + self.log.push("delivered but not tunnelled".to_owned()); + return; + }; + let (Some(source), Some(port)) = (carried.ip_source(), carried.transport_src_port()) + else { + return; + }; + let now = (source, port.get()); + self.delivered += 1; + match self.given { + None => { + self.log.push(format!("left as {source}:{}", port.get())); + self.given = Some(now); + } + Some(first) => assert_eq!( + now, + first, + "packets of one flow were given different public tuples, so the flow was \ + allocated for more than once. {}", + self.describe() + ), + } + } + + fn finished(&self) -> bool { + self.to_send == 0 && self.in_flight == 0 + } + + fn checked(&self) -> bool { + self.delivered >= 2 + } + + fn describe(&self) -> String { + format!( + "[blast {}:{} -> {}:{} | {} left, {} in flight, {} delivered | {}]", + self.src, + self.sport, + self.dst, + self.dport, + self.to_send, + self.in_flight, + self.delivered, + if self.log.is_empty() { + "nothing yet".to_owned() + } else { + self.log.join("; ") + } + ) + } + } + + pub(super) struct Inbound { + path: Path, + from: IpAddr, + external: IpAddr, + external_port: u16, + internal: IpAddr, + internal_port: u16, + sport: u16, + state: InboundState, + log: Vec, + } + + #[derive(Debug, PartialEq, Eq)] + enum InboundState { + Reaching, + AwaitingArrival, + Answering, + AwaitingAnswer, + Closed, + Abandoned, + } + + impl Inbound { + pub(super) fn new( + path: Path, + from: IpAddr, + external: IpAddr, + external_port: u16, + internal: IpAddr, + internal_port: u16, + sport: u16, + ) -> Self { + Self { + path, + from, + external, + external_port, + internal, + internal_port, + sport, + state: InboundState::Reaching, + log: Vec::new(), + } + } + + fn judge_arrival(&mut self, got: &Packet) { + if !matches!(verdict(got), Verdict::Delivered { .. }) { + self.log.push(format!("not forwarded: {:?}", verdict(got))); + self.state = InboundState::Abandoned; + return; + } + let arrived = inside(got).expect("a forwarded packet was not tunnelled"); + assert_eq!( + arrived.ip_destination(), + Some(self.internal), + "reached the wrong host. {}", + self.describe() + ); + assert_eq!( + arrived.transport_dst_port().map(std::num::NonZero::get), + Some(self.internal_port), + "reached the right host on the wrong port. {}", + self.describe() + ); + self.log.push("arrived inside".to_owned()); + self.state = InboundState::Answering; + } + + fn judge_answer(&mut self, got: &Packet) { + assert!( + matches!(verdict(got), Verdict::Delivered { .. }), + "the service's answer did not get out: {:?}. {}", + verdict(got), + self.describe() + ); + let answered = inside(got).expect("a delivered answer was not tunnelled"); + assert_eq!( + answered.ip_source(), + Some(self.external), + "the answer was sourced from the internal address, which the outside host never \ + addressed. {}", + self.describe() + ); + assert_eq!( + answered.transport_src_port().map(std::num::NonZero::get), + Some(self.external_port), + "the answer came from the right address on the wrong port. {}", + self.describe() + ); + assert_eq!( + answered.ip_destination(), + Some(self.from), + "the answer did not reach the host that made the request. {}", + self.describe() + ); + self.log.push("answered".to_owned()); + self.state = InboundState::Closed; + } + } + + impl Load for Inbound { + fn next(&mut self) -> Option> { + match self.state { + InboundState::Reaching => { + let request = udp(self.from, self.external, self.sport, self.external_port)?; + self.state = InboundState::AwaitingArrival; + Some(tunnelled_from(self.path.to, &request)) + } + InboundState::Answering => { + let answer = udp(self.internal, self.from, self.internal_port, self.sport)?; + self.state = InboundState::AwaitingAnswer; + Some(tunnelled_from(self.path.from, &answer)) + } + _ => None, + } + } + + fn observe(&mut self, got: &Packet) { + match self.state { + InboundState::AwaitingArrival => self.judge_arrival(got), + InboundState::AwaitingAnswer => self.judge_answer(got), + ref state => panic!( + "a load was given an answer it was not waiting for (state {state:?}). {}", + self.describe() + ), + } + } + + fn finished(&self) -> bool { + matches!(self.state, InboundState::Closed | InboundState::Abandoned) + } + + fn checked(&self) -> bool { + self.state == InboundState::Closed + } + + fn describe(&self) -> String { + format!( + "[inbound {}:{} -> {}:{} (expects {}:{}) | {}]", + self.from, + self.sport, + self.external, + self.external_port, + self.internal, + self.internal_port, + if self.log.is_empty() { + "nothing yet".to_owned() + } else { + self.log.join("; ") + } + ) + } + } + + fn payload_of(packet: &Packet) -> Vec { + let bytes = packet + .clone() + .serialize() + .expect("a packet in hand serializes"); + let headers = packet.headers().size().get() as usize; + bytes.as_ref().get(headers..).unwrap_or_default().to_vec() + } + + fn ttl_of(packet: &Packet) -> Option { + packet.try_ipv4().map(net::ipv4::Ipv4::ttl) + } + + fn one( + pipeline: &mut DynPipeline, + packet: Packet, + ) -> Packet { + let mut out: Vec<_> = pipeline.process(std::iter::once(packet)).collect(); + assert_eq!(out.len(), 1, "the pipeline did not return the packet"); + out.pop().unwrap_or_else(|| unreachable!()) + } +} diff --git a/dataplane/src/packet_processor/ipforward.rs b/dataplane/src/packet_processor/ipforward.rs index 483882e416..fd78f045a4 100644 --- a/dataplane/src/packet_processor/ipforward.rs +++ b/dataplane/src/packet_processor/ipforward.rs @@ -154,7 +154,8 @@ impl IpForwarder { debug!("Next fib/vrf is {next_vrf}"); /* At this point decapsulation has already happened and `Packet` refers to - the innner packet. */ + the innner packet. Annotate the incoming vni and the corresponding vrf to + make lookups from */ if !packet.headers().vlan().is_empty() { debug!( diff --git a/dataplane/src/packet_processor/mod.rs b/dataplane/src/packet_processor/mod.rs index 4e9190c5ce..971c4d2b27 100644 --- a/dataplane/src/packet_processor/mod.rs +++ b/dataplane/src/packet_processor/mod.rs @@ -2,6 +2,8 @@ // Copyright Open Network Fabric Authors mod egress; +#[cfg(test)] +mod fuzz; mod ingress; mod ipforward; diff --git a/justfile b/justfile index 5dd7a02f82..e33658365c 100644 --- a/justfile +++ b/justfile @@ -52,6 +52,14 @@ bolero_coverage_test_time_ms := env("BOLERO_COVERAGE_TEST_TIME_MS", "15000") # comma-separated list of cargo features to enable (e.g. "shuttle") features := "" +fuzz_max_input_length := env("FUZZ_MAX_INPUT_LENGTH", "65536") + +fuzz_corpus_root := env("FUZZ_CORPUS_ROOT", justfile_directory() / ".fuzz-corpus") + +fuzz_jobs := env("FUZZ_JOBS", `echo $(( ($(nproc) + 1) / 2 ))`) + +fuzz_len_control := env("FUZZ_LEN_CONTROL", "0") + # whether to include default cargo features for this workspace (set to "false" to disable) default_features := "true" @@ -197,7 +205,13 @@ fuzz target time="60s" *args="": # asan does not need that, and skipping the std rebuild keeps it far quicker. # `sanitize=NONE` drops instrumentation altogether, which buys roughly four times # the executions per second in exchange for only catching what the test asserts. + corpus_dir="{{ fuzz_corpus_root }}/$(printf '%s' '{{ target }}' | tr -c 'A-Za-z0-9_.-' '_')" + mkdir -p "${corpus_dir}" cargo bolero test '{{ target }}' --rustc-bootstrap -T '{{ time }}' \ + --corpus-dir "${corpus_dir}" \ + -l '{{ fuzz_max_input_length }}' \ + -E='-len_control={{ fuzz_len_control }}' \ + -j '{{ fuzz_jobs }}' \ {{ if sanitize != "" { "--sanitizer " + sanitize } else { "" } }} \ {{ if sanitize == "thread" { "--build-std" } else { "" } }} \ {{ _cargo_feature_flags }} {{ args }} diff --git a/net/Cargo.toml b/net/Cargo.toml index 5059c423fd..174bd1ad3e 100644 --- a/net/Cargo.toml +++ b/net/Cargo.toml @@ -12,6 +12,7 @@ netdevsim = [] bolero = ["dep:bolero", "id/bolero"] test_buffer = [] +test_meta = [] builder = [] [target.'cfg(unix)'.dependencies] diff --git a/net/src/packet/meta.rs b/net/src/packet/meta.rs index 826c8fcb21..6df494e299 100644 --- a/net/src/packet/meta.rs +++ b/net/src/packet/meta.rs @@ -151,6 +151,14 @@ pub struct PacketMeta { pub dscp: Option, /* Dscp to preserve for egress traffic */ pub ecn: Option, /* Ecn to preserve for egress traffic */ pub flow_key: Option>, /* the flow key to use for NAT flow creation */ + #[cfg(feature = "test_meta")] + pub test: Option>, +} + +#[cfg(feature = "test_meta")] +#[derive(Debug, Default, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)] +pub struct TestMeta { + pub id: u64, } impl PacketMeta { #[must_use] diff --git a/net/src/packet/test_utils.rs b/net/src/packet/test_utils.rs index deb7d4bb6a..47f309a325 100644 --- a/net/src/packet/test_utils.rs +++ b/net/src/packet/test_utils.rs @@ -536,9 +536,23 @@ pub fn build_test_vxlan_ipv4_packet_carrying( dscp: Dscp, ecn: Ecn, inner_bytes: &[u8], +) -> Result, InvalidPacket> { + build_test_vxlan_ipv4_packet_carrying_vni( + Vni::new_checked(100).unwrap_or_else(|_| unreachable!()), + dscp, + ecn, + inner_bytes, + ) +} + +#[must_use] +pub fn build_test_vxlan_ipv4_packet_carrying_vni( + vni: Vni, + dscp: Dscp, + ecn: Ecn, + inner_bytes: &[u8], ) -> Result, InvalidPacket> { // VXLAN header bytes - let vni = Vni::new_checked(100).unwrap(); let vxlan = Vxlan::new(vni); let vxlan_len = Vxlan::MIN_LENGTH.get() as usize; diff --git a/pipeline/src/pipeline.rs b/pipeline/src/pipeline.rs index c0cdf110b2..2e9896869d 100644 --- a/pipeline/src/pipeline.rs +++ b/pipeline/src/pipeline.rs @@ -245,7 +245,7 @@ mod test { let mut pipeline = DynPipeline::new(); let mut stages = DynStageGenerator::new(); - let num_stages = 999; + let num_stages = 500; for _ in 0..num_stages { pipeline = pipeline.add_stage_dyn(stages.next().unwrap()); diff --git a/routing/src/interfaces/iftablerw.rs b/routing/src/interfaces/iftablerw.rs index c6aabc6c4e..afd912c0c4 100644 --- a/routing/src/interfaces/iftablerw.rs +++ b/routing/src/interfaces/iftablerw.rs @@ -169,9 +169,12 @@ impl IfTableWriter { ) -> Result<(), RouterError> { // FIXME(fredi): this can be significantly simplified let fibkey = self.interface_attach_check(ifindex, vrfid, vrftable)?; + self.attach_interface_to_fib(ifindex, fibkey); + Ok(()) + } + pub(crate) fn attach_interface_to_fib(&mut self, ifindex: InterfaceIndex, fibkey: FibKey) { self.0.append(IfTableChange::Attach((ifindex, fibkey))); self.0.publish(); - Ok(()) } pub fn detach_interface(&mut self, ifindex: InterfaceIndex) { self.0.append(IfTableChange::Detach(ifindex)); diff --git a/routing/src/lib.rs b/routing/src/lib.rs index 82496d120b..3791867c8a 100644 --- a/routing/src/lib.rs +++ b/routing/src/lib.rs @@ -46,12 +46,7 @@ pub use rib::encapsulation::{ pub use rib::vrf::{RouterVrfConfig, VrfId}; #[cfg(any(test, feature = "testing"))] -pub mod testing { - pub use crate::fib::fibobjects::FibGroup; - pub use crate::fib::fibtype::{Fib, FibReader, FibWriter}; - pub use crate::rib::nexthop::{FwAction, NhopKey}; - pub use crate::rib::vrf::RouteOrigin; -} +pub mod testing; pub use bmp::spawn_bmp_server; pub use router::ctl::RouterCtlSender; diff --git a/routing/src/testing.rs b/routing/src/testing.rs new file mode 100644 index 0000000000..bbb637208e --- /dev/null +++ b/routing/src/testing.rs @@ -0,0 +1,157 @@ +// SPDX-License-Identifier: Apache-2.0 +// Copyright Open Network Fabric Authors + +pub use crate::fib::fibobjects::FibGroup; +pub use crate::fib::fibtype::{Fib, FibReader, FibWriter}; +pub use crate::rib::nexthop::{FwAction, NhopKey}; +pub use crate::rib::vrf::RouteOrigin; + +use std::collections::BTreeMap; +use std::net::IpAddr; + +use lpm::prefix::Prefix; +use net::eth::mac::{Mac, SourceMac}; +use net::interface::InterfaceIndex; +use net::vxlan::Vni; + +use crate::atable::adjacency::Adjacency; +use crate::atable::atablerw::{AtableReader, AtableWriter}; +use crate::evpn::Vtep; +use crate::fib::fibtable::{FibTableReader, FibTableWriter}; +use crate::fib::fibtype::FibKey; +use crate::interfaces::iftablerw::{IfTableReader, IfTableWriter}; +use crate::interfaces::interface::{IfDataEthernet, IfState, IfType, RouterInterfaceConfig}; +use crate::rib::vrf::VrfId; + +pub struct RouterTables { + fib_table: FibTableWriter, + fibs: BTreeMap, + interfaces: IfTableWriter, + adjacencies: AtableWriter, + fib_reader: FibTableReader, + if_reader: IfTableReader, + adj_reader: AtableReader, +} + +impl Default for RouterTables { + fn default() -> Self { + Self::new() + } +} + +impl RouterTables { + #[must_use] + pub fn new() -> Self { + let (fib_table, fib_reader) = FibTableWriter::new(); + let (interfaces, if_reader) = IfTableWriter::new(); + let (adjacencies, adj_reader) = AtableWriter::new(); + Self { + fib_table, + fibs: BTreeMap::new(), + interfaces, + adjacencies, + fib_reader, + if_reader, + adj_reader, + } + } + + pub fn vrf(&mut self, vrfid: VrfId, vni: Option) -> &mut Self { + let fib = self.fib_table.add_fib(vrfid, vni); + self.fibs.insert(vrfid, fib); + self + } + + pub fn vtep(&mut self, vrfid: VrfId, vtep: Vtep) -> &mut Self { + let fib = self.fib_mut(vrfid); + fib.set_vtep(vtep); + fib.publish(); + self + } + + pub fn nexthop(&mut self, vrfid: VrfId, key: &NhopKey, group: &FibGroup) -> &mut Self { + self.fib_mut(vrfid).register_fibgroup(key, group, true); + self + } + + pub fn route(&mut self, vrfid: VrfId, prefix: Prefix, keys: Vec) -> &mut Self { + self.fib_mut(vrfid).add_fibroute(prefix, keys, true); + self + } + + pub fn route_via( + &mut self, + vrfid: VrfId, + prefix: Prefix, + key: NhopKey, + group: &FibGroup, + ) -> &mut Self { + self.nexthop(vrfid, &key, group); + self.route(vrfid, prefix, vec![key]) + } + + /// # Panics + /// + /// Panics if `ifindex` is already in use. + pub fn interface(&mut self, ifindex: InterfaceIndex, name: &str, mac: SourceMac) -> &mut Self { + let mut config = RouterInterfaceConfig::new(name, ifindex); + config.set_iftype(IfType::Ethernet(IfDataEthernet { mac })); + config.set_admin_state(IfState::Up); + self.interfaces + .add_interface(config) + .expect("interface index is already in use"); + self.interfaces.set_iface_oper_state(ifindex, IfState::Up); + self + } + + pub fn interface_state( + &mut self, + ifindex: InterfaceIndex, + admin: IfState, + oper: IfState, + ) -> &mut Self { + self.interfaces.set_iface_admin_state(ifindex, admin); + self.interfaces.set_iface_oper_state(ifindex, oper); + self + } + + /// # Panics + /// + /// Panics if there is no fib for `vrfid`; call `vrf` before `attach`. + pub fn attach(&mut self, ifindex: InterfaceIndex, vrfid: VrfId) -> &mut Self { + assert!( + self.fibs.contains_key(&vrfid), + "no fib for vrf {vrfid}: call `vrf` before `attach`" + ); + self.interfaces + .attach_interface_to_fib(ifindex, FibKey::from_vrfid(vrfid)); + self + } + + pub fn adjacency(&mut self, address: IpAddr, ifindex: InterfaceIndex, mac: Mac) -> &mut Self { + self.adjacencies + .add_adjacency(Adjacency::new(address, ifindex, mac), true); + self + } + + #[must_use] + pub fn interfaces(&self) -> IfTableReader { + self.if_reader.clone() + } + + #[must_use] + pub fn fibs(&self) -> FibTableReader { + self.fib_reader.clone() + } + + #[must_use] + pub fn adjacencies(&self) -> AtableReader { + self.adj_reader.clone() + } + + fn fib_mut(&mut self, vrfid: VrfId) -> &mut FibWriter { + self.fibs + .get_mut(&vrfid) + .unwrap_or_else(|| panic!("no fib for vrf {vrfid}: call `vrf` first")) + } +}