diff --git a/backend/src/api/network_api.py b/backend/src/api/network_api.py index cc15138..269927d 100644 --- a/backend/src/api/network_api.py +++ b/backend/src/api/network_api.py @@ -86,10 +86,10 @@ class BridgeRemoveRequest(BaseModel): class BridgeLinkStateEnableRequest(BaseModel): """Payload for enabling bridge member link-state propagation.""" - poll_interval_seconds: float = Field( - settings.bridge_link_state_poll_interval_seconds, + recovery_holdoff_seconds: float = Field( + settings.bridge_link_state_recovery_holdoff_seconds, gt=0, - description="Polling interval for member link-state checks.", + description="Holdoff before sibling interfaces are restored after recovery.", ) @@ -109,11 +109,12 @@ class BridgeLinkStateWatcherStatus(BaseModel): bridge: str = Field(..., description="Bridge interface name.") active: bool = Field(..., description="Whether the watcher thread is currently active.") - poll_interval_seconds: Optional[float] = Field(None, description="Watcher polling interval.") - last_poll_ts: Optional[float] = Field(None, description="Unix timestamp of the last poll.") + event_driven: Optional[bool] = Field(None, description="Whether the watcher is driven by netlink link events.") + last_event_ts: Optional[float] = Field(None, description="Unix timestamp of the last processed event/state evaluation.") last_error: Optional[str] = Field(None, description="Most recent watcher error, if any.") 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.") members: Dict[str, BridgeMemberLinkStateInfo] = Field(default_factory=dict, description="Per-member link-state snapshot.") message: Optional[str] = Field(None, description="Optional informational message.") @@ -437,7 +438,7 @@ def enable_bridge_link_state_watcher( status = bridge_link_state_manager.enable( bridge_name=bridge_name, - poll_interval_seconds=req.poll_interval_seconds, + recovery_holdoff_seconds=req.recovery_holdoff_seconds, ) return _watcher_status_response(status) diff --git a/backend/src/config.py b/backend/src/config.py index 4ada9a1..5cb20ab 100644 --- a/backend/src/config.py +++ b/backend/src/config.py @@ -52,7 +52,6 @@ class BackendSettings: sniffer_buffer_drain_interval_seconds: float sniffer_thread_join_timeout_seconds: float bridge_bpf_build_dir: str - bridge_link_state_poll_interval_seconds: float bridge_link_state_thread_join_timeout_seconds: float bridge_link_state_recovery_holdoff_seconds: float telemetry_process_stop_timeout_seconds: float @@ -89,10 +88,6 @@ def load_settings() -> BackendSettings: sniffer_buffer_drain_interval_seconds=_env_float("BACKEND_SNIFFER_BUFFER_DRAIN_INTERVAL_SECONDS", 5.0), sniffer_thread_join_timeout_seconds=_env_float("BACKEND_SNIFFER_THREAD_JOIN_TIMEOUT_SECONDS", 2.0), bridge_bpf_build_dir=_env_str("BACKEND_BRIDGE_BPF_BUILD_DIR", "/tmp/mitm-bpf"), - bridge_link_state_poll_interval_seconds=_env_float( - "BACKEND_BRIDGE_LINK_STATE_POLL_INTERVAL_SECONDS", - 0.25, - ), bridge_link_state_thread_join_timeout_seconds=_env_float( "BACKEND_BRIDGE_LINK_STATE_THREAD_JOIN_TIMEOUT_SECONDS", 2.0, diff --git a/backend/src/utilities/bridge_link_state_manager.py b/backend/src/utilities/bridge_link_state_manager.py index 437bdc3..0e997ec 100644 --- a/backend/src/utilities/bridge_link_state_manager.py +++ b/backend/src/utilities/bridge_link_state_manager.py @@ -1,8 +1,10 @@ -"""Background watcher that propagates bridge member link failures to sibling ports.""" +"""Event-driven watcher that propagates bridge member link failures to sibling ports.""" from __future__ import annotations import logging +import os +import select import threading import time from dataclasses import dataclass @@ -57,11 +59,14 @@ class MemberLinkState: class BridgeLinkStateWatcher: """Watch one bridge and mirror member failures to the other bridge members.""" - def __init__(self, bridge_name: str, poll_interval_seconds: float) -> None: + def __init__(self, bridge_name: str, recovery_holdoff_seconds: float) -> None: self.bridge_name = bridge_name - self.poll_interval_seconds = poll_interval_seconds + 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, @@ -71,7 +76,7 @@ class BridgeLinkStateWatcher: self._suppressed_members: dict[str, bool] = {} self._settle_deadlines: dict[str, float] = {} self._all_clear_since: Optional[float] = None - self._last_poll_ts: Optional[float] = None + 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] = {} @@ -87,6 +92,7 @@ class BridgeLinkStateWatcher: 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) @@ -98,6 +104,12 @@ class BridgeLinkStateWatcher: 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: @@ -108,34 +120,60 @@ class BridgeLinkStateWatcher: return { "bridge": self.bridge_name, "active": self._running and self._thread.is_alive() and not self._stop_event.is_set(), - "poll_interval_seconds": self.poll_interval_seconds, - "last_poll_ts": self._last_poll_ts, + "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), - "recovery_holdoff_seconds": settings.bridge_link_state_recovery_holdoff_seconds, + "recovery_holdoff_seconds": self.recovery_holdoff_seconds, "members": members, } def _run(self) -> None: - while not self._stop_event.is_set(): - try: - self._poll_once() - with self._lock: - self._last_error = None - except Exception as exc: - logger.exception("Bridge link-state poll failed for bridge=%s", self.bridge_name) - with self._lock: - self._last_error = str(exc) - self._stop_event.wait(self.poll_interval_seconds) + with IPRoute() as ipr: + ipr.bind() + self._evaluate_bridge_state(reason="watcher_started") - def _poll_once(self) -> None: + 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): + self._evaluate_bridge_state(reason="netlink_event") + continue + + self._evaluate_bridge_state(reason="recovery_deadline") + + def _evaluate_bridge_state(self, reason: str) -> 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() with self._lock: - self._last_poll_ts = now + self._last_event_ts = now + self._last_error = None self._member_states = states self._suppressed_members = { iface: restore_up @@ -149,7 +187,8 @@ class BridgeLinkStateWatcher: } if len(states) < 2: - self._all_clear_since = None + with self._lock: + self._all_clear_since = None self._restore_suppressed_members(reason="bridge_has_fewer_than_two_members") return @@ -176,22 +215,44 @@ class BridgeLinkStateWatcher: with self._lock: if self._all_clear_since is None: self._all_clear_since = now - should_restore = now - self._all_clear_since >= settings.bridge_link_state_recovery_holdoff_seconds + should_restore = now - self._all_clear_since >= self.recovery_holdoff_seconds if not should_restore: - remaining = max( - 0.0, - settings.bridge_link_state_recovery_holdoff_seconds - (now - self._all_clear_since), - ) + remaining = max(0.0, self.recovery_holdoff_seconds - (now - self._all_clear_since)) self._last_action = f"waiting {remaining:.2f}s before restoring suppressed members" if should_restore: - self._restore_suppressed_members(reason="all_members_recovered") + self._restore_suppressed_members(reason=reason) return with self._lock: self._all_clear_since = None - self._restore_suppressed_members(reason="all_members_recovered") + self._restore_suppressed_members(reason=reason) + + 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, deadline - time.time()) + 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) -> MemberLinkState: return MemberLinkState( @@ -261,7 +322,7 @@ class BridgeLinkStateWatcher: if target_up: with self._lock: - self._settle_deadlines[ifname] = time.time() + settings.bridge_link_state_recovery_holdoff_seconds + self._settle_deadlines[ifname] = time.time() + self.recovery_holdoff_seconds class BridgeLinkStateManager: @@ -271,16 +332,16 @@ class BridgeLinkStateManager: self._watchers: dict[str, BridgeLinkStateWatcher] = {} self._lock = threading.Lock() - def enable(self, bridge_name: str, poll_interval_seconds: Optional[float] = None) -> dict[str, Any]: + 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_interval = poll_interval_seconds or settings.bridge_link_state_poll_interval_seconds + 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_interval) + watcher = BridgeLinkStateWatcher(bridge_name, normalized_holdoff) watcher.start() with self._lock: diff --git a/frontend/src/pages/Network.tsx b/frontend/src/pages/Network.tsx index 0ba3c4c..d37cdcb 100644 --- a/frontend/src/pages/Network.tsx +++ b/frontend/src/pages/Network.tsx @@ -284,11 +284,11 @@ export default function Network() { function openEnableWatcherModal(bridgeName: string) { setWatcherBridgeName(bridgeName); - watcherForm.setFieldsValue({ poll_interval_seconds: 0.25 }); + watcherForm.setFieldsValue({ recovery_holdoff_seconds: 1.0 }); setWatcherModalVisible(true); } - function handleEnableWatcher(values: { poll_interval_seconds: number }) { + function handleEnableWatcher(values: { recovery_holdoff_seconds: number }) { if (!watcherBridgeName) { return; } @@ -434,8 +434,9 @@ export default function Network() { > {bridge.members.map((member) => member.ifname).join(', ') || '—'} - - {watcher?.poll_interval_seconds ? `${watcher.poll_interval_seconds}s` : '—'} + {watcher?.event_driven ? 'netlink events' : '—'} + + {watcher?.recovery_holdoff_seconds ? `${watcher.recovery_holdoff_seconds}s` : '—'} {watcher?.suppressed_members.length ? watcher.suppressed_members.join(', ') : 'none'} @@ -501,16 +502,21 @@ export default function Network() { }} onOk={() => watcherForm.submit()} > -
+ - When enabled, the backend suppresses sibling bridge ports if one member loses link and restores them after recovery. + The backend reacts to kernel link events and waits for this holdoff before restoring sibling ports after recovery.
diff --git a/frontend/src/types/network.ts b/frontend/src/types/network.ts index febba3a..b4bbee8 100644 --- a/frontend/src/types/network.ts +++ b/frontend/src/types/network.ts @@ -52,7 +52,7 @@ export interface BridgeRemoveRequest { } export interface BridgeLinkStateEnableRequest { - poll_interval_seconds: number; + recovery_holdoff_seconds: number; } export interface BridgeMemberLinkStateInfo { @@ -67,11 +67,12 @@ export interface BridgeMemberLinkStateInfo { export interface BridgeLinkStateWatcherStatus { bridge: string; active: boolean; - poll_interval_seconds?: number | null; - last_poll_ts?: number | null; + event_driven?: boolean | null; + last_event_ts?: number | null; last_error?: string | null; last_action?: string | null; suppressed_members: string[]; + recovery_holdoff_seconds?: number | null; members: Record; message?: string | null; }