995 lines
38 KiB
Python
995 lines
38 KiB
Python
"""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
|
|
from typing import Any, Dict, Optional
|
|
|
|
from pyroute2 import IPRoute
|
|
|
|
from src.config import settings
|
|
from src.utilities.interface_bridge_helpers import (
|
|
check_interface_exists,
|
|
get_bridge_ports_once,
|
|
read_interface_admin_up,
|
|
read_interface_carrier,
|
|
read_interface_ethernet_profile,
|
|
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 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:
|
|
"""Return whether the member currently looks healthy enough to forward."""
|
|
if self.admin_up is not True:
|
|
return False
|
|
if self.carrier_up is False:
|
|
return False
|
|
if self.operstate in {"down", "lowerlayerdown", "notpresent"}:
|
|
return False
|
|
return True
|
|
|
|
def to_dict(self, suppressed: bool = False) -> Dict[str, Any]:
|
|
"""Serialize member state for API responses."""
|
|
return {
|
|
"ifname": self.ifname,
|
|
"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,
|
|
}
|
|
|
|
|
|
class BridgeLinkStateWatcher:
|
|
"""Watch one bridge and mirror member failures to the other bridge members."""
|
|
|
|
def __init__(self, bridge_name: str, recovery_holdoff_seconds: float) -> None:
|
|
self.bridge_name = bridge_name
|
|
self.recovery_holdoff_seconds = recovery_holdoff_seconds
|
|
self._stop_event = threading.Event()
|
|
self._lock = threading.Lock()
|
|
self._wake_r, self._wake_w = os.pipe()
|
|
os.set_blocking(self._wake_r, False)
|
|
os.set_blocking(self._wake_w, False)
|
|
self._thread = threading.Thread(
|
|
target=self._run,
|
|
daemon=True,
|
|
name=f"bridge-link-state-{bridge_name}",
|
|
)
|
|
self._running = False
|
|
self._suppressed_members: dict[str, bool] = {}
|
|
self._failing_since: dict[str, float] = {}
|
|
self._settle_deadlines: dict[str, float] = {}
|
|
self._managed_event_deadlines: dict[str, float] = {}
|
|
self._config_change_deadlines: dict[str, float] = {}
|
|
self._sync_attempt_deadlines: dict[tuple[str, str], tuple[str, float]] = {}
|
|
self._all_clear_since: Optional[float] = None
|
|
self._degraded = False
|
|
self._last_event_ts: Optional[float] = None
|
|
self._last_error: Optional[str] = None
|
|
self._last_action: Optional[str] = None
|
|
self._member_states: dict[str, MemberLinkState] = {}
|
|
|
|
def start(self) -> None:
|
|
"""Start the watcher thread."""
|
|
with self._lock:
|
|
if self._running:
|
|
return
|
|
self._running = True
|
|
self._thread.start()
|
|
|
|
def stop(self) -> None:
|
|
"""Stop the watcher and restore interfaces that were suppressed by it."""
|
|
self._stop_event.set()
|
|
self._wake_thread()
|
|
if self._thread.is_alive():
|
|
self._thread.join(timeout=settings.bridge_link_state_thread_join_timeout_seconds)
|
|
|
|
try:
|
|
self._restore_suppressed_members(reason="watcher_stopped")
|
|
except Exception:
|
|
logger.exception("Failed to restore suppressed members for bridge=%s", self.bridge_name)
|
|
|
|
with self._lock:
|
|
self._running = False
|
|
|
|
for fd in (self._wake_r, self._wake_w):
|
|
try:
|
|
os.close(fd)
|
|
except OSError:
|
|
pass
|
|
|
|
def status(self) -> dict[str, Any]:
|
|
"""Return a JSON-serializable snapshot of the watcher state."""
|
|
with self._lock:
|
|
members = {
|
|
ifname: state.to_dict(suppressed=ifname in self._suppressed_members)
|
|
for ifname, state in sorted(self._member_states.items())
|
|
}
|
|
return {
|
|
"bridge": self.bridge_name,
|
|
"active": self._running and self._thread.is_alive() and not self._stop_event.is_set(),
|
|
"event_driven": True,
|
|
"last_event_ts": self._last_event_ts,
|
|
"last_error": self._last_error,
|
|
"last_action": self._last_action,
|
|
"suppressed_members": sorted(self._suppressed_members),
|
|
"failure_holdoff_seconds": settings.bridge_link_state_failure_holdoff_seconds,
|
|
"recovery_holdoff_seconds": self.recovery_holdoff_seconds,
|
|
"degraded_recheck_seconds": settings.bridge_link_state_degraded_recheck_seconds,
|
|
"members": members,
|
|
}
|
|
|
|
def _run(self) -> None:
|
|
with IPRoute() as ipr:
|
|
ipr.bind()
|
|
self._evaluate_bridge_state(reason="watcher_started")
|
|
|
|
while not self._stop_event.is_set():
|
|
try:
|
|
timeout = self._next_wait_timeout()
|
|
ready, _, _ = select.select([ipr, self._wake_r], [], [], timeout)
|
|
except Exception as exc:
|
|
logger.exception("Bridge link-state select failed for bridge=%s", self.bridge_name)
|
|
with self._lock:
|
|
self._last_error = str(exc)
|
|
continue
|
|
|
|
if self._stop_event.is_set():
|
|
break
|
|
|
|
if self._wake_r in ready:
|
|
self._drain_wake_pipe()
|
|
continue
|
|
|
|
if ipr in ready:
|
|
try:
|
|
messages = ipr.get()
|
|
except Exception as exc:
|
|
logger.exception("Bridge link-state netlink read failed for bridge=%s", self.bridge_name)
|
|
with self._lock:
|
|
self._last_error = str(exc)
|
|
continue
|
|
|
|
if any(msg.get("event") in {"RTM_NEWLINK", "RTM_DELLINK"} for msg in messages):
|
|
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, source_ifname: Optional[str] = None) -> None:
|
|
with self._lock:
|
|
previous_states = dict(self._member_states)
|
|
|
|
members = [iface for iface in get_bridge_ports_once(self.bridge_name) if check_interface_exists(iface)]
|
|
states = {
|
|
iface: self._read_member_state(iface, previous_states.get(iface))
|
|
for iface in members
|
|
}
|
|
now = time.time()
|
|
|
|
with self._lock:
|
|
self._last_event_ts = now
|
|
self._last_error = None
|
|
self._member_states = states
|
|
self._suppressed_members = {
|
|
iface: restore_up
|
|
for iface, restore_up in self._suppressed_members.items()
|
|
if iface in states
|
|
}
|
|
self._failing_since = {
|
|
iface: first_seen
|
|
for iface, first_seen in self._failing_since.items()
|
|
if iface in states
|
|
}
|
|
self._settle_deadlines = {
|
|
iface: deadline
|
|
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
|
|
}
|
|
self._config_change_deadlines = {
|
|
iface: deadline
|
|
for iface, deadline in self._config_change_deadlines.items()
|
|
if iface in states and deadline > now
|
|
}
|
|
self._sync_attempt_deadlines = {
|
|
key: value
|
|
for key, value in self._sync_attempt_deadlines.items()
|
|
if key[0] in states and key[1] in states and value[1] > now
|
|
}
|
|
|
|
if len(states) < 2:
|
|
with self._lock:
|
|
self._degraded = False
|
|
self._failing_since = {}
|
|
self._all_clear_since = None
|
|
self._restore_suppressed_members(reason="bridge_has_fewer_than_two_members")
|
|
return
|
|
|
|
with self._lock:
|
|
suppressed_snapshot = set(self._suppressed_members)
|
|
settling_snapshot = {iface for iface, deadline in self._settle_deadlines.items() if deadline > now}
|
|
|
|
failing_members = sorted(
|
|
ifname
|
|
for ifname, state in states.items()
|
|
if ifname not in suppressed_snapshot and ifname not in settling_snapshot and not state.link_ready
|
|
)
|
|
changed_members = sorted(
|
|
ifname
|
|
for ifname in states
|
|
if self._config_changed(ifname, states, previous_states)
|
|
)
|
|
if changed_members:
|
|
self._mark_config_changes(changed_members, now)
|
|
|
|
if failing_members:
|
|
config_source = self._find_config_sync_source(states, previous_states, preferred_ifname=source_ifname)
|
|
if config_source is not None and self._has_config_mismatch(config_source, states):
|
|
with self._lock:
|
|
self._failing_since = {}
|
|
self._sync_member_configuration(config_source, states)
|
|
return
|
|
|
|
transition_members = sorted(
|
|
{
|
|
*changed_members,
|
|
*( [source_ifname] if source_ifname is not None else [] ),
|
|
*failing_members,
|
|
}
|
|
)
|
|
if self._config_transition_active(transition_members, now):
|
|
with self._lock:
|
|
self._failing_since = {}
|
|
self._set_last_action(
|
|
f"waiting for config transition to settle on {transition_members} before suppressing"
|
|
)
|
|
return
|
|
|
|
matured_failing_members = self._track_failing_members(failing_members, now)
|
|
if not matured_failing_members:
|
|
self._set_last_action(f"waiting before suppressing transient failures on {failing_members}")
|
|
return
|
|
|
|
with self._lock:
|
|
self._degraded = True
|
|
self._all_clear_since = None
|
|
self._suppress_other_members(states, matured_failing_members)
|
|
return
|
|
|
|
if suppressed_snapshot:
|
|
with self._lock:
|
|
self._failing_since = {}
|
|
should_restore = False
|
|
waiting_for_restore = False
|
|
with self._lock:
|
|
self._degraded = True
|
|
if self._all_clear_since is None:
|
|
self._all_clear_since = now
|
|
should_restore = now - self._all_clear_since >= self.recovery_holdoff_seconds
|
|
if not should_restore:
|
|
waiting_for_restore = True
|
|
|
|
if waiting_for_restore:
|
|
self._set_last_action("waiting before restoring suppressed members")
|
|
|
|
if should_restore:
|
|
self._restore_suppressed_members(reason=reason)
|
|
return
|
|
|
|
with self._lock:
|
|
self._degraded = False
|
|
self._failing_since = {}
|
|
self._all_clear_since = None
|
|
|
|
config_source = self._find_config_sync_source(states, previous_states, preferred_ifname=source_ifname)
|
|
if config_source is not None and self._has_config_mismatch(config_source, states):
|
|
self._sync_member_configuration(config_source, 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 _find_config_sync_source(
|
|
self,
|
|
states: dict[str, MemberLinkState],
|
|
previous_states: dict[str, MemberLinkState],
|
|
preferred_ifname: Optional[str] = None,
|
|
) -> Optional[str]:
|
|
"""Pick a healthy member whose configuration should be mirrored to siblings."""
|
|
changed_candidates = [
|
|
ifname
|
|
for ifname in sorted(states)
|
|
if self._config_changed(ifname, states, previous_states)
|
|
]
|
|
if preferred_ifname in changed_candidates:
|
|
changed_candidates.remove(preferred_ifname)
|
|
changed_candidates.insert(0, preferred_ifname)
|
|
|
|
if changed_candidates:
|
|
for ifname in changed_candidates:
|
|
if self._is_changed_source_eligible(ifname, states):
|
|
return ifname
|
|
self._set_last_action(
|
|
f"waiting for changed configuration on {changed_candidates} to settle before syncing"
|
|
)
|
|
return None
|
|
|
|
candidates: list[str] = []
|
|
if preferred_ifname:
|
|
candidates.append(preferred_ifname)
|
|
candidates.extend(ifname for ifname in sorted(states) if ifname != preferred_ifname)
|
|
|
|
seen: set[str] = set()
|
|
for ifname in candidates:
|
|
if ifname in seen:
|
|
continue
|
|
seen.add(ifname)
|
|
if self._is_source_eligible(ifname, states):
|
|
return ifname
|
|
|
|
return None
|
|
|
|
def _is_changed_source_eligible(self, ifname: str, states: dict[str, MemberLinkState]) -> bool:
|
|
"""Allow a changed member to drive sync even while the link is transiently renegotiating."""
|
|
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)
|
|
if state is None:
|
|
return False
|
|
if state.link_ready:
|
|
return True
|
|
return self._state_has_usable_config(state)
|
|
|
|
def _config_changed(
|
|
self,
|
|
ifname: str,
|
|
states: dict[str, MemberLinkState],
|
|
previous_states: dict[str, MemberLinkState],
|
|
) -> bool:
|
|
"""Return whether this member's MTU or link profile changed since the last snapshot."""
|
|
current = states.get(ifname)
|
|
previous = previous_states.get(ifname)
|
|
if current is None or previous is None:
|
|
return False
|
|
|
|
return current.mtu != previous.mtu or self._profiles_differ(
|
|
current.ethernet_profile,
|
|
previous.ethernet_profile,
|
|
)
|
|
|
|
def _state_has_usable_config(self, state: MemberLinkState) -> bool:
|
|
"""Return whether this snapshot contains configuration that can be mirrored."""
|
|
if state.mtu is not None:
|
|
return True
|
|
|
|
profile = state.ethernet_profile
|
|
if profile is None:
|
|
return False
|
|
if profile.autoneg is True:
|
|
return True
|
|
return profile.autoneg is False and profile.speed_mbps is not None and profile.duplex is not None
|
|
|
|
def _has_config_mismatch(self, source_ifname: str, states: dict[str, MemberLinkState]) -> bool:
|
|
"""Return whether any sibling differs from the chosen source configuration."""
|
|
source = states.get(source_ifname)
|
|
if source is None:
|
|
return False
|
|
|
|
for target_ifname, target in states.items():
|
|
if target_ifname == source_ifname:
|
|
continue
|
|
if source.mtu is not None and target.mtu is not None and source.mtu != target.mtu:
|
|
return True
|
|
if source.ethernet_profile is not None and source.ethernet_profile != target.ethernet_profile:
|
|
return True
|
|
|
|
return False
|
|
|
|
def _track_failing_members(self, failing_members: list[str], now: float) -> list[str]:
|
|
"""Record first-seen timestamps and return failures that exceeded the holdoff."""
|
|
with self._lock:
|
|
self._degraded = True
|
|
tracked = {
|
|
ifname: self._failing_since.get(ifname, now)
|
|
for ifname in failing_members
|
|
}
|
|
self._failing_since = tracked
|
|
matured = [
|
|
ifname
|
|
for ifname, first_seen in tracked.items()
|
|
if now - first_seen >= settings.bridge_link_state_failure_holdoff_seconds
|
|
]
|
|
|
|
return sorted(matured)
|
|
|
|
def _mark_config_changes(self, ifnames: list[str], now: float) -> None:
|
|
"""Keep short-lived link flaps from config changes from being treated as failures."""
|
|
holdoff = max(
|
|
settings.bridge_link_state_failure_holdoff_seconds,
|
|
self.recovery_holdoff_seconds,
|
|
)
|
|
with self._lock:
|
|
for ifname in ifnames:
|
|
self._config_change_deadlines[ifname] = now + holdoff
|
|
|
|
def _config_transition_active(self, ifnames: list[str], now: float) -> bool:
|
|
"""Return whether any listed member is still inside the config-change grace window."""
|
|
with self._lock:
|
|
return any(self._config_change_deadlines.get(ifname, 0.0) > now for ifname in ifnames)
|
|
|
|
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))
|
|
if self._failing_since:
|
|
earliest_deadline = min(
|
|
first_seen + settings.bridge_link_state_failure_holdoff_seconds
|
|
for first_seen in self._failing_since.values()
|
|
)
|
|
return max(0.0, min(earliest_deadline - time.time(), settings.bridge_link_state_degraded_recheck_seconds))
|
|
if self._degraded:
|
|
return settings.bridge_link_state_degraded_recheck_seconds
|
|
return None
|
|
|
|
def _wake_thread(self) -> None:
|
|
"""Wake the event loop from another thread."""
|
|
try:
|
|
os.write(self._wake_w, b"\x00")
|
|
except OSError:
|
|
pass
|
|
|
|
def _drain_wake_pipe(self) -> None:
|
|
"""Drain pending wake-up bytes from the control pipe."""
|
|
try:
|
|
while os.read(self._wake_r, 4096):
|
|
pass
|
|
except BlockingIOError:
|
|
return
|
|
except OSError:
|
|
return
|
|
|
|
def _read_member_state(
|
|
self,
|
|
ifname: str,
|
|
previous_state: Optional[MemberLinkState] = None,
|
|
) -> MemberLinkState:
|
|
admin_up = read_interface_admin_up(ifname)
|
|
carrier_up = read_interface_carrier(ifname)
|
|
operstate = read_interface_operstate(ifname)
|
|
link_ready = self._link_looks_ready(admin_up, carrier_up, operstate)
|
|
|
|
return MemberLinkState(
|
|
ifname=ifname,
|
|
admin_up=admin_up,
|
|
carrier_up=carrier_up,
|
|
operstate=operstate,
|
|
mtu=read_interface_mtu(ifname),
|
|
ethernet_profile=self._read_ethernet_profile(
|
|
ifname,
|
|
previous_state=previous_state,
|
|
link_ready=link_ready,
|
|
),
|
|
)
|
|
|
|
def _read_ethernet_profile(
|
|
self,
|
|
ifname: str,
|
|
previous_state: Optional[MemberLinkState] = None,
|
|
link_ready: Optional[bool] = None,
|
|
) -> Optional[EthernetProfile]:
|
|
"""Read ethtool speed/duplex/autoneg for one interface if supported."""
|
|
profile = read_interface_ethernet_profile(ifname)
|
|
if profile is None:
|
|
current_profile = None
|
|
else:
|
|
current_profile = EthernetProfile(
|
|
speed_mbps=profile.get("speed_mbps"),
|
|
duplex=profile.get("duplex"),
|
|
autoneg=profile.get("autoneg"),
|
|
)
|
|
|
|
if self._should_keep_previous_profile(ifname, current_profile, previous_state, link_ready):
|
|
return previous_state.ethernet_profile if previous_state is not None else None
|
|
|
|
previous_profile = previous_state.ethernet_profile if previous_state is not None else None
|
|
return self._materialize_profile(current_profile, previous_profile)
|
|
|
|
def _should_keep_previous_profile(
|
|
self,
|
|
ifname: str,
|
|
current_profile: Optional[EthernetProfile],
|
|
previous_state: Optional[MemberLinkState],
|
|
link_ready: Optional[bool],
|
|
) -> bool:
|
|
"""Keep the previous non-null profile during transient renegotiation windows."""
|
|
if previous_state is None or previous_state.ethernet_profile is None:
|
|
return False
|
|
if current_profile is not None and self._partial_profile_is_usable(current_profile):
|
|
return False
|
|
if link_ready is True:
|
|
return False
|
|
|
|
now = time.time()
|
|
with self._lock:
|
|
settling = self._settle_deadlines.get(ifname, 0.0) > now
|
|
managed = self._managed_event_deadlines.get(ifname, 0.0) > now
|
|
|
|
return settling or managed or link_ready is False
|
|
|
|
def _link_looks_ready(
|
|
self,
|
|
admin_up: Optional[bool],
|
|
carrier_up: Optional[bool],
|
|
operstate: Optional[str],
|
|
) -> bool:
|
|
"""Evaluate link health from sysfs fields before the MemberLinkState is built."""
|
|
if admin_up is not True:
|
|
return False
|
|
if carrier_up is False:
|
|
return False
|
|
if operstate in {"down", "lowerlayerdown", "notpresent"}:
|
|
return False
|
|
return True
|
|
|
|
def _profiles_differ(
|
|
self,
|
|
current_profile: Optional[EthernetProfile],
|
|
previous_profile: Optional[EthernetProfile],
|
|
) -> bool:
|
|
"""Ignore transient non-null to null drops when detecting config changes."""
|
|
if current_profile is None:
|
|
return False
|
|
if previous_profile is None:
|
|
return True
|
|
return current_profile != previous_profile
|
|
|
|
def _partial_profile_is_usable(self, profile: Optional[EthernetProfile]) -> bool:
|
|
"""Return whether a profile contains enough data to drive synchronization."""
|
|
if profile is None:
|
|
return False
|
|
if profile.autoneg is True:
|
|
return True
|
|
return profile.autoneg is False and profile.speed_mbps is not None and profile.duplex is not None
|
|
|
|
def _materialize_profile(
|
|
self,
|
|
current_profile: Optional[EthernetProfile],
|
|
previous_profile: Optional[EthernetProfile],
|
|
) -> Optional[EthernetProfile]:
|
|
"""Fill transiently missing profile fields from the last stable snapshot."""
|
|
if current_profile is None:
|
|
return previous_profile
|
|
if previous_profile is None:
|
|
return current_profile
|
|
if self._partial_profile_is_usable(current_profile):
|
|
return current_profile
|
|
|
|
return EthernetProfile(
|
|
speed_mbps=current_profile.speed_mbps if current_profile.speed_mbps is not None else previous_profile.speed_mbps,
|
|
duplex=current_profile.duplex if current_profile.duplex is not None else previous_profile.duplex,
|
|
autoneg=current_profile.autoneg if current_profile.autoneg is not None else previous_profile.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]
|
|
changes: list[str] = []
|
|
|
|
for target_ifname, target in sorted(states.items()):
|
|
if target_ifname == source_ifname:
|
|
continue
|
|
|
|
mtu_change = self._sync_member_mtu(source_ifname, source, target_ifname, target)
|
|
profile_change = self._sync_member_ethernet_profile(source_ifname, source, target_ifname, target)
|
|
if mtu_change:
|
|
changes.append(mtu_change)
|
|
if profile_change:
|
|
changes.append(profile_change)
|
|
|
|
if changes:
|
|
self._set_last_action(f"synchronized from {source_ifname}: {'; '.join(changes)}")
|
|
else:
|
|
self._set_last_action(f"no configuration mismatch detected after event on {source_ifname}")
|
|
|
|
def _sync_member_mtu(
|
|
self,
|
|
source_ifname: str,
|
|
source: MemberLinkState,
|
|
target_ifname: str,
|
|
target: MemberLinkState,
|
|
) -> Optional[str]:
|
|
"""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 None
|
|
|
|
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 f"{target_ifname} mtu={source.mtu}"
|
|
|
|
def _sync_member_ethernet_profile(
|
|
self,
|
|
source_ifname: str,
|
|
source: MemberLinkState,
|
|
target_ifname: str,
|
|
target: MemberLinkState,
|
|
) -> Optional[str]:
|
|
"""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 None
|
|
if not self._should_attempt_sync(source_ifname, target_ifname, source_profile):
|
|
return None
|
|
|
|
change_summary: Optional[str] = None
|
|
commands: list[tuple[list[str], str]] = []
|
|
if source_profile.autoneg is True:
|
|
if source_profile.speed_mbps is not None and source_profile.duplex is not None:
|
|
# A remote peer can pull this port down to a lower negotiated speed while autoneg
|
|
# stays enabled locally. Mirror that effective mode while keeping autoneg enabled
|
|
# on the sibling so its far-end peer can still negotiate successfully.
|
|
commands.append(
|
|
(
|
|
[
|
|
_ETHTOOL_BIN,
|
|
"-s",
|
|
target_ifname,
|
|
"speed",
|
|
str(source_profile.speed_mbps),
|
|
"duplex",
|
|
source_profile.duplex,
|
|
"autoneg",
|
|
"on",
|
|
],
|
|
(
|
|
f"{target_ifname} link={source_profile.speed_mbps}Mb/"
|
|
f"{source_profile.duplex}/autoneg-on"
|
|
),
|
|
)
|
|
)
|
|
else:
|
|
commands.append(([_ETHTOOL_BIN, "-s", target_ifname, "autoneg", "on"], f"{target_ifname} link=autoneg-on"))
|
|
elif source_profile.autoneg is False and source_profile.speed_mbps is not None and source_profile.duplex is not None:
|
|
# Do not force autoneg off on sibling ports. The far-end peer on that segment may
|
|
# still rely on autoneg, and forcing a fixed mode here can leave the link down.
|
|
commands.append(
|
|
(
|
|
[
|
|
_ETHTOOL_BIN,
|
|
"-s",
|
|
target_ifname,
|
|
"speed",
|
|
str(source_profile.speed_mbps),
|
|
"duplex",
|
|
source_profile.duplex,
|
|
"autoneg",
|
|
"on",
|
|
],
|
|
(
|
|
f"{target_ifname} link={source_profile.speed_mbps}Mb/"
|
|
f"{source_profile.duplex}/autoneg-on-safe"
|
|
),
|
|
)
|
|
)
|
|
commands.append(
|
|
(
|
|
[_ETHTOOL_BIN, "-s", target_ifname, "autoneg", "on"],
|
|
f"{target_ifname} link=autoneg-on-safe",
|
|
)
|
|
)
|
|
else:
|
|
return None
|
|
|
|
last_error: Optional[Exception] = None
|
|
for cmd, summary in commands:
|
|
try:
|
|
subprocess.run(cmd, capture_output=True, text=True, check=True)
|
|
change_summary = summary
|
|
break
|
|
except (FileNotFoundError, subprocess.CalledProcessError) as exc:
|
|
last_error = exc
|
|
logger.debug(
|
|
"Failed to synchronize ethtool profile from %s to %s on bridge=%s with %s: %s",
|
|
source_ifname,
|
|
target_ifname,
|
|
self.bridge_name,
|
|
cmd,
|
|
exc,
|
|
)
|
|
|
|
if change_summary is None:
|
|
if last_error is not None:
|
|
self._record_sync_attempt(source_ifname, target_ifname, source_profile)
|
|
self._set_last_action(
|
|
f"failed to sync link profile from {source_ifname} to {target_ifname}: {last_error}"
|
|
)
|
|
return None
|
|
|
|
self._record_sync_attempt(source_ifname, target_ifname, source_profile)
|
|
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 change_summary
|
|
|
|
def _sync_attempt_signature(self, source_profile: EthernetProfile) -> str:
|
|
"""Serialize the desired mirrored link state for retry deduplication."""
|
|
return f"{source_profile.speed_mbps}:{source_profile.duplex}:{source_profile.autoneg}"
|
|
|
|
def _should_attempt_sync(
|
|
self,
|
|
source_ifname: str,
|
|
target_ifname: str,
|
|
source_profile: EthernetProfile,
|
|
) -> bool:
|
|
"""Avoid replaying the same sync on every degraded-state recheck."""
|
|
signature = self._sync_attempt_signature(source_profile)
|
|
now = time.time()
|
|
with self._lock:
|
|
cached = self._sync_attempt_deadlines.get((source_ifname, target_ifname))
|
|
if cached is None:
|
|
return True
|
|
cached_signature, deadline = cached
|
|
return cached_signature != signature or deadline <= now
|
|
|
|
def _record_sync_attempt(
|
|
self,
|
|
source_ifname: str,
|
|
target_ifname: str,
|
|
source_profile: EthernetProfile,
|
|
) -> None:
|
|
"""Rate-limit repeated sync attempts for the same desired link profile."""
|
|
signature = self._sync_attempt_signature(source_profile)
|
|
cooldown = max(
|
|
settings.bridge_link_state_degraded_recheck_seconds * 4,
|
|
self.recovery_holdoff_seconds,
|
|
)
|
|
with self._lock:
|
|
self._sync_attempt_deadlines[(source_ifname, target_ifname)] = (signature, time.time() + cooldown)
|
|
|
|
def _suppress_other_members(self, states: dict[str, MemberLinkState], failing_members: list[str]) -> None:
|
|
desired_suppressed = set(states) - set(failing_members)
|
|
|
|
with self._lock:
|
|
current_suppressed = dict(self._suppressed_members)
|
|
|
|
next_suppressed: dict[str, bool] = {}
|
|
changed_members: list[str] = []
|
|
|
|
for ifname in sorted(desired_suppressed):
|
|
restore_up = current_suppressed.get(ifname, states[ifname].admin_up is True)
|
|
if ifname not in current_suppressed and states[ifname].admin_up is True:
|
|
self._set_interface_admin_state(ifname, target_up=False)
|
|
changed_members.append(ifname)
|
|
next_suppressed[ifname] = restore_up
|
|
|
|
with self._lock:
|
|
self._suppressed_members = next_suppressed
|
|
action = (
|
|
f"suppressed {changed_members} because failing members={failing_members}"
|
|
if changed_members
|
|
else f"holding suppressed members because failing members={failing_members}"
|
|
)
|
|
self._set_last_action(action)
|
|
|
|
def _restore_suppressed_members(self, reason: str) -> None:
|
|
with self._lock:
|
|
suppressed = dict(self._suppressed_members)
|
|
|
|
if not suppressed:
|
|
self._set_last_action(f"no_restore_needed ({reason})")
|
|
return
|
|
|
|
restored_members: list[str] = []
|
|
for ifname, restore_up in sorted(suppressed.items()):
|
|
if not check_interface_exists(ifname):
|
|
continue
|
|
if restore_up:
|
|
self._set_interface_admin_state(ifname, target_up=True)
|
|
restored_members.append(ifname)
|
|
|
|
with self._lock:
|
|
self._suppressed_members = {}
|
|
self._all_clear_since = None
|
|
self._set_last_action(f"restored {restored_members} ({reason})")
|
|
|
|
def _set_interface_admin_state(self, ifname: str, target_up: bool) -> None:
|
|
if not check_interface_exists(ifname):
|
|
raise RuntimeError(f"Interface {ifname} disappeared while propagating link state")
|
|
|
|
state_name = "up" if target_up else "down"
|
|
with IPRoute() as ipr:
|
|
indices = ipr.link_lookup(ifname=ifname)
|
|
if not indices:
|
|
raise RuntimeError(f"Interface {ifname} not found while propagating link state")
|
|
ipr.link("set", index=indices[0], state=state_name)
|
|
|
|
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."""
|
|
changed = False
|
|
with self._lock:
|
|
changed = action != self._last_action
|
|
self._last_action = action
|
|
|
|
if changed:
|
|
self._publish_network_state_update()
|
|
|
|
def _publish_network_state_update(self) -> None:
|
|
"""Publish a network snapshot after a meaningful watcher state change."""
|
|
try:
|
|
from src.api.network_api import publish_network_state_update
|
|
|
|
publish_network_state_update(f"bridge_watcher:{self.bridge_name}")
|
|
except Exception:
|
|
logger.exception("Failed to publish network state for bridge=%s", self.bridge_name)
|
|
|
|
|
|
class BridgeLinkStateManager:
|
|
"""Track bridge link-state watchers keyed by bridge name."""
|
|
|
|
def __init__(self) -> None:
|
|
self._watchers: dict[str, BridgeLinkStateWatcher] = {}
|
|
self._lock = threading.Lock()
|
|
|
|
def enable(self, bridge_name: str, recovery_holdoff_seconds: Optional[float] = None) -> dict[str, Any]:
|
|
"""Start or replace the watcher for the given bridge."""
|
|
normalized_holdoff = recovery_holdoff_seconds or settings.bridge_link_state_recovery_holdoff_seconds
|
|
with self._lock:
|
|
old_watcher = self._watchers.pop(bridge_name, None)
|
|
|
|
if old_watcher is not None:
|
|
old_watcher.stop()
|
|
|
|
watcher = BridgeLinkStateWatcher(bridge_name, normalized_holdoff)
|
|
watcher.start()
|
|
|
|
with self._lock:
|
|
self._watchers[bridge_name] = watcher
|
|
|
|
return watcher.status()
|
|
|
|
def disable(self, bridge_name: str) -> dict[str, Any]:
|
|
"""Stop the watcher for one bridge."""
|
|
with self._lock:
|
|
watcher = self._watchers.pop(bridge_name, None)
|
|
|
|
if watcher is None:
|
|
return {
|
|
"bridge": bridge_name,
|
|
"active": False,
|
|
"message": "watcher not enabled",
|
|
}
|
|
|
|
watcher.stop()
|
|
status = watcher.status()
|
|
status["active"] = False
|
|
return status
|
|
|
|
def get_status(self, bridge_name: str) -> Optional[dict[str, Any]]:
|
|
"""Return the current watcher status for one bridge, if present."""
|
|
with self._lock:
|
|
watcher = self._watchers.get(bridge_name)
|
|
return watcher.status() if watcher is not None else None
|
|
|
|
def list_statuses(self) -> list[dict[str, Any]]:
|
|
"""Return the current status for all bridge watchers."""
|
|
with self._lock:
|
|
watchers = list(self._watchers.values())
|
|
return [watcher.status() for watcher in sorted(watchers, key=lambda item: item.bridge_name)]
|
|
|
|
def stop(self) -> None:
|
|
"""Stop all running watchers."""
|
|
with self._lock:
|
|
watchers = list(self._watchers.values())
|
|
self._watchers.clear()
|
|
|
|
for watcher in watchers:
|
|
watcher.stop()
|
|
|
|
|
|
bridge_link_state_manager = BridgeLinkStateManager()
|