use std::{
collections::{HashMap, VecDeque},
pin::Pin,
task::{Context, Poll},
};
use tokio_stream::Stream;
use super::{events::NetworkEvent, reflector::Store, uevent::Uevent};
use crate::{LinkMessage, Result};
const NET_SUBSYSTEM: &str = "net";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NetdevAnnotation {
pub devpath: String,
pub interface: Option<String>,
pub driver: Option<String>,
pub devtype: Option<String>,
pub seqnum: Option<u64>,
}
impl NetdevAnnotation {
fn from_uevent(event: &Uevent) -> Option<(u32, Self)> {
let ifindex = event.env.get("IFINDEX")?.parse().ok()?;
Some((
ifindex,
Self {
devpath: event.devpath.clone(),
interface: event.env.get("INTERFACE").cloned(),
driver: event.driver().map(str::to_string),
devtype: event.devtype().map(str::to_string),
seqnum: event.seqnum(),
},
))
}
}
#[derive(Debug, Clone)]
pub struct NetdevInfo {
link: LinkMessage,
annotation: Option<NetdevAnnotation>,
}
impl NetdevInfo {
pub fn ifindex(&self) -> u32 {
self.link.ifindex()
}
pub fn name(&self) -> Option<&str> {
self.link.name()
}
pub fn link(&self) -> &LinkMessage {
&self.link
}
pub fn annotation(&self) -> Option<&NetdevAnnotation> {
self.annotation.as_ref()
}
pub fn driver(&self) -> Option<&str> {
self.annotation.as_ref()?.driver.as_deref()
}
pub fn devpath(&self) -> Option<&str> {
Some(self.annotation.as_ref()?.devpath.as_str())
}
pub fn is_fully_attributed(&self) -> bool {
let Some(annotation) = &self.annotation else {
return false;
};
match (annotation.interface.as_deref(), self.link.name()) {
(Some(a), Some(b)) => a == b,
_ => true,
}
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum NetdevEvent {
Added(NetdevInfo),
Changed(NetdevInfo),
DriverBound {
ifindex: u32,
annotation: NetdevAnnotation,
},
DriverUnbound {
ifindex: u32,
annotation: NetdevAnnotation,
},
Removed {
ifindex: u32,
name: Option<String>,
},
}
impl NetdevEvent {
pub fn info(&self) -> Option<&NetdevInfo> {
match self {
Self::Added(info) | Self::Changed(info) => Some(info),
_ => None,
}
}
pub fn ifindex(&self) -> u32 {
match self {
Self::Added(info) | Self::Changed(info) => info.ifindex(),
Self::DriverBound { ifindex, .. }
| Self::DriverUnbound { ifindex, .. }
| Self::Removed { ifindex, .. } => *ifindex,
}
}
}
pub struct NetdevLifecycle<L, U> {
links: Option<L>,
uevents: Option<U>,
known: HashMap<u32, LinkMessage>,
annotations: HashMap<u32, NetdevAnnotation>,
pending: VecDeque<NetdevEvent>,
store: Option<Store<u32, NetdevInfo>>,
prefer_uevents: bool,
}
impl<L, U> NetdevLifecycle<L, U>
where
L: Stream<Item = Result<NetworkEvent>> + Unpin,
U: Stream<Item = Result<Uevent>> + Unpin,
{
pub fn new(links: L, uevents: U) -> Self {
Self {
links: Some(links),
uevents: Some(uevents),
known: HashMap::new(),
annotations: HashMap::new(),
pending: VecDeque::new(),
store: None,
prefer_uevents: false,
}
}
pub fn with_store(mut self, store: Store<u32, NetdevInfo>) -> Self {
self.store = Some(store);
self
}
pub fn len(&self) -> usize {
self.known.len()
}
pub fn is_empty(&self) -> bool {
self.known.is_empty()
}
fn info(&self, ifindex: u32) -> Option<NetdevInfo> {
Some(NetdevInfo {
link: self.known.get(&ifindex)?.clone(),
annotation: self.annotations.get(&ifindex).cloned(),
})
}
fn on_link_event(&mut self, event: NetworkEvent) {
match event {
NetworkEvent::NewLink(link) => {
let ifindex = link.ifindex();
let first_sighting = !self.known.contains_key(&ifindex);
if first_sighting {
if let Some(stale) = self.annotations.get(&ifindex)
&& let (Some(annotated), Some(actual)) =
(stale.interface.as_deref(), link.name())
&& annotated != actual
{
tracing::debug!(
ifindex,
annotated,
actual,
"discarding uevent annotation from a recycled ifindex"
);
self.annotations.remove(&ifindex);
}
}
self.known.insert(ifindex, link);
let Some(info) = self.info(ifindex) else {
return;
};
self.emit(if first_sighting {
NetdevEvent::Added(info)
} else {
NetdevEvent::Changed(info)
});
}
NetworkEvent::DelLink(link) => {
let ifindex = link.ifindex();
let name = self
.known
.remove(&ifindex)
.and_then(|l| l.name().map(str::to_string))
.or_else(|| link.name().map(str::to_string));
self.annotations.remove(&ifindex);
self.emit(NetdevEvent::Removed { ifindex, name });
}
_ => {}
}
}
fn on_uevent(&mut self, event: Uevent) {
if event.subsystem != NET_SUBSYSTEM {
return;
}
let Some((ifindex, annotation)) = NetdevAnnotation::from_uevent(&event) else {
tracing::debug!(devpath = %event.devpath, action = %event.action,
"net uevent without IFINDEX=; cannot join");
return;
};
match event.action.as_str() {
"remove" => {
self.annotations.remove(&ifindex);
}
"bind" | "unbind" => {
let bound = event.action == "bind";
self.annotations.insert(ifindex, annotation.clone());
self.emit(if bound {
NetdevEvent::DriverBound {
ifindex,
annotation,
}
} else {
NetdevEvent::DriverUnbound {
ifindex,
annotation,
}
});
}
_ => {
let changed = self.annotations.get(&ifindex) != Some(&annotation);
self.annotations.insert(ifindex, annotation);
if changed && let Some(info) = self.info(ifindex) {
self.emit(NetdevEvent::Changed(info));
}
}
}
}
fn emit(&mut self, event: NetdevEvent) {
if let Some(store) = &self.store {
match &event {
NetdevEvent::Added(info) | NetdevEvent::Changed(info) => {
store.upsert(info.ifindex(), info.clone());
}
NetdevEvent::Removed { ifindex, .. } => {
store.remove(ifindex);
}
_ => {}
}
}
self.pending.push_back(event);
}
fn poll_links(&mut self, cx: &mut Context<'_>) -> Poll<Option<Result<()>>> {
let Some(links) = self.links.as_mut() else {
return Poll::Ready(None);
};
match Pin::new(links).poll_next(cx) {
Poll::Ready(Some(Ok(event))) => {
self.on_link_event(event);
Poll::Ready(Some(Ok(())))
}
Poll::Ready(Some(Err(e))) => Poll::Ready(Some(Err(e))),
Poll::Ready(None) => {
self.links = None;
Poll::Ready(None)
}
Poll::Pending => Poll::Pending,
}
}
fn poll_uevents(&mut self, cx: &mut Context<'_>) -> Poll<Option<Result<()>>> {
let Some(uevents) = self.uevents.as_mut() else {
return Poll::Ready(None);
};
match Pin::new(uevents).poll_next(cx) {
Poll::Ready(Some(Ok(event))) => {
self.on_uevent(event);
Poll::Ready(Some(Ok(())))
}
Poll::Ready(Some(Err(e))) => Poll::Ready(Some(Err(e))),
Poll::Ready(None) => {
self.uevents = None;
Poll::Ready(None)
}
Poll::Pending => Poll::Pending,
}
}
}
impl<L, U> Unpin for NetdevLifecycle<L, U> {}
impl<L, U> Stream for NetdevLifecycle<L, U>
where
L: Stream<Item = Result<NetworkEvent>> + Unpin,
U: Stream<Item = Result<Uevent>> + Unpin,
{
type Item = Result<NetdevEvent>;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let this = self.get_mut();
loop {
if let Some(event) = this.pending.pop_front() {
return Poll::Ready(Some(Ok(event)));
}
if this.links.is_none() && this.uevents.is_none() {
return Poll::Ready(None);
}
this.prefer_uevents = !this.prefer_uevents;
let (first, second) = if this.prefer_uevents {
(
Self::poll_uevents as fn(&mut Self, &mut Context<'_>) -> _,
Self::poll_links as fn(&mut Self, &mut Context<'_>) -> _,
)
} else {
(
Self::poll_links as fn(&mut Self, &mut Context<'_>) -> _,
Self::poll_uevents as fn(&mut Self, &mut Context<'_>) -> _,
)
};
let a = first(this, cx);
if let Poll::Ready(Some(Err(e))) = a {
return Poll::Ready(Some(Err(e)));
}
let b = second(this, cx);
if let Poll::Ready(Some(Err(e))) = b {
return Poll::Ready(Some(Err(e)));
}
let produced = matches!(a, Poll::Ready(Some(Ok(()))))
|| matches!(b, Poll::Ready(Some(Ok(()))));
if !produced {
if this.links.is_none() && this.uevents.is_none() {
return Poll::Ready(None);
}
return Poll::Pending;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::{cell::RefCell, rc::Rc, task::Waker};
fn link(ifindex: u32, name: &str) -> LinkMessage {
crate::netlink::messages::LinkMessageBuilder::new()
.ifindex(ifindex as i32)
.name(name)
.build()
}
fn uevent(action: &str, ifindex: u32, interface: &str, driver: Option<&str>) -> Uevent {
let mut env = HashMap::new();
env.insert("ACTION".to_string(), action.to_string());
env.insert("SUBSYSTEM".to_string(), "net".to_string());
env.insert("IFINDEX".to_string(), ifindex.to_string());
env.insert("INTERFACE".to_string(), interface.to_string());
if let Some(driver) = driver {
env.insert("DRIVER".to_string(), driver.to_string());
}
Uevent {
action: action.to_string(),
devpath: format!("/devices/virtual/net/{interface}"),
subsystem: "net".to_string(),
env,
}
}
#[derive(Clone)]
struct Feed<T> {
items: Rc<RefCell<VecDeque<Result<T>>>>,
closed: Rc<RefCell<bool>>,
}
impl<T> Feed<T> {
fn new() -> Self {
Self {
items: Rc::new(RefCell::new(VecDeque::new())),
closed: Rc::new(RefCell::new(false)),
}
}
fn push(&self, item: Result<T>) {
self.items.borrow_mut().push_back(item);
}
fn close(&self) {
*self.closed.borrow_mut() = true;
}
}
impl<T: Unpin> Stream for Feed<T> {
type Item = Result<T>;
fn poll_next(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
if let Some(item) = self.items.borrow_mut().pop_front() {
Poll::Ready(Some(item))
} else if *self.closed.borrow() {
Poll::Ready(None)
} else {
Poll::Pending
}
}
}
struct Harness {
lifecycle: NetdevLifecycle<Feed<NetworkEvent>, Feed<Uevent>>,
links: Feed<NetworkEvent>,
uevents: Feed<Uevent>,
}
impl Harness {
fn new() -> Self {
Self::with_store(None)
}
fn with_store(store: Option<Store<u32, NetdevInfo>>) -> Self {
let links = Feed::new();
let uevents = Feed::new();
let mut lifecycle = NetdevLifecycle::new(links.clone(), uevents.clone());
if let Some(store) = store {
lifecycle = lifecycle.with_store(store);
}
Self {
lifecycle,
links,
uevents,
}
}
fn drain(&mut self) -> Vec<NetdevEvent> {
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
let mut out = Vec::new();
loop {
match Pin::new(&mut self.lifecycle).poll_next(&mut cx) {
Poll::Ready(Some(Ok(event))) => out.push(event),
Poll::Ready(Some(Err(e))) => panic!("unexpected stream error: {e}"),
Poll::Ready(None) | Poll::Pending => return out,
}
}
}
fn feed_link(&mut self, event: NetworkEvent) -> Vec<NetdevEvent> {
self.links.push(Ok(event));
self.drain()
}
fn feed_uevent(&mut self, event: Uevent) -> Vec<NetdevEvent> {
self.uevents.push(Ok(event));
self.drain()
}
}
#[test]
fn rtnetlink_alone_produces_added_without_an_annotation() {
let mut h = Harness::new();
let events = h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
assert_eq!(events.len(), 1, "{events:?}");
let NetdevEvent::Added(info) = &events[0] else {
panic!("expected Added, got {:?}", events[0]);
};
assert_eq!(info.ifindex(), 3);
assert_eq!(info.name(), Some("veth0"));
assert!(info.annotation().is_none());
assert!(!info.is_fully_attributed());
}
#[test]
fn late_uevent_surfaces_as_changed() {
let mut h = Harness::new();
assert!(matches!(
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")))[..],
[NetdevEvent::Added(_)]
));
let events = h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
assert_eq!(events.len(), 1, "{events:?}");
let NetdevEvent::Changed(info) = &events[0] else {
panic!("expected Changed, got {:?}", events[0]);
};
assert_eq!(info.driver(), Some("veth"));
assert_eq!(info.devpath(), Some("/devices/virtual/net/veth0"));
assert!(info.is_fully_attributed());
}
#[test]
fn early_uevent_rides_along_on_added() {
let mut h = Harness::new();
let events = h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
assert!(events.is_empty(), "{events:?}");
let events = h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
assert_eq!(events.len(), 1, "{events:?}");
let NetdevEvent::Added(info) = &events[0] else {
panic!("expected Added, got {:?}", events[0]);
};
assert_eq!(info.driver(), Some("veth"));
assert!(info.is_fully_attributed());
}
#[test]
fn a_redundant_uevent_emits_nothing() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
assert_eq!(h.feed_uevent(uevent("add", 3, "veth0", Some("veth"))).len(), 1);
let events = h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
assert!(events.is_empty(), "{events:?}");
}
#[test]
fn removal_reports_the_last_known_name() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
let events = h.feed_link(NetworkEvent::DelLink(link(3, "veth0")));
assert_eq!(events.len(), 1, "{events:?}");
let NetdevEvent::Removed { ifindex, name } = &events[0] else {
panic!("expected Removed, got {:?}", events[0]);
};
assert_eq!(*ifindex, 3);
assert_eq!(name.as_deref(), Some("veth0"));
}
#[test]
fn recycled_ifindex_does_not_inherit_a_stale_annotation() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
h.feed_link(NetworkEvent::DelLink(link(3, "veth0")));
h.feed_uevent(uevent("change", 3, "veth0", Some("veth")));
let events = h.feed_link(NetworkEvent::NewLink(link(3, "eth9")));
assert_eq!(events.len(), 1, "{events:?}");
let NetdevEvent::Added(info) = &events[0] else {
panic!("expected Added, got {:?}", events[0]);
};
assert_eq!(info.name(), Some("eth9"));
assert!(
info.annotation().is_none(),
"stale annotation attached: {:?}",
info.annotation()
);
}
#[test]
fn rename_keeps_the_annotation_but_reports_it_unverified() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
let events = h.feed_link(NetworkEvent::NewLink(link(3, "veth1")));
let NetdevEvent::Changed(info) = &events[0] else {
panic!("expected Changed, got {:?}", events[0]);
};
assert_eq!(info.name(), Some("veth1"));
assert_eq!(info.driver(), Some("veth"));
assert!(!info.is_fully_attributed());
let events = h.feed_uevent(uevent("move", 3, "veth1", Some("veth")));
assert!(events[0].info().unwrap().is_fully_attributed());
}
#[test]
fn bind_and_unbind_are_emitted_without_an_rtnetlink_partner() {
let mut h = Harness::new();
let events = h.feed_uevent(uevent("bind", 7, "enp4s0", Some("igb")));
assert!(
matches!(events[..], [NetdevEvent::DriverBound { ifindex: 7, .. }]),
"{events:?}"
);
let events = h.feed_uevent(uevent("unbind", 7, "enp4s0", Some("igb")));
assert!(
matches!(events[..], [NetdevEvent::DriverUnbound { ifindex: 7, .. }]),
"{events:?}"
);
}
#[test]
fn a_uevent_remove_retires_the_annotation_without_removing_the_device() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
assert_eq!(h.feed_uevent(uevent("add", 3, "veth0", Some("veth"))).len(), 1);
let events = h.feed_uevent(uevent("remove", 3, "veth0", Some("veth")));
assert!(events.is_empty(), "{events:?}");
let events = h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
let info = events[0].info().unwrap();
assert!(info.annotation().is_none());
assert!(!info.is_fully_attributed());
}
#[test]
fn an_annotation_without_a_name_still_counts_as_attributed() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
let mut ev = uevent("add", 3, "veth0", Some("veth"));
ev.env.remove("INTERFACE");
let events = h.feed_uevent(ev);
let info = events[0].info().unwrap();
assert!(info.annotation().is_some());
assert_eq!(info.annotation().unwrap().interface, None);
assert!(info.is_fully_attributed());
}
#[test]
fn the_annotation_is_the_uevent_environment_verbatim() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
let mut ev = uevent("add", 3, "veth0", Some("veth"));
ev.env.insert("DEVTYPE".to_string(), "veth".to_string());
ev.env.insert("SEQNUM".to_string(), "4242".to_string());
let events = h.feed_uevent(ev);
let annotation = events[0].info().unwrap().annotation().unwrap();
assert_eq!(annotation.devpath, "/devices/virtual/net/veth0");
assert_eq!(annotation.interface.as_deref(), Some("veth0"));
assert_eq!(annotation.driver.as_deref(), Some("veth"));
assert_eq!(annotation.devtype.as_deref(), Some("veth"));
assert_eq!(annotation.seqnum, Some(4242));
}
#[test]
fn device_count_tracks_rtnetlink_only() {
let mut h = Harness::new();
assert!(h.lifecycle.is_empty());
h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
assert!(h.lifecycle.is_empty());
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
h.feed_link(NetworkEvent::NewLink(link(4, "veth1")));
assert_eq!(h.lifecycle.len(), 2);
h.feed_link(NetworkEvent::DelLink(link(3, "veth0")));
assert_eq!(h.lifecycle.len(), 1);
}
#[test]
fn removal_of_an_unknown_device_is_still_reported() {
let mut h = Harness::new();
let events = h.feed_link(NetworkEvent::DelLink(link(9, "ghost0")));
let NetdevEvent::Removed { ifindex, name } = &events[0] else {
panic!("expected Removed, got {:?}", events[0]);
};
assert_eq!(*ifindex, 9);
assert_eq!(name.as_deref(), Some("ghost0"));
}
#[test]
fn non_net_uevents_and_ifindexless_uevents_are_ignored() {
let mut h = Harness::new();
let mut usb = uevent("add", 1, "irrelevant", None);
usb.subsystem = "usb".to_string();
assert!(h.feed_uevent(usb).is_empty());
let mut no_ifindex = uevent("add", 3, "veth0", None);
no_ifindex.env.remove("IFINDEX");
assert!(h.feed_uevent(no_ifindex).is_empty());
}
#[test]
fn ifindex_is_available_on_every_variant() {
let mut h = Harness::new();
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
let events = h.feed_uevent(uevent("bind", 3, "veth0", Some("veth")));
assert_eq!(events[0].ifindex(), 3);
let events = h.feed_link(NetworkEvent::DelLink(link(3, "veth0")));
assert_eq!(events[0].ifindex(), 3);
}
#[test]
fn store_mirrors_the_joined_state() {
let store: Store<u32, NetdevInfo> = Store::new();
let mut h = Harness::with_store(Some(store.clone()));
h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
h.feed_link(NetworkEvent::NewLink(link(4, "veth1")));
h.feed_uevent(uevent("add", 3, "veth0", Some("veth")));
assert_eq!(store.len(), 2);
assert_eq!(store.get(&3).unwrap().driver(), Some("veth"));
assert!(store.get(&4).unwrap().annotation().is_none());
h.feed_link(NetworkEvent::DelLink(link(3, "veth0")));
assert_eq!(store.len(), 1);
assert!(!store.contains_key(&3));
}
#[test]
fn store_ignores_driver_events_for_unknown_devices() {
let store: Store<u32, NetdevInfo> = Store::new();
let mut h = Harness::with_store(Some(store.clone()));
h.feed_uevent(uevent("bind", 7, "enp4s0", Some("igb")));
assert!(store.is_empty());
}
#[test]
fn errors_from_either_source_propagate() {
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
let mut h = Harness::new();
h.links
.push(Err(crate::Error::InvalidMessage("boom".into())));
assert!(matches!(
Pin::new(&mut h.lifecycle).poll_next(&mut cx),
Poll::Ready(Some(Err(_)))
));
let mut h = Harness::new();
h.uevents
.push(Err(crate::Error::InvalidMessage("boom".into())));
assert!(matches!(
Pin::new(&mut h.lifecycle).poll_next(&mut cx),
Poll::Ready(Some(Err(_)))
));
}
#[test]
fn one_source_ending_does_not_end_or_spin_the_join() {
let mut h = Harness::new();
h.uevents.close();
let events = h.feed_link(NetworkEvent::NewLink(link(3, "veth0")));
assert!(matches!(events[..], [NetdevEvent::Added(_)]), "{events:?}");
assert!(h.drain().is_empty());
h.links.close();
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
assert!(matches!(
Pin::new(&mut h.lifecycle).poll_next(&mut cx),
Poll::Ready(None)
));
}
#[test]
fn link_events_other_than_new_and_del_are_ignored() {
let mut h = Harness::new();
let events = h.feed_link(NetworkEvent::NewRoute(Default::default()));
assert!(events.is_empty(), "{events:?}");
assert!(h.lifecycle.is_empty());
}
}