include more link layer mirrors
This commit is contained in:
@@ -100,6 +100,11 @@ class BridgeMemberLinkStateInfo(BaseModel):
|
||||
admin_up: Optional[bool] = Field(None, description="Whether the interface has IFF_UP set.")
|
||||
carrier_up: Optional[bool] = Field(None, description="Whether the interface currently reports carrier.")
|
||||
operstate: Optional[str] = Field(None, description="Kernel operational state string.")
|
||||
mtu: Optional[int] = Field(None, description="Current interface MTU.")
|
||||
ethernet_profile: Optional[Dict[str, Any]] = Field(
|
||||
None,
|
||||
description="Current speed/duplex/autoneg profile when available.",
|
||||
)
|
||||
link_ready: bool = Field(..., description="Whether the member currently looks usable for forwarding.")
|
||||
suppressed: bool = Field(..., description="Whether the watcher administratively suppressed this member.")
|
||||
|
||||
@@ -115,6 +120,10 @@ class BridgeLinkStateWatcherStatus(BaseModel):
|
||||
last_action: Optional[str] = Field(None, description="Most recent propagation action.")
|
||||
suppressed_members: List[str] = Field(default_factory=list, description="Members currently forced down by the watcher.")
|
||||
recovery_holdoff_seconds: Optional[float] = Field(None, description="Configured recovery holdoff before restoring siblings.")
|
||||
degraded_recheck_seconds: Optional[float] = Field(
|
||||
None,
|
||||
description="Low-rate fallback recheck interval while the bridge is degraded.",
|
||||
)
|
||||
members: Dict[str, BridgeMemberLinkStateInfo] = Field(default_factory=dict, description="Per-member link-state snapshot.")
|
||||
message: Optional[str] = Field(None, description="Optional informational message.")
|
||||
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
"""Event-driven watcher that propagates bridge member link failures to sibling ports."""
|
||||
"""Event-driven watcher that propagates bridge member failures and selected link settings."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import select
|
||||
import subprocess
|
||||
import threading
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
@@ -18,20 +19,42 @@ from src.utilities.interface_bridge_helpers import (
|
||||
get_bridge_ports_once,
|
||||
read_interface_admin_up,
|
||||
read_interface_carrier,
|
||||
read_interface_mtu,
|
||||
read_interface_operstate,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("bridge_link_state_manager")
|
||||
|
||||
_ETHTOOL_BIN = "/usr/sbin/ethtool" if os.path.exists("/usr/sbin/ethtool") else "ethtool"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class EthernetProfile:
|
||||
"""Subset of ethtool settings that can be mirrored across peer ports."""
|
||||
|
||||
speed_mbps: Optional[int]
|
||||
duplex: Optional[str]
|
||||
autoneg: Optional[bool]
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
"""Serialize the profile for API responses."""
|
||||
return {
|
||||
"speed_mbps": self.speed_mbps,
|
||||
"duplex": self.duplex,
|
||||
"autoneg": self.autoneg,
|
||||
}
|
||||
|
||||
|
||||
@dataclass
|
||||
class MemberLinkState:
|
||||
"""Current sysfs-derived state for one bridge member."""
|
||||
"""Current state for one bridge member."""
|
||||
|
||||
ifname: str
|
||||
admin_up: Optional[bool]
|
||||
carrier_up: Optional[bool]
|
||||
operstate: Optional[str]
|
||||
mtu: Optional[int]
|
||||
ethernet_profile: Optional[EthernetProfile]
|
||||
|
||||
@property
|
||||
def link_ready(self) -> bool:
|
||||
@@ -51,6 +74,8 @@ class MemberLinkState:
|
||||
"admin_up": self.admin_up,
|
||||
"carrier_up": self.carrier_up,
|
||||
"operstate": self.operstate,
|
||||
"mtu": self.mtu,
|
||||
"ethernet_profile": self.ethernet_profile.to_dict() if self.ethernet_profile is not None else None,
|
||||
"link_ready": self.link_ready,
|
||||
"suppressed": suppressed,
|
||||
}
|
||||
@@ -75,6 +100,7 @@ class BridgeLinkStateWatcher:
|
||||
self._running = False
|
||||
self._suppressed_members: dict[str, bool] = {}
|
||||
self._settle_deadlines: dict[str, float] = {}
|
||||
self._managed_event_deadlines: dict[str, float] = {}
|
||||
self._all_clear_since: Optional[float] = None
|
||||
self._degraded = False
|
||||
self._last_event_ts: Optional[float] = None
|
||||
@@ -163,12 +189,13 @@ class BridgeLinkStateWatcher:
|
||||
continue
|
||||
|
||||
if any(msg.get("event") in {"RTM_NEWLINK", "RTM_DELLINK"} for msg in messages):
|
||||
self._evaluate_bridge_state(reason="netlink_event")
|
||||
source_ifname = self._pick_source_interface(messages)
|
||||
self._evaluate_bridge_state(reason="netlink_event", source_ifname=source_ifname)
|
||||
continue
|
||||
|
||||
self._evaluate_bridge_state(reason="recovery_deadline")
|
||||
|
||||
def _evaluate_bridge_state(self, reason: str) -> None:
|
||||
def _evaluate_bridge_state(self, reason: str, source_ifname: Optional[str] = None) -> None:
|
||||
members = [iface for iface in get_bridge_ports_once(self.bridge_name) if check_interface_exists(iface)]
|
||||
states = {iface: self._read_member_state(iface) for iface in members}
|
||||
now = time.time()
|
||||
@@ -187,6 +214,11 @@ class BridgeLinkStateWatcher:
|
||||
for iface, deadline in self._settle_deadlines.items()
|
||||
if iface in states and deadline > now
|
||||
}
|
||||
self._managed_event_deadlines = {
|
||||
iface: deadline
|
||||
for iface, deadline in self._managed_event_deadlines.items()
|
||||
if iface in states and deadline > now
|
||||
}
|
||||
|
||||
if len(states) < 2:
|
||||
with self._lock:
|
||||
@@ -197,9 +229,7 @@ class BridgeLinkStateWatcher:
|
||||
|
||||
with self._lock:
|
||||
suppressed_snapshot = set(self._suppressed_members)
|
||||
settling_snapshot = {
|
||||
iface for iface, deadline in self._settle_deadlines.items() if deadline > now
|
||||
}
|
||||
settling_snapshot = {iface for iface, deadline in self._settle_deadlines.items() if deadline > now}
|
||||
|
||||
failing_members = sorted(
|
||||
ifname
|
||||
@@ -233,17 +263,42 @@ class BridgeLinkStateWatcher:
|
||||
self._degraded = False
|
||||
self._all_clear_since = None
|
||||
|
||||
self._restore_suppressed_members(reason=reason)
|
||||
if source_ifname in states and self._is_source_eligible(source_ifname, states):
|
||||
self._sync_member_configuration(source_ifname, states)
|
||||
return
|
||||
|
||||
self._set_last_action(f"no_restore_needed ({reason})")
|
||||
|
||||
def _pick_source_interface(self, messages: list[dict[str, Any]]) -> Optional[str]:
|
||||
"""Pick the most relevant bridge member from a batch of netlink messages."""
|
||||
members = set(get_bridge_ports_once(self.bridge_name))
|
||||
source_ifname: Optional[str] = None
|
||||
|
||||
for message in messages:
|
||||
attrs = dict(message.get("attrs", []))
|
||||
ifname = attrs.get("IFLA_IFNAME")
|
||||
if ifname in members:
|
||||
source_ifname = ifname
|
||||
|
||||
return source_ifname
|
||||
|
||||
def _is_source_eligible(self, ifname: str, states: dict[str, MemberLinkState]) -> bool:
|
||||
"""Return whether this member should be trusted as the configuration source."""
|
||||
with self._lock:
|
||||
ignored = self._managed_event_deadlines.get(ifname, 0.0) > time.time()
|
||||
suppressed = ifname in self._suppressed_members
|
||||
if ignored or suppressed:
|
||||
return False
|
||||
|
||||
state = states.get(ifname)
|
||||
return state is not None and state.link_ready
|
||||
|
||||
def _next_wait_timeout(self) -> Optional[float]:
|
||||
"""Return how long the watcher may sleep before the next restore deadline."""
|
||||
with self._lock:
|
||||
if self._suppressed_members and self._all_clear_since is not None:
|
||||
deadline = self._all_clear_since + self.recovery_holdoff_seconds
|
||||
return max(
|
||||
0.0,
|
||||
min(deadline - time.time(), settings.bridge_link_state_degraded_recheck_seconds),
|
||||
)
|
||||
return max(0.0, min(deadline - time.time(), settings.bridge_link_state_degraded_recheck_seconds))
|
||||
if self._degraded:
|
||||
return settings.bridge_link_state_degraded_recheck_seconds
|
||||
return None
|
||||
@@ -271,8 +326,155 @@ class BridgeLinkStateWatcher:
|
||||
admin_up=read_interface_admin_up(ifname),
|
||||
carrier_up=read_interface_carrier(ifname),
|
||||
operstate=read_interface_operstate(ifname),
|
||||
mtu=read_interface_mtu(ifname),
|
||||
ethernet_profile=self._read_ethernet_profile(ifname),
|
||||
)
|
||||
|
||||
def _read_ethernet_profile(self, ifname: str) -> Optional[EthernetProfile]:
|
||||
"""Read ethtool speed/duplex/autoneg for one interface if supported."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
[_ETHTOOL_BIN, ifname],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
check=True,
|
||||
)
|
||||
except (FileNotFoundError, subprocess.CalledProcessError):
|
||||
return None
|
||||
|
||||
values: dict[str, str] = {}
|
||||
for line in result.stdout.splitlines():
|
||||
if ":" not in line:
|
||||
continue
|
||||
key, value = line.split(":", 1)
|
||||
values[key.strip()] = value.strip()
|
||||
|
||||
speed_value = values.get("Speed")
|
||||
duplex_value = values.get("Duplex")
|
||||
autoneg_value = values.get("Auto-negotiation")
|
||||
|
||||
speed_mbps: Optional[int] = None
|
||||
if speed_value and speed_value.endswith("Mb/s"):
|
||||
try:
|
||||
speed_mbps = int(speed_value[:-4], 10)
|
||||
except ValueError:
|
||||
speed_mbps = None
|
||||
|
||||
duplex: Optional[str] = None
|
||||
if duplex_value and duplex_value.lower() in {"full", "half"}:
|
||||
duplex = duplex_value.lower()
|
||||
|
||||
autoneg: Optional[bool] = None
|
||||
if autoneg_value:
|
||||
lowered = autoneg_value.lower()
|
||||
if lowered in {"on", "off"}:
|
||||
autoneg = lowered == "on"
|
||||
|
||||
if speed_mbps is None and duplex is None and autoneg is None:
|
||||
return None
|
||||
|
||||
return EthernetProfile(
|
||||
speed_mbps=speed_mbps,
|
||||
duplex=duplex,
|
||||
autoneg=autoneg,
|
||||
)
|
||||
|
||||
def _sync_member_configuration(self, source_ifname: str, states: dict[str, MemberLinkState]) -> None:
|
||||
"""Mirror MTU and ethtool link settings from one healthy member to its siblings."""
|
||||
source = states[source_ifname]
|
||||
changed_members: list[str] = []
|
||||
|
||||
for target_ifname, target in sorted(states.items()):
|
||||
if target_ifname == source_ifname:
|
||||
continue
|
||||
|
||||
mtu_changed = self._sync_member_mtu(source_ifname, source, target_ifname, target)
|
||||
profile_changed = self._sync_member_ethernet_profile(source_ifname, source, target_ifname, target)
|
||||
if mtu_changed or profile_changed:
|
||||
changed_members.append(target_ifname)
|
||||
|
||||
if changed_members:
|
||||
self._set_last_action(f"synchronized configuration from {source_ifname} to {changed_members}")
|
||||
|
||||
def _sync_member_mtu(
|
||||
self,
|
||||
source_ifname: str,
|
||||
source: MemberLinkState,
|
||||
target_ifname: str,
|
||||
target: MemberLinkState,
|
||||
) -> bool:
|
||||
"""Mirror MTU when the source member differs from the target."""
|
||||
if source.mtu is None or target.mtu is None or source.mtu == target.mtu:
|
||||
return False
|
||||
|
||||
with IPRoute() as ipr:
|
||||
indices = ipr.link_lookup(ifname=target_ifname)
|
||||
if not indices:
|
||||
raise RuntimeError(f"Interface {target_ifname} not found while synchronizing MTU")
|
||||
ipr.link("set", index=indices[0], mtu=source.mtu)
|
||||
|
||||
self._mark_managed_change(target_ifname)
|
||||
logger.info(
|
||||
"Synchronized MTU from %s to %s on bridge=%s: %s",
|
||||
source_ifname,
|
||||
target_ifname,
|
||||
self.bridge_name,
|
||||
source.mtu,
|
||||
)
|
||||
return True
|
||||
|
||||
def _sync_member_ethernet_profile(
|
||||
self,
|
||||
source_ifname: str,
|
||||
source: MemberLinkState,
|
||||
target_ifname: str,
|
||||
target: MemberLinkState,
|
||||
) -> bool:
|
||||
"""Mirror ethtool speed/duplex/autoneg from the source member to the target."""
|
||||
source_profile = source.ethernet_profile
|
||||
target_profile = target.ethernet_profile
|
||||
if source_profile is None or target_profile is None or source_profile == target_profile:
|
||||
return False
|
||||
|
||||
cmd = [_ETHTOOL_BIN, "-s", target_ifname]
|
||||
if source_profile.autoneg is True:
|
||||
cmd.extend(["autoneg", "on"])
|
||||
elif source_profile.autoneg is False and source_profile.speed_mbps is not None and source_profile.duplex is not None:
|
||||
cmd.extend(
|
||||
[
|
||||
"speed",
|
||||
str(source_profile.speed_mbps),
|
||||
"duplex",
|
||||
source_profile.duplex,
|
||||
"autoneg",
|
||||
"off",
|
||||
]
|
||||
)
|
||||
else:
|
||||
return False
|
||||
|
||||
try:
|
||||
subprocess.run(cmd, capture_output=True, text=True, check=True)
|
||||
except (FileNotFoundError, subprocess.CalledProcessError) as exc:
|
||||
logger.debug(
|
||||
"Failed to synchronize ethtool profile from %s to %s on bridge=%s: %s",
|
||||
source_ifname,
|
||||
target_ifname,
|
||||
self.bridge_name,
|
||||
exc,
|
||||
)
|
||||
return False
|
||||
|
||||
self._mark_managed_change(target_ifname)
|
||||
logger.info(
|
||||
"Synchronized ethtool profile from %s to %s on bridge=%s: %s",
|
||||
source_ifname,
|
||||
target_ifname,
|
||||
self.bridge_name,
|
||||
source_profile,
|
||||
)
|
||||
return True
|
||||
|
||||
def _suppress_other_members(self, states: dict[str, MemberLinkState], failing_members: list[str]) -> None:
|
||||
desired_suppressed = set(states) - set(failing_members)
|
||||
|
||||
@@ -303,8 +505,7 @@ class BridgeLinkStateWatcher:
|
||||
suppressed = dict(self._suppressed_members)
|
||||
|
||||
if not suppressed:
|
||||
with self._lock:
|
||||
self._last_action = f"no_restore_needed ({reason})"
|
||||
self._set_last_action(f"no_restore_needed ({reason})")
|
||||
return
|
||||
|
||||
restored_members: list[str] = []
|
||||
@@ -334,6 +535,17 @@ class BridgeLinkStateWatcher:
|
||||
if target_up:
|
||||
with self._lock:
|
||||
self._settle_deadlines[ifname] = time.time() + self.recovery_holdoff_seconds
|
||||
self._mark_managed_change(ifname)
|
||||
|
||||
def _mark_managed_change(self, ifname: str) -> None:
|
||||
"""Ignore immediate follow-up netlink events caused by our own changes."""
|
||||
with self._lock:
|
||||
self._managed_event_deadlines[ifname] = time.time() + self.recovery_holdoff_seconds
|
||||
|
||||
def _set_last_action(self, action: str) -> None:
|
||||
"""Update the watcher action text."""
|
||||
with self._lock:
|
||||
self._last_action = action
|
||||
|
||||
|
||||
class BridgeLinkStateManager:
|
||||
|
||||
@@ -73,6 +73,19 @@ def read_interface_admin_up(iface: str) -> Optional[bool]:
|
||||
return None
|
||||
|
||||
|
||||
def read_interface_mtu(iface: str) -> Optional[int]:
|
||||
"""Return the interface MTU if available."""
|
||||
value = _read_sysfs_text(f"/sys/class/net/{iface}/mtu")
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
try:
|
||||
return int(value, 10)
|
||||
except ValueError:
|
||||
logger.warning("Unexpected MTU value for %s: %r", iface, value)
|
||||
return None
|
||||
|
||||
|
||||
def _read_bridge_ports_from_sysfs(bridge: str) -> List[str]:
|
||||
"""
|
||||
Read bridge member interfaces from sysfs. Internal helper that always reads.
|
||||
|
||||
Reference in New Issue
Block a user