test nic state link api
This commit is contained in:
@@ -1,11 +1,15 @@
|
||||
"""Network inspection and bridge management endpoints."""
|
||||
|
||||
from typing import List, Optional
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException
|
||||
from pydantic import BaseModel, Field
|
||||
from pyroute2 import IPRoute, NDB
|
||||
|
||||
from src.config import settings
|
||||
from src.utilities.bridge_link_state_manager import bridge_link_state_manager
|
||||
from src.utilities.interface_bridge_helpers import get_bridge_ports_once
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
ip: IPRoute | None = None
|
||||
@@ -79,6 +83,41 @@ class BridgeRemoveRequest(BaseModel):
|
||||
name: str
|
||||
|
||||
|
||||
class BridgeLinkStateEnableRequest(BaseModel):
|
||||
"""Payload for enabling bridge member link-state propagation."""
|
||||
|
||||
poll_interval_seconds: float = Field(
|
||||
settings.bridge_link_state_poll_interval_seconds,
|
||||
gt=0,
|
||||
description="Polling interval for member link-state checks.",
|
||||
)
|
||||
|
||||
|
||||
class BridgeMemberLinkStateInfo(BaseModel):
|
||||
"""Current link-state snapshot for one bridge member."""
|
||||
|
||||
ifname: str = Field(..., description="Interface name.")
|
||||
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.")
|
||||
link_ready: bool = Field(..., description="Whether the member currently looks usable for forwarding.")
|
||||
suppressed: bool = Field(..., description="Whether the watcher administratively suppressed this member.")
|
||||
|
||||
|
||||
class BridgeLinkStateWatcherStatus(BaseModel):
|
||||
"""Status for one bridge link-state propagation watcher."""
|
||||
|
||||
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.")
|
||||
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.")
|
||||
members: Dict[str, BridgeMemberLinkStateInfo] = Field(default_factory=dict, description="Per-member link-state snapshot.")
|
||||
message: Optional[str] = Field(None, description="Optional informational message.")
|
||||
|
||||
|
||||
def init_network_api() -> None:
|
||||
"""Initialize lazy pyroute2 clients."""
|
||||
global ip, ndb
|
||||
@@ -91,6 +130,7 @@ def init_network_api() -> None:
|
||||
def shutdown_network_api() -> None:
|
||||
"""Close pyroute2 clients if they were initialized."""
|
||||
global ip, ndb
|
||||
bridge_link_state_manager.stop()
|
||||
if ip:
|
||||
ip.close()
|
||||
ip = None
|
||||
@@ -167,6 +207,11 @@ def bridge_exists(name: str, ip_route: IPRoute) -> bool:
|
||||
return bool(ip_route.link_lookup(ifname=name))
|
||||
|
||||
|
||||
def _watcher_status_response(payload: dict[str, Any]) -> BridgeLinkStateWatcherStatus:
|
||||
"""Convert an internal watcher status dictionary to the API response model."""
|
||||
return BridgeLinkStateWatcherStatus.model_validate(payload)
|
||||
|
||||
|
||||
@router.get("/interfaces", response_model=List[InterfaceInfo])
|
||||
def get_interfaces(ip: IPRoute = Depends(get_iproute)) -> List[InterfaceInfo]:
|
||||
"""List host interfaces with addresses and decoded flags."""
|
||||
@@ -346,3 +391,64 @@ def remove_bridge(req: BridgeRemoveRequest, ip: IPRoute = Depends(get_iproute))
|
||||
"status": "ok",
|
||||
"deleted": req.name,
|
||||
}
|
||||
|
||||
|
||||
@router.get("/bridge/link-state-watchers", response_model=List[BridgeLinkStateWatcherStatus])
|
||||
def list_bridge_link_state_watchers() -> List[BridgeLinkStateWatcherStatus]:
|
||||
"""List all bridge member link-state propagation watchers."""
|
||||
return [_watcher_status_response(status) for status in bridge_link_state_manager.list_statuses()]
|
||||
|
||||
|
||||
@router.get("/bridge/{bridge_name}/link-state-watcher", response_model=BridgeLinkStateWatcherStatus)
|
||||
def get_bridge_link_state_watcher(
|
||||
bridge_name: str,
|
||||
ip: IPRoute = Depends(get_iproute),
|
||||
) -> BridgeLinkStateWatcherStatus:
|
||||
"""Return the watcher status for one bridge."""
|
||||
if not bridge_exists(bridge_name, ip):
|
||||
raise HTTPException(status_code=404, detail=f"Bridge {bridge_name} not found")
|
||||
|
||||
status = bridge_link_state_manager.get_status(bridge_name)
|
||||
if status is None:
|
||||
return BridgeLinkStateWatcherStatus(
|
||||
bridge=bridge_name,
|
||||
active=False,
|
||||
message="watcher not enabled",
|
||||
)
|
||||
return _watcher_status_response(status)
|
||||
|
||||
|
||||
@router.post("/bridge/{bridge_name}/link-state-watcher/enable", response_model=BridgeLinkStateWatcherStatus)
|
||||
def enable_bridge_link_state_watcher(
|
||||
bridge_name: str,
|
||||
req: BridgeLinkStateEnableRequest,
|
||||
ip: IPRoute = Depends(get_iproute),
|
||||
) -> BridgeLinkStateWatcherStatus:
|
||||
"""Enable member link-state propagation for a bridge."""
|
||||
if not bridge_exists(bridge_name, ip):
|
||||
raise HTTPException(status_code=404, detail=f"Bridge {bridge_name} not found")
|
||||
|
||||
members = get_bridge_ports_once(bridge_name)
|
||||
if len(members) < 2:
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail=f"Bridge {bridge_name} must have at least two member interfaces",
|
||||
)
|
||||
|
||||
status = bridge_link_state_manager.enable(
|
||||
bridge_name=bridge_name,
|
||||
poll_interval_seconds=req.poll_interval_seconds,
|
||||
)
|
||||
return _watcher_status_response(status)
|
||||
|
||||
|
||||
@router.post("/bridge/{bridge_name}/link-state-watcher/disable", response_model=BridgeLinkStateWatcherStatus)
|
||||
def disable_bridge_link_state_watcher(
|
||||
bridge_name: str,
|
||||
ip: IPRoute = Depends(get_iproute),
|
||||
) -> BridgeLinkStateWatcherStatus:
|
||||
"""Disable member link-state propagation for a bridge."""
|
||||
if not bridge_exists(bridge_name, ip):
|
||||
raise HTTPException(status_code=404, detail=f"Bridge {bridge_name} not found")
|
||||
|
||||
return _watcher_status_response(bridge_link_state_manager.disable(bridge_name))
|
||||
|
||||
@@ -52,6 +52,8 @@ 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
|
||||
telemetry_process_stop_timeout_seconds: float
|
||||
telemetry_reader_join_timeout_seconds: float
|
||||
tshark_enabled: bool
|
||||
@@ -86,6 +88,14 @@ 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,
|
||||
),
|
||||
telemetry_process_stop_timeout_seconds=_env_float("BACKEND_TELEMETRY_PROCESS_STOP_TIMEOUT_SECONDS", 3.0),
|
||||
telemetry_reader_join_timeout_seconds=_env_float("BACKEND_TELEMETRY_READER_JOIN_TIMEOUT_SECONDS", 2.0),
|
||||
tshark_enabled=_env_bool("BACKEND_TSHARK_ENABLED", True),
|
||||
|
||||
289
backend/src/utilities/bridge_link_state_manager.py
Normal file
289
backend/src/utilities/bridge_link_state_manager.py
Normal file
@@ -0,0 +1,289 @@
|
||||
"""Background watcher that propagates bridge member link failures to sibling ports."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
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_operstate,
|
||||
)
|
||||
|
||||
logger = logging.getLogger("bridge_link_state_manager")
|
||||
|
||||
|
||||
@dataclass
|
||||
class MemberLinkState:
|
||||
"""Current sysfs-derived state for one bridge member."""
|
||||
|
||||
ifname: str
|
||||
admin_up: Optional[bool]
|
||||
carrier_up: Optional[bool]
|
||||
operstate: Optional[str]
|
||||
|
||||
@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,
|
||||
"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, poll_interval_seconds: float) -> None:
|
||||
self.bridge_name = bridge_name
|
||||
self.poll_interval_seconds = poll_interval_seconds
|
||||
self._stop_event = threading.Event()
|
||||
self._lock = threading.Lock()
|
||||
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._last_poll_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()
|
||||
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
|
||||
|
||||
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(),
|
||||
"poll_interval_seconds": self.poll_interval_seconds,
|
||||
"last_poll_ts": self._last_poll_ts,
|
||||
"last_error": self._last_error,
|
||||
"last_action": self._last_action,
|
||||
"suppressed_members": sorted(self._suppressed_members),
|
||||
"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)
|
||||
|
||||
def _poll_once(self) -> 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._member_states = states
|
||||
self._suppressed_members = {
|
||||
iface: restore_up
|
||||
for iface, restore_up in self._suppressed_members.items()
|
||||
if iface in states
|
||||
}
|
||||
|
||||
if len(states) < 2:
|
||||
self._restore_suppressed_members(reason="bridge_has_fewer_than_two_members")
|
||||
return
|
||||
|
||||
suppressed_snapshot = set(self._suppressed_members)
|
||||
failing_members = sorted(
|
||||
ifname
|
||||
for ifname, state in states.items()
|
||||
if ifname not in suppressed_snapshot and not state.link_ready
|
||||
)
|
||||
|
||||
if failing_members:
|
||||
self._suppress_other_members(states, failing_members)
|
||||
return
|
||||
|
||||
self._restore_suppressed_members(reason="all_members_recovered")
|
||||
|
||||
def _read_member_state(self, ifname: str) -> MemberLinkState:
|
||||
return MemberLinkState(
|
||||
ifname=ifname,
|
||||
admin_up=read_interface_admin_up(ifname),
|
||||
carrier_up=read_interface_carrier(ifname),
|
||||
operstate=read_interface_operstate(ifname),
|
||||
)
|
||||
|
||||
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._last_action = action
|
||||
|
||||
def _restore_suppressed_members(self, reason: str) -> None:
|
||||
with self._lock:
|
||||
suppressed = dict(self._suppressed_members)
|
||||
|
||||
if not suppressed:
|
||||
with self._lock:
|
||||
self._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._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)
|
||||
|
||||
|
||||
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, poll_interval_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
|
||||
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.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()
|
||||
@@ -1,4 +1,4 @@
|
||||
from typing import List, Dict, Optional
|
||||
from typing import List, Optional
|
||||
import logging
|
||||
import os
|
||||
|
||||
@@ -30,6 +30,49 @@ def check_interface_up(iface: str) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _read_sysfs_text(path: str) -> Optional[str]:
|
||||
"""Read one sysfs file and return stripped text, or `None` if unavailable."""
|
||||
try:
|
||||
with open(path, "r") as f:
|
||||
return f.read().strip()
|
||||
except FileNotFoundError:
|
||||
return None
|
||||
except Exception:
|
||||
logger.exception("Error reading sysfs path %s", path)
|
||||
return None
|
||||
|
||||
|
||||
def read_interface_operstate(iface: str) -> Optional[str]:
|
||||
"""Return the kernel operstate string for an interface."""
|
||||
return _read_sysfs_text(f"/sys/class/net/{iface}/operstate")
|
||||
|
||||
|
||||
def read_interface_carrier(iface: str) -> Optional[bool]:
|
||||
"""Return the carrier state for an interface if available."""
|
||||
value = _read_sysfs_text(f"/sys/class/net/{iface}/carrier")
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
try:
|
||||
return int(value, 10) == 1
|
||||
except ValueError:
|
||||
logger.warning("Unexpected carrier value for %s: %r", iface, value)
|
||||
return None
|
||||
|
||||
|
||||
def read_interface_admin_up(iface: str) -> Optional[bool]:
|
||||
"""Return whether the interface has the IFF_UP flag set."""
|
||||
value = _read_sysfs_text(f"/sys/class/net/{iface}/flags")
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
try:
|
||||
return bool(int(value, 0) & 0x1)
|
||||
except ValueError:
|
||||
logger.warning("Unexpected flags 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