websockets and network interface datat added
All checks were successful
Build and Deploy MITM Webserver / traffic_target (push) Successful in 0s
Build and Deploy MITM Webserver / build (push) Successful in 9s

This commit is contained in:
2026-03-09 16:45:01 +01:00
parent e919853d07
commit 40d8b15417
8 changed files with 352 additions and 107 deletions

View File

@@ -1,16 +1,21 @@
"""Network inspection and bridge management endpoints."""
import asyncio
import logging
from typing import Any, Dict, List, Optional
from fastapi import APIRouter, Depends, HTTPException
from fastapi import APIRouter, Depends, HTTPException, WebSocket, WebSocketDisconnect
from pydantic import BaseModel, Field
from pyroute2 import IPRoute, NDB
from starlette.websockets import WebSocketState
import src.shared_objects as shared
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
from src.utilities.interface_bridge_helpers import get_bridge_ports_once, read_interface_ethernet_profile
router = APIRouter()
logger = logging.getLogger("network_router")
ip: IPRoute | None = None
ndb: NDB | None = None
@@ -33,6 +38,10 @@ class InterfaceInfo(BaseModel):
mac: Optional[str] = Field(None, description="MAC address.")
mtu: int = Field(..., description="Maximum transmission unit.")
flags: List[str] = Field(..., description="Decoded interface flags.")
ethernet_profile: Optional[Dict[str, Any]] = Field(
None,
description="Current speed/duplex/autoneg profile when available.",
)
addresses: List[InterfaceAddress] = Field(..., description="Assigned IP addresses.")
@@ -57,6 +66,10 @@ class BridgeInterfaceInfo(BaseModel):
ifname: str = Field(..., description="Interface name.")
state: Optional[str] = Field(None, description="Operational state.")
mtu: Optional[int] = Field(None, description="Interface MTU.")
ethernet_profile: Optional[Dict[str, Any]] = Field(
None,
description="Current speed/duplex/autoneg profile when available.",
)
class BridgeInfo(BaseModel):
@@ -124,10 +137,23 @@ class BridgeLinkStateWatcherStatus(BaseModel):
None,
description="Low-rate fallback recheck interval while the bridge is degraded.",
)
failure_holdoff_seconds: Optional[float] = Field(
None,
description="Transient failure debounce before suppressing siblings.",
)
members: Dict[str, BridgeMemberLinkStateInfo] = Field(default_factory=dict, description="Per-member link-state snapshot.")
message: Optional[str] = Field(None, description="Optional informational message.")
class FullStateResponse(BaseModel):
"""Interfaces, routes, bridges, and active watcher state."""
interfaces: List[InterfaceInfo] = Field(default_factory=list)
routes: List[RouteInfo] = Field(default_factory=list)
bridges: List[BridgeInfo] = Field(default_factory=list)
watchers: List[BridgeLinkStateWatcherStatus] = Field(default_factory=list)
def init_network_api() -> None:
"""Initialize lazy pyroute2 clients."""
global ip, ndb
@@ -222,6 +248,40 @@ def _watcher_status_response(payload: dict[str, Any]) -> BridgeLinkStateWatcherS
return BridgeLinkStateWatcherStatus.model_validate(payload)
def _build_full_state_response(ip_route: Optional[IPRoute] = None) -> FullStateResponse:
"""Build the current full network snapshot used by HTTP and websocket consumers."""
ip_instance = ip_route or get_iproute()
return FullStateResponse(
interfaces=get_interfaces(ip_instance),
routes=get_routes(ip_instance),
bridges=get_bridges(),
watchers=list_bridge_link_state_watchers(),
)
def build_full_state_payload(ip_route: Optional[IPRoute] = None) -> dict[str, Any]:
"""Return the current network snapshot as a JSON-safe dictionary."""
return _build_full_state_response(ip_route).model_dump(mode="json")
def publish_network_state_update(reason: str) -> None:
"""Publish the latest network snapshot to websocket subscribers."""
broadcaster = getattr(shared, "network_broadcaster", None)
if broadcaster is None:
return
try:
broadcaster.sync_publish(
{
"type": "network_state",
"reason": reason,
"snapshot": build_full_state_payload(),
}
)
except Exception:
logger.exception("Failed to publish network state update")
@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."""
@@ -247,6 +307,7 @@ def get_interfaces(ip: IPRoute = Depends(get_iproute)) -> List[InterfaceInfo]:
mac=attrs.get("IFLA_ADDRESS"),
mtu=attrs.get("IFLA_MTU"),
flags=parse_flags(link.get("flags", 0)),
ethernet_profile=read_interface_ethernet_profile(attrs.get("IFLA_IFNAME")),
addresses=parse_addresses(addrs),
)
)
@@ -310,6 +371,7 @@ def get_raw_links(ip: IPRoute = Depends(get_iproute)) -> List[InterfaceInfo]:
mac=attrs.get("IFLA_ADDRESS"),
mtu=attrs.get("IFLA_MTU", 0),
flags=[],
ethernet_profile=read_interface_ethernet_profile(attrs.get("IFLA_IFNAME", "unknown")),
addresses=parse_addresses(addrs),
)
)
@@ -336,6 +398,7 @@ def get_bridges() -> List[BridgeInfo]:
ifname=iface.ifname,
state=getattr(iface, "operstate", None),
mtu=getattr(iface, "mtu", None),
ethernet_profile=read_interface_ethernet_profile(iface.ifname),
)
)
@@ -353,14 +416,10 @@ def get_bridges() -> List[BridgeInfo]:
return bridges_list
@router.get("/full-state")
def full_state(ip: IPRoute = Depends(get_iproute)) -> dict:
@router.get("/full-state", response_model=FullStateResponse)
def full_state(ip: IPRoute = Depends(get_iproute)) -> FullStateResponse:
"""Return interfaces, routes, and bridges in one response."""
return {
"interfaces": get_interfaces(ip),
"routes": get_routes(ip),
"bridges": get_bridges(),
}
return _build_full_state_response(ip)
@router.post("/bridge/create")
@@ -380,6 +439,8 @@ def create_bridge(req: BridgeCreateRequest, ip: IPRoute = Depends(get_iproute))
ip.link("set", index=idx, state="up")
ip.link("set", index=idx, master=br_idx)
publish_network_state_update("bridge_created")
return {
"status": "ok",
"bridge": req.name,
@@ -397,6 +458,8 @@ def remove_bridge(req: BridgeRemoveRequest, ip: IPRoute = Depends(get_iproute))
ip.link("set", index=br_idx, state="down")
ip.link("del", index=br_idx)
publish_network_state_update("bridge_removed")
return {
"status": "ok",
"deleted": req.name,
@@ -449,6 +512,7 @@ def enable_bridge_link_state_watcher(
bridge_name=bridge_name,
recovery_holdoff_seconds=req.recovery_holdoff_seconds,
)
publish_network_state_update("bridge_watcher_enabled")
return _watcher_status_response(status)
@@ -461,4 +525,70 @@ def disable_bridge_link_state_watcher(
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))
status = _watcher_status_response(bridge_link_state_manager.disable(bridge_name))
publish_network_state_update("bridge_watcher_disabled")
return status
@router.websocket("/ws/state")
async def websocket_network_state(ws: WebSocket) -> None:
"""Stream full network snapshots to websocket clients whenever the backend publishes updates."""
await ws.accept()
broadcaster = getattr(shared, "network_broadcaster", None)
if broadcaster is None:
await ws.send_json({"error": "network broadcaster not available"})
await ws.close()
return
queue: Optional[asyncio.Queue] = None
try:
await ws.send_json({"type": "network_state", "reason": "initial", "snapshot": build_full_state_payload()})
queue = await broadcaster.subscribe()
while True:
queue_task = asyncio.create_task(queue.get())
receive_task = asyncio.create_task(ws.receive())
done, pending = await asyncio.wait({queue_task, receive_task}, return_when=asyncio.FIRST_COMPLETED)
for task in pending:
task.cancel()
if pending:
await asyncio.gather(*pending, return_exceptions=True)
if receive_task in done:
try:
inbound = receive_task.result()
except WebSocketDisconnect:
break
except Exception:
break
if inbound.get("type") == "websocket.disconnect":
break
if queue_task not in done:
if ws.client_state is not WebSocketState.CONNECTED:
break
continue
message = queue_task.result()
if isinstance(message, dict) and message.get("type") == "__broadcaster_shutdown__":
break
try:
await ws.send_json(message)
except Exception:
break
except WebSocketDisconnect:
pass
finally:
if queue is not None:
try:
await broadcaster.unsubscribe(queue)
except Exception:
logger.exception("Failed to unsubscribe network websocket queue")
try:
await ws.close()
except Exception:
pass

View File

@@ -61,6 +61,11 @@ async def on_startup() -> None:
except Exception:
logging.exception("Failed to create/attach broadcaster")
try:
shared_objects.network_broadcaster = PacketBroadcaster(loop, queue_maxsize=settings.broadcaster_queue_maxsize)
except Exception:
logging.exception("Failed to create network broadcaster")
try:
from src import network_sniffer as sniffer
@@ -114,6 +119,12 @@ async def shutdown_event() -> None:
except Exception:
logging.exception("Failed to close packet broadcaster during shutdown")
try:
if shared_objects.network_broadcaster is not None:
await shared_objects.network_broadcaster.close()
except Exception:
logging.exception("Failed to close network broadcaster during shutdown")
try:
web_db = getattr(shared_objects, "db", None)
if web_db is not None:
@@ -123,6 +134,7 @@ async def shutdown_event() -> None:
shared_objects.db = None
shared_objects.broadcaster = None
shared_objects.network_broadcaster = None
shared_objects.web_loop = None

View File

@@ -6,3 +6,4 @@ from typing import Any, Optional
db: Any = None
web_loop: Optional[asyncio.AbstractEventLoop] = None
broadcaster: Any = None
network_broadcaster: Any = None

View File

@@ -19,6 +19,7 @@ from src.utilities.interface_bridge_helpers import (
get_bridge_ports_once,
read_interface_admin_up,
read_interface_carrier,
read_interface_ethernet_profile,
read_interface_mtu,
read_interface_operstate,
)
@@ -255,9 +256,7 @@ class BridgeLinkStateWatcher:
matured_failing_members = self._track_failing_members(failing_members, now)
if not matured_failing_members:
self._set_last_action(
f"waiting {settings.bridge_link_state_failure_holdoff_seconds:.2f}s before suppressing due to failures={failing_members}"
)
self._set_last_action(f"waiting before suppressing transient failures on {failing_members}")
return
with self._lock:
@@ -270,14 +269,17 @@ class BridgeLinkStateWatcher:
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:
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"
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)
@@ -413,69 +415,36 @@ class BridgeLinkStateWatcher:
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:
profile = read_interface_ethernet_profile(ifname)
if profile is None:
return None
return EthernetProfile(
speed_mbps=speed_mbps,
duplex=duplex,
autoneg=autoneg,
speed_mbps=profile.get("speed_mbps"),
duplex=profile.get("duplex"),
autoneg=profile.get("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] = []
changes: 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)
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 changed_members:
self._set_last_action(f"synchronized configuration from {source_ifname} to {changed_members}")
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,
@@ -483,10 +452,10 @@ class BridgeLinkStateWatcher:
source: MemberLinkState,
target_ifname: str,
target: MemberLinkState,
) -> bool:
) -> 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 False
return None
with IPRoute() as ipr:
indices = ipr.link_lookup(ifname=target_ifname)
@@ -502,7 +471,7 @@ class BridgeLinkStateWatcher:
self.bridge_name,
source.mtu,
)
return True
return f"{target_ifname} mtu={source.mtu}"
def _sync_member_ethernet_profile(
self,
@@ -510,12 +479,12 @@ class BridgeLinkStateWatcher:
source: MemberLinkState,
target_ifname: str,
target: MemberLinkState,
) -> bool:
) -> 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 False
return None
cmd = [_ETHTOOL_BIN, "-s", target_ifname]
if source_profile.autoneg is True:
@@ -532,7 +501,7 @@ class BridgeLinkStateWatcher:
]
)
else:
return False
return None
try:
subprocess.run(cmd, capture_output=True, text=True, check=True)
@@ -544,7 +513,7 @@ class BridgeLinkStateWatcher:
self.bridge_name,
exc,
)
return False
return None
self._mark_managed_change(target_ifname)
logger.info(
@@ -554,7 +523,11 @@ class BridgeLinkStateWatcher:
self.bridge_name,
source_profile,
)
return True
return (
f"{target_ifname} link="
f"{source_profile.speed_mbps or 'unknown'}Mb/{source_profile.duplex or 'unknown'}/"
f"{'autoneg-on' if source_profile.autoneg else 'autoneg-off'}"
)
def _suppress_other_members(self, states: dict[str, MemberLinkState], failing_members: list[str]) -> None:
desired_suppressed = set(states) - set(failing_members)
@@ -579,7 +552,7 @@ class BridgeLinkStateWatcher:
if changed_members
else f"holding suppressed members because failing members={failing_members}"
)
self._last_action = action
self._set_last_action(action)
def _restore_suppressed_members(self, reason: str) -> None:
with self._lock:
@@ -600,7 +573,7 @@ class BridgeLinkStateWatcher:
with self._lock:
self._suppressed_members = {}
self._all_clear_since = None
self._last_action = f"restored {restored_members} ({reason})"
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):
@@ -625,9 +598,23 @@ class BridgeLinkStateWatcher:
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."""

View File

@@ -1,9 +1,11 @@
from typing import List, Optional
from typing import Any, Dict, List, Optional
import logging
import os
import subprocess
# ---- Logging ----------------------------------------------------------
logger = logging.getLogger("af_packet_sniffer")
_ETHTOOL_BIN = "/usr/sbin/ethtool" if os.path.exists("/usr/sbin/ethtool") else "ethtool"
# -------------------------
# Interface / bridge helpers
@@ -86,6 +88,55 @@ def read_interface_mtu(iface: str) -> Optional[int]:
return None
def read_interface_ethernet_profile(iface: str) -> Optional[Dict[str, Any]]:
"""Return the current speed/duplex/autoneg profile when ethtool supports it."""
try:
result = subprocess.run(
[_ETHTOOL_BIN, iface],
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_mbps: Optional[int] = None
speed_value = values.get("Speed")
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
duplex_value = values.get("Duplex")
if duplex_value and duplex_value.lower() in {"full", "half"}:
duplex = duplex_value.lower()
autoneg: Optional[bool] = None
autoneg_value = values.get("Auto-negotiation")
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 {
"speed_mbps": speed_mbps,
"duplex": duplex,
"autoneg": autoneg,
}
def _read_bridge_ports_from_sysfs(bridge: str) -> List[str]:
"""
Read bridge member interfaces from sysfs. Internal helper that always reads.