diff --git a/netx_api/config.py b/netx_api/config.py index ee00f2e..dac3a2e 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -77,6 +77,9 @@ class Settings(BaseSettings): config_sync_scheduler_tick_sec: int = 60 # After process start / unexpected restart, wait before any scheduled sync. config_sync_startup_grace_sec: int = 3600 + # Port traffic monitoring (CLI rate bit/s samples) + port_traffic_scheduler_enabled: bool = True + port_traffic_scheduler_tick_sec: int = 15 # Managed NE exec: max CLI commands per request (lab can raise; hard-capped in ne_exec). ne_exec_max_commands: int = 5 # WebCRT interactive terminal sessions diff --git a/netx_api/main.py b/netx_api/main.py index 5659920..5c379fb 100644 --- a/netx_api/main.py +++ b/netx_api/main.py @@ -27,6 +27,7 @@ from .db import Base, SessionLocal, engine, get_db from .collection_router import router as collection_router from .cli_router import router as cli_router from .config_sync_router import router as config_sync_router +from .port_traffic_router import router as port_traffic_router from .managed_ne_router import router as managed_ne_router from .webcrt_router import router as webcrt_router from .topology_router import router as topology_router @@ -135,6 +136,7 @@ app.include_router(managed_ne_router) app.include_router(cli_router) app.include_router(collection_router) app.include_router(config_sync_router) +app.include_router(port_traffic_router) app.include_router(webcrt_router) app.include_router(topology_router) parser_cfg = load_parser_config() @@ -836,11 +838,15 @@ def on_startup() -> None: _schedule_log.info("startup: resumed %s pending ne collection runs", resumed) from .config_sync_recovery import recover_config_sync_on_startup from .config_sync_service import ensure_policy + from .port_traffic_recovery import recover_port_traffic_on_startup ensure_policy(db) cfg_resumed = recover_config_sync_on_startup(db) if cfg_resumed: _schedule_log.info("startup: resumed %s config_sync task(s) from interrupted cycle", cfg_resumed) + pt_cleared = recover_port_traffic_on_startup(db) + if pt_cleared: + _schedule_log.info("startup: cleared %s port_traffic stuck collect_running flag(s)", pt_cleared) except Exception: _schedule_log.exception("startup: ne collection / config_sync recovery failed") finally: @@ -851,6 +857,12 @@ def on_startup() -> None: start_config_sync_scheduler() except Exception: _schedule_log.exception("startup: config_sync scheduler init failed") + try: + from .port_traffic_scheduler import start_port_traffic_scheduler + + start_port_traffic_scheduler() + except Exception: + _schedule_log.exception("startup: port_traffic scheduler init failed") # Best-effort schema evolution for new columns (no migrations framework). # Safe for Postgres (IF NOT EXISTS); ignored on failure. try: diff --git a/netx_api/models.py b/netx_api/models.py index 7d01fd8..d8527a9 100644 --- a/netx_api/models.py +++ b/netx_api/models.py @@ -3,7 +3,7 @@ from __future__ import annotations from datetime import datetime from uuid import uuid4 -from sqlalchemy import Boolean, DateTime, Float, ForeignKey, Integer, LargeBinary, String, Text, UniqueConstraint +from sqlalchemy import BigInteger, Boolean, DateTime, Float, ForeignKey, Integer, LargeBinary, String, Text, UniqueConstraint from sqlalchemy.dialects.postgresql import JSONB from sqlalchemy.orm import Mapped, mapped_column, relationship from sqlalchemy.types import JSON @@ -588,3 +588,65 @@ class NeConfigHistory(Base): collected_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) cycle_id: Mapped[str] = mapped_column(String(64), default="") task_id: Mapped[str] = mapped_column(String(64), default="") + + +class PortTrafficTask(Base): + """Port traffic monitoring job definition.""" + + __tablename__ = "port_traffic_task" + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + title: Mapped[str] = mapped_column(String(256), default="") + status: Mapped[str] = mapped_column(String(32), default="draft", index=True) # draft|running|paused|stopped + interval_sec: Mapped[int] = mapped_column(Integer, default=60) + retention_days: Mapped[int] = mapped_column(Integer, default=7) + concurrency: Mapped[int] = mapped_column(Integer, default=5) + collect_running: Mapped[bool] = mapped_column(Boolean, default=False) + last_collect_started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + last_collect_ended_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + last_error: Mapped[str] = mapped_column(String(1024), default="") + created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + updated_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + + +class PortTrafficTarget(Base): + """Monitored interface under a port traffic task.""" + + __tablename__ = "port_traffic_target" + __table_args__ = ( + UniqueConstraint("task_id", "source", "target_id", "ifname", name="uq_port_traffic_target_if"), + ) + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + task_id: Mapped[str] = mapped_column(String(64), index=True) + source: Mapped[str] = mapped_column(String(32), default="managed", index=True) + target_id: Mapped[str] = mapped_column(String(128), index=True) + ne_name: Mapped[str] = mapped_column(String(256), default="") + ne_ip: Mapped[str] = mapped_column(String(128), default="") + vendor: Mapped[str] = mapped_column(String(64), default="") + ifname: Mapped[str] = mapped_column(String(128), default="") + if_description: Mapped[str] = mapped_column(String(512), default="") + bw_bps: Mapped[int] = mapped_column(BigInteger, default=0) + status: Mapped[str] = mapped_column(String(32), default="active", index=True) # active|disabled + last_error: Mapped[str] = mapped_column(String(1024), default="") + last_sample_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True) + created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow) + + +class PortTrafficSample(Base): + """Time-series sample for a monitored interface.""" + + __tablename__ = "port_traffic_sample" + __table_args__ = (UniqueConstraint("target_row_id", "ts", name="uq_port_traffic_sample_ts"),) + + id: Mapped[str] = mapped_column(String(64), primary_key=True, default=lambda: uuid4().hex) + target_row_id: Mapped[str] = mapped_column(String(64), index=True) + ts: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow, index=True) + in_bps: Mapped[float] = mapped_column(Float, default=0.0) + out_bps: Mapped[float] = mapped_column(Float, default=0.0) + in_util_pct: Mapped[float] = mapped_column(Float, default=0.0) + out_util_pct: Mapped[float] = mapped_column(Float, default=0.0) + bw_bps: Mapped[int] = mapped_column(BigInteger, default=0) + rate_period_sec: Mapped[int] = mapped_column(Integer, default=0) + raw_ok: Mapped[bool] = mapped_column(Boolean, default=True) + message: Mapped[str] = mapped_column(String(512), default="") diff --git a/netx_api/port_traffic_commands.py b/netx_api/port_traffic_commands.py new file mode 100644 index 0000000..1f6b681 --- /dev/null +++ b/netx_api/port_traffic_commands.py @@ -0,0 +1,29 @@ +"""Vendor → port traffic CLI command matrix (ZTE first).""" + +from __future__ import annotations + +from dataclasses import dataclass + +from .config_sync_commands import normalize_vendor_key + + +@dataclass(frozen=True) +class PortTrafficCommands: + brief: str + detail_template: str # format with ifname= + vendor_key: str = "other" + + +def commands_for_vendor(vendor: str, device_type: str = "") -> PortTrafficCommands | None: + key = normalize_vendor_key(vendor, device_type) + if key == "zte": + return PortTrafficCommands( + brief="show interface brief", + detail_template="show interface {ifname}", + vendor_key=key, + ) + return None + + +def detail_command(cmds: PortTrafficCommands, ifname: str) -> str: + return cmds.detail_template.format(ifname=ifname) diff --git a/netx_api/port_traffic_parsers.py b/netx_api/port_traffic_parsers.py new file mode 100644 index 0000000..ec36467 --- /dev/null +++ b/netx_api/port_traffic_parsers.py @@ -0,0 +1,237 @@ +"""Parsers for ZTE show interface brief / detail (rate bit/s).""" + +from __future__ import annotations + +import re +from dataclasses import dataclass +from typing import Any + + +_BW_UNIT = { + "k": 1_000, + "m": 1_000_000, + "g": 1_000_000_000, + "t": 1_000_000_000_000, +} + +_RE_BW_COMPACT = re.compile(r"^(\d+(?:\.\d+)?)\s*([kKmMgGtT])(?:bit)?s?$", re.I) +_RE_BW_DETAIL = re.compile( + r"\bBW\s+(\d+(?:\.\d+)?)\s*([kKmMgGtT])?\s*(?:G?bit|bit)/s\b", + re.I, +) +_RE_RATE_PERIOD = re.compile(r"Rate\s+period\s*:\s*(\d+)\s*s", re.I) +_RE_INPUT_BPS = re.compile(r"^\s*Input\s*:\s*([\d.]+)\s*bit/s", re.I | re.M) +_RE_OUTPUT_BPS = re.compile(r"^\s*Output\s*:\s*([\d.]+)\s*bit/s", re.I | re.M) +_RE_UTIL = re.compile( + r"Intf\s+utilization\s*:\s*input\s*([\d.]+)%\s*output\s*([\d.]+)%", + re.I, +) +_RE_IF_UP = re.compile(r"^(\S+)\s+is\s+(up|down)\b", re.I | re.M) +_RE_DESC = re.compile(r"^\s*Description:\s*(.+?)\s*$", re.I | re.M) + +# Fixed-width columns from ZTE `show interface brief` header. +_COL_IF = (0, 24) +_COL_ATTR = (24, 35) +_COL_MODE = (35, 48) +_COL_BW = (48, 54) +_COL_ADMIN = (54, 60) +_COL_PHY = (60, 66) +_COL_PROT = (66, 72) +_COL_DESC = (72, None) + + +@dataclass(frozen=True) +class BriefPort: + ifname: str + attribute: str = "" + mode: str = "" + bw_raw: str = "" + bw_bps: int = 0 + admin: str = "" + phy: str = "" + prot: str = "" + description: str = "" + + +@dataclass(frozen=True) +class DetailRates: + ifname: str = "" + admin_oper: str = "" + description: str = "" + bw_bps: int = 0 + rate_period_sec: int = 0 + in_bps: float = 0.0 + out_bps: float = 0.0 + in_util_pct: float = 0.0 + out_util_pct: float = 0.0 + + +def parse_bw_to_bps(raw: str) -> int: + text = (raw or "").strip() + if not text or text.upper() == "N/A": + return 0 + m = _RE_BW_COMPACT.match(text.replace(" ", "")) + if m: + value = float(m.group(1)) + mult = _BW_UNIT[m.group(2).lower()] + return int(value * mult) + m2 = _RE_BW_DETAIL.search(text) + if m2: + value = float(m2.group(1)) + unit = (m2.group(2) or "g").lower() + # "BW 1 Gbit/s" — unit letter may be in "Gbit" when group2 empty after "1 " + if m2.group(2) is None and "gbit" in text.lower(): + unit = "g" + elif m2.group(2) is None and "mbit" in text.lower(): + unit = "m" + elif m2.group(2) is None and "kbit" in text.lower(): + unit = "k" + mult = _BW_UNIT.get(unit, 1_000_000_000) + return int(value * mult) + # Plain "BW 1000000000" unlikely; try digits only + digits = re.sub(r"[^\d.]", "", text) + if digits: + try: + return int(float(digits)) + except ValueError: + return 0 + return 0 + + +def _slice(line: str, start: int, end: int | None) -> str: + if end is None: + return line[start:].rstrip() if len(line) > start else "" + if len(line) <= start: + return "" + return line[start:end].strip() + + +def parse_zte_interface_brief(text: str) -> list[BriefPort]: + """Parse ZTE `show interface brief` into port rows.""" + lines = (text or "").replace("\r\n", "\n").replace("\r", "\n").split("\n") + out: list[BriefPort] = [] + started = False + for raw in lines: + line = raw.rstrip() + if not line.strip(): + continue + if line.lstrip().startswith("Interface") and "Admin" in line and "Description" in line: + started = True + continue + if not started: + continue + if " is " in line and "ifindex" in line.lower(): + break + ifname = _slice(line, *_COL_IF).split()[0] if _slice(line, *_COL_IF) else "" + if not ifname or ifname.lower() == "interface": + continue + bw_raw = _slice(line, *_COL_BW) + out.append( + BriefPort( + ifname=ifname, + attribute=_slice(line, *_COL_ATTR), + mode=_slice(line, *_COL_MODE), + bw_raw=bw_raw, + bw_bps=parse_bw_to_bps(bw_raw), + admin=_slice(line, *_COL_ADMIN).lower(), + phy=_slice(line, *_COL_PHY).lower(), + prot=_slice(line, *_COL_PROT).lower(), + description=_slice(line, *_COL_DESC).strip(), + ) + ) + return out + + +def parse_zte_interface_detail(text: str) -> DetailRates: + """Parse ZTE `show interface {ifname}` rate / util / BW.""" + blob = text or "" + ifname = "" + admin_oper = "" + m_if = _RE_IF_UP.search(blob) + if m_if: + ifname = m_if.group(1) + admin_oper = m_if.group(2).lower() + desc = "" + m_desc = _RE_DESC.search(blob) + if m_desc: + desc = m_desc.group(1).strip() + + bw_bps = 0 + m_bw = _RE_BW_DETAIL.search(blob) + if m_bw: + bw_bps = parse_bw_to_bps(m_bw.group(0)) + else: + # Fallback line scan + for line in blob.splitlines(): + if re.search(r"\bBW\b", line, re.I): + bw_bps = parse_bw_to_bps(line) + if bw_bps: + break + + rate_period = 0 + m_rp = _RE_RATE_PERIOD.search(blob) + if m_rp: + rate_period = int(m_rp.group(1)) + + # Prefer Rate period block: first Input/Output after "Rate period" + in_bps = 0.0 + out_bps = 0.0 + rp_idx = blob.lower().find("rate period") + rate_blob = blob[rp_idx:] if rp_idx >= 0 else blob + # Stop before Peak rate to avoid peak Input/Output + peak_idx = rate_blob.lower().find("peak rate") + if peak_idx >= 0: + rate_blob = rate_blob[:peak_idx] + m_in = _RE_INPUT_BPS.search(rate_blob) + m_out = _RE_OUTPUT_BPS.search(rate_blob) + if m_in: + in_bps = float(m_in.group(1)) + if m_out: + out_bps = float(m_out.group(1)) + + in_util = 0.0 + out_util = 0.0 + m_util = _RE_UTIL.search(blob) + if m_util: + in_util = float(m_util.group(1)) + out_util = float(m_util.group(2)) + + return DetailRates( + ifname=ifname, + admin_oper=admin_oper, + description=desc, + bw_bps=bw_bps, + rate_period_sec=rate_period, + in_bps=in_bps, + out_bps=out_bps, + in_util_pct=in_util, + out_util_pct=out_util, + ) + + +def brief_port_to_dict(row: BriefPort) -> dict[str, Any]: + return { + "ifname": row.ifname, + "attribute": row.attribute, + "mode": row.mode, + "bw_raw": row.bw_raw, + "bw_bps": row.bw_bps, + "admin": row.admin, + "phy": row.phy, + "prot": row.prot, + "description": row.description, + } + + +def detail_to_dict(row: DetailRates) -> dict[str, Any]: + return { + "ifname": row.ifname, + "admin_oper": row.admin_oper, + "description": row.description, + "bw_bps": row.bw_bps, + "rate_period_sec": row.rate_period_sec, + "in_bps": row.in_bps, + "out_bps": row.out_bps, + "in_util_pct": row.in_util_pct, + "out_util_pct": row.out_util_pct, + } diff --git a/netx_api/port_traffic_recovery.py b/netx_api/port_traffic_recovery.py new file mode 100644 index 0000000..8f68020 --- /dev/null +++ b/netx_api/port_traffic_recovery.py @@ -0,0 +1,31 @@ +"""Startup recovery for interrupted port traffic collect rounds.""" + +from __future__ import annotations + +import logging +from datetime import datetime + +from sqlalchemy.orm import Session + +from .models import PortTrafficTask + +_log = logging.getLogger("netx.port_traffic.recovery") + + +def recover_port_traffic_on_startup(db: Session) -> int: + """Clear stuck collect_running flags so scheduler can resume.""" + now = datetime.utcnow() + rows = db.query(PortTrafficTask).filter(PortTrafficTask.collect_running.is_(True)).all() + n = 0 + for task in rows: + task.collect_running = False + if not task.last_collect_ended_at: + task.last_collect_ended_at = now + if not task.last_error: + task.last_error = "requeued_after_restart" + task.updated_at = now + n += 1 + if n: + db.commit() + _log.info("port_traffic recovery cleared collect_running on %s task(s)", n) + return n diff --git a/netx_api/port_traffic_router.py b/netx_api/port_traffic_router.py new file mode 100644 index 0000000..480b30b --- /dev/null +++ b/netx_api/port_traffic_router.py @@ -0,0 +1,211 @@ +"""HTTP API for port traffic monitoring.""" + +from __future__ import annotations + +from datetime import datetime + +from fastapi import APIRouter, Depends, Query, Request +from sqlalchemy.orm import Session + +from .auth_service import write_audit +from .db import get_db +from .port_traffic_schemas import ( + DiscoverPortsRequest, + PortTrafficTaskCreate, + PortTrafficTaskUpdate, + PortTrafficTargetsPut, +) +from .port_traffic_service import ( + create_task, + dashboard, + delete_task, + discover_ports, + get_samples, + get_task, + list_targets, + list_tasks, + put_targets, + set_task_status, + update_task, +) + +router = APIRouter(prefix="/v1/port-traffic", tags=["port-traffic"]) + + +def _actor(request: Request) -> tuple[str, str]: + user = getattr(request.state, "auth_user", None) + if not user: + return "", "" + return str(getattr(user, "id", "") or ""), str(getattr(user, "username", "") or "") + + +@router.get("/dashboard") +def api_dashboard(db: Session = Depends(get_db)): + return dashboard(db).model_dump() + + +@router.get("/tasks") +def api_list_tasks( + page: int = Query(default=1, ge=1), + page_size: int = Query(default=20, ge=1, le=100), + db: Session = Depends(get_db), +): + return list_tasks(db, page=page, page_size=page_size) + + +@router.post("/tasks") +def api_create_task( + body: PortTrafficTaskCreate, + request: Request, + db: Session = Depends(get_db), +): + out = create_task(db, body) + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.create", + actor_user_id=uid, + actor_username=uname, + method="POST", + path="/v1/port-traffic/tasks", + status_code=200, + detail={"id": out.id, "title": out.title, "start_now": body.start_now}, + ) + return out.model_dump() + + +@router.get("/tasks/{task_id}") +def api_get_task(task_id: str, db: Session = Depends(get_db)): + return get_task(db, task_id).model_dump() + + +@router.patch("/tasks/{task_id}") +def api_patch_task( + task_id: str, + body: PortTrafficTaskUpdate, + request: Request, + db: Session = Depends(get_db), +): + out = update_task(db, task_id, body) + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.update", + actor_user_id=uid, + actor_username=uname, + method="PATCH", + path=f"/v1/port-traffic/tasks/{task_id}", + status_code=200, + detail=body.model_dump(exclude_unset=True), + ) + return out.model_dump() + + +@router.delete("/tasks/{task_id}") +def api_delete_task(task_id: str, request: Request, db: Session = Depends(get_db)): + out = delete_task(db, task_id) + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.delete", + actor_user_id=uid, + actor_username=uname, + method="DELETE", + path=f"/v1/port-traffic/tasks/{task_id}", + status_code=200, + detail={"id": task_id}, + ) + return out + + +@router.post("/tasks/{task_id}/start") +def api_start_task(task_id: str, request: Request, db: Session = Depends(get_db)): + out = set_task_status(db, task_id, "running") + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.start", + actor_user_id=uid, + actor_username=uname, + method="POST", + path=f"/v1/port-traffic/tasks/{task_id}/start", + status_code=200, + detail={"id": task_id}, + ) + return out.model_dump() + + +@router.post("/tasks/{task_id}/pause") +def api_pause_task(task_id: str, request: Request, db: Session = Depends(get_db)): + out = set_task_status(db, task_id, "paused") + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.pause", + actor_user_id=uid, + actor_username=uname, + method="POST", + path=f"/v1/port-traffic/tasks/{task_id}/pause", + status_code=200, + detail={"id": task_id}, + ) + return out.model_dump() + + +@router.post("/tasks/{task_id}/stop") +def api_stop_task(task_id: str, request: Request, db: Session = Depends(get_db)): + out = set_task_status(db, task_id, "stopped") + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.stop", + actor_user_id=uid, + actor_username=uname, + method="POST", + path=f"/v1/port-traffic/tasks/{task_id}/stop", + status_code=200, + detail={"id": task_id}, + ) + return out.model_dump() + + +@router.get("/tasks/{task_id}/targets") +def api_list_targets(task_id: str, db: Session = Depends(get_db)): + return {"items": [t.model_dump() for t in list_targets(db, task_id)]} + + +@router.put("/tasks/{task_id}/targets") +def api_put_targets( + task_id: str, + body: PortTrafficTargetsPut, + request: Request, + db: Session = Depends(get_db), +): + items = put_targets(db, task_id, body) + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.targets.put", + actor_user_id=uid, + actor_username=uname, + method="PUT", + path=f"/v1/port-traffic/tasks/{task_id}/targets", + status_code=200, + detail={"id": task_id, "count": len(items)}, + ) + return {"items": [t.model_dump() for t in items]} + + +@router.post("/discover/ports") +def api_discover_ports(body: DiscoverPortsRequest, db: Session = Depends(get_db)): + return discover_ports(db, body).model_dump() + + +@router.get("/samples") +def api_samples( + target_id: str = Query(..., description="port_traffic_target row id"), + from_ts: datetime | None = Query(default=None, alias="from"), + to_ts: datetime | None = Query(default=None, alias="to"), + db: Session = Depends(get_db), +): + return get_samples(db, target_row_id=target_id, from_ts=from_ts, to_ts=to_ts).model_dump() diff --git a/netx_api/port_traffic_runner.py b/netx_api/port_traffic_runner.py new file mode 100644 index 0000000..eb06d26 --- /dev/null +++ b/netx_api/port_traffic_runner.py @@ -0,0 +1,258 @@ +"""Port traffic collection worker: claim task round, sample interfaces via CLI.""" + +from __future__ import annotations + +import logging +import time +from concurrent.futures import ThreadPoolExecutor, as_completed +from datetime import datetime +from threading import Lock +from typing import Any +from uuid import uuid4 + +from fastapi import HTTPException + +from .cli_resolve import resolve_cli_target +from .config import settings +from .db import SessionLocal +from .models import PortTrafficSample, PortTrafficTarget, PortTrafficTask +from .ne_session_factory import close_netmiko_connection, open_netmiko_connection +from .port_traffic_commands import commands_for_vendor, detail_command +from .port_traffic_parsers import parse_zte_interface_detail + +_log = logging.getLogger("netx.port_traffic.runner") +_pools: dict[str, ThreadPoolExecutor] = {} +_pools_lock = Lock() + + +def _utcnow() -> datetime: + return datetime.utcnow() + + +def _format_error(exc: BaseException) -> str: + return f"{type(exc).__name__}: {exc}"[:1020] + + +def _pool_for_task(task_id: str, concurrency: int) -> ThreadPoolExecutor: + with _pools_lock: + pool = _pools.get(task_id) + if pool is None: + workers = max(1, min(20, int(concurrency or 5))) + pool = ThreadPoolExecutor(max_workers=workers, thread_name_prefix=f"pt-{task_id[:8]}") + _pools[task_id] = pool + return pool + + +def _release_pool(task_id: str) -> None: + with _pools_lock: + pool = _pools.pop(task_id, None) + if pool is not None: + try: + pool.shutdown(wait=False, cancel_futures=False) + except TypeError: + pool.shutdown(wait=False) + except Exception: + _log.exception("port_traffic pool shutdown failed task=%s", task_id) + + +def _set_target_error(target_row_id: str, message: str) -> None: + db = SessionLocal() + try: + row = db.get(PortTrafficTarget, target_row_id) + if row: + row.last_error = message[:1020] + db.commit() + finally: + db.close() + + +def _claim_collect_round(task_id: str) -> list[str] | None: + """Mark task collect_running and return active target row ids, or None if skip.""" + for attempt in range(8): + db = SessionLocal() + try: + task = db.get(PortTrafficTask, task_id) + if not task: + return None + if str(task.status or "") != "running": + return None + if bool(task.collect_running): + return None + ended = task.last_collect_ended_at + interval = max(15, int(task.interval_sec or 60)) + if ended is not None: + elapsed = (_utcnow() - ended).total_seconds() + if elapsed < interval: + return None + targets = ( + db.query(PortTrafficTarget) + .filter( + PortTrafficTarget.task_id == task_id, + PortTrafficTarget.status == "active", + ) + .all() + ) + if not targets: + return None + task.collect_running = True + task.last_collect_started_at = _utcnow() + task.last_error = "" + task.updated_at = _utcnow() + db.commit() + return [str(t.id) for t in targets] + except Exception: + db.rollback() + _log.exception("port_traffic claim failed task=%s attempt=%s", task_id, attempt) + time.sleep(0.05 * (attempt + 1)) + finally: + db.close() + return None + + +def _finish_collect_round(task_id: str, *, error: str = "") -> None: + db = SessionLocal() + try: + task = db.get(PortTrafficTask, task_id) + if not task: + return + task.collect_running = False + task.last_collect_ended_at = _utcnow() + if error: + task.last_error = error[:1020] + task.updated_at = _utcnow() + db.commit() + finally: + db.close() + _release_pool(task_id) + + +def _run_show(creds: dict[str, Any], command: str, read_timeout: int) -> str: + conn = open_netmiko_connection(creds, session_timeout=read_timeout + 60) + try: + return str(conn.send_command(command_string=command, read_timeout=read_timeout) or "") + finally: + close_netmiko_connection(conn) + + +def _sample_one_target(target_row_id: str) -> None: + creds: dict[str, Any] | None = None + cmd = "" + per_cmd = int(settings.ne_collect_read_timeout_sec or 120) + cap = int(settings.ne_collect_run_timeout_cap_sec or 600) + + db = SessionLocal() + try: + row = db.get(PortTrafficTarget, target_row_id) + if not row or str(row.status or "") != "active": + return + source = str(row.source or "").strip().lower() + target_id = str(row.target_id or "").strip() + ifname = str(row.ifname or "").strip() + vendor_hint = str(row.vendor or "") + try: + if source == "managed": + creds, device = resolve_cli_target(db, managed_ne_id=target_id) + elif source == "ume": + creds, device = resolve_cli_target(db, ume_ne_id=target_id) + else: + row.last_error = "invalid_source" + db.commit() + return + except HTTPException as exc: + row.last_error = str(exc.detail or "resolve_failed")[:1020] + db.commit() + return + except Exception as exc: + row.last_error = _format_error(exc) + db.commit() + return + + vendor = str(device.get("vendor") or vendor_hint or "") + device_type = str(device.get("device_type") or "") + cmds = commands_for_vendor(vendor, device_type) + if cmds is None: + row.last_error = "unsupported_vendor" + db.commit() + return + cmd = detail_command(cmds, ifname) + finally: + db.close() + + if not creds or not cmd: + return + + budget = min(cap, per_cmd + 90) + try: + with ThreadPoolExecutor(max_workers=1) as pool: + fut = pool.submit(_run_show, creds, cmd, per_cmd) + raw = fut.result(timeout=budget) + except Exception as exc: + _set_target_error(target_row_id, _format_error(exc)) + return + + parsed = parse_zte_interface_detail(raw) + if parsed.in_bps == 0 and parsed.out_bps == 0 and parsed.bw_bps == 0 and not parsed.ifname: + _set_target_error(target_row_id, "parse_empty") + return + + now = _utcnow() + db = SessionLocal() + try: + row = db.get(PortTrafficTarget, target_row_id) + if not row: + return + bw = int(parsed.bw_bps or row.bw_bps or 0) + if bw and not row.bw_bps: + row.bw_bps = bw + db.add( + PortTrafficSample( + id=uuid4().hex, + target_row_id=target_row_id, + ts=now, + in_bps=float(parsed.in_bps), + out_bps=float(parsed.out_bps), + in_util_pct=float(parsed.in_util_pct), + out_util_pct=float(parsed.out_util_pct), + bw_bps=bw, + rate_period_sec=int(parsed.rate_period_sec or 0), + raw_ok=True, + message="", + ) + ) + row.last_error = "" + row.last_sample_at = now + db.commit() + except Exception: + db.rollback() + _log.exception("port_traffic sample save failed target=%s", target_row_id) + finally: + db.close() + + +def dispatch_collect(task_id: str) -> int: + """Claim and sample all active targets for a running task. Returns target count.""" + target_ids = _claim_collect_round(task_id) + if not target_ids: + return 0 + + db = SessionLocal() + try: + task = db.get(PortTrafficTask, task_id) + concurrency = int(task.concurrency or 5) if task else 5 + finally: + db.close() + + pool = _pool_for_task(task_id, concurrency) + futures = [pool.submit(_sample_one_target, tid) for tid in target_ids] + errors = 0 + try: + for fut in as_completed(futures): + try: + fut.result() + except Exception: + errors += 1 + _log.exception("port_traffic target worker failed task=%s", task_id) + finally: + err_msg = f"{errors}_target_errors" if errors else "" + _finish_collect_round(task_id, error=err_msg) + return len(target_ids) diff --git a/netx_api/port_traffic_scheduler.py b/netx_api/port_traffic_scheduler.py new file mode 100644 index 0000000..3690cf9 --- /dev/null +++ b/netx_api/port_traffic_scheduler.py @@ -0,0 +1,96 @@ +"""Background scheduler for port traffic collection + retention purge.""" + +from __future__ import annotations + +import logging +import threading +from datetime import datetime + +from .config import settings +from .db import SessionLocal +from .models import PortTrafficTask +from .port_traffic_runner import dispatch_collect +from .port_traffic_service import purge_expired_samples + +_log = logging.getLogger("netx.port_traffic.scheduler") +_stop = threading.Event() +_thread: threading.Thread | None = None +_purge_counter = 0 + + +def _utcnow() -> datetime: + return datetime.utcnow() + + +def try_dispatch_due_tasks() -> int: + """Dispatch collect rounds for due running tasks. Returns number started.""" + db = SessionLocal() + try: + tasks = ( + db.query(PortTrafficTask) + .filter(PortTrafficTask.status == "running", PortTrafficTask.collect_running.is_(False)) + .all() + ) + due_ids: list[str] = [] + now = _utcnow() + for task in tasks: + interval = max(15, int(task.interval_sec or 60)) + ended = task.last_collect_ended_at + if ended is None: + due_ids.append(str(task.id)) + continue + if (now - ended).total_seconds() >= interval: + due_ids.append(str(task.id)) + finally: + db.close() + + started = 0 + for tid in due_ids: + try: + n = dispatch_collect(tid) + if n: + started += 1 + _log.info("port_traffic collect started task=%s targets=%s", tid, n) + except Exception: + _log.exception("port_traffic dispatch failed task=%s", tid) + return started + + +def _loop() -> None: + global _purge_counter + tick = max(5, int(settings.port_traffic_scheduler_tick_sec or 15)) + _log.info("port_traffic scheduler started tick=%ss", tick) + while not _stop.is_set(): + try: + if bool(settings.port_traffic_scheduler_enabled): + try_dispatch_due_tasks() + _purge_counter += 1 + # Retention purge roughly every ~20 ticks + if _purge_counter >= 20: + _purge_counter = 0 + db = SessionLocal() + try: + purge_expired_samples(db) + finally: + db.close() + except Exception: + _log.exception("port_traffic scheduler tick failed") + _stop.wait(tick) + _log.info("port_traffic scheduler stopped") + + +def start_port_traffic_scheduler() -> None: + global _thread + if not bool(settings.port_traffic_scheduler_enabled): + _log.info("port_traffic scheduler disabled by settings") + return + if _thread and _thread.is_alive(): + return + _stop.clear() + _thread = threading.Thread(target=_loop, name="port-traffic-scheduler", daemon=True) + _thread.start() + _log.info("started thread %s alive=%s", _thread.name, _thread.is_alive()) + + +def stop_port_traffic_scheduler() -> None: + _stop.set() diff --git a/netx_api/port_traffic_schemas.py b/netx_api/port_traffic_schemas.py new file mode 100644 index 0000000..1e0f5ad --- /dev/null +++ b/netx_api/port_traffic_schemas.py @@ -0,0 +1,131 @@ +"""Pydantic schemas for port traffic monitoring API.""" + +from __future__ import annotations + +from datetime import datetime +from typing import Literal + +from pydantic import BaseModel, Field + + +class PortTrafficNeRef(BaseModel): + source: Literal["managed", "ume"] + id: str + ne_name: str = "" + ne_ip: str = "" + vendor: str = "" + + +class PortTrafficTargetIn(BaseModel): + source: Literal["managed", "ume"] + target_id: str + ne_name: str = "" + ne_ip: str = "" + vendor: str = "" + ifname: str + if_description: str = "" + bw_bps: int = 0 + + +class PortTrafficTaskCreate(BaseModel): + title: str = Field(min_length=1, max_length=256) + interval_sec: int = Field(default=60, ge=15, le=3600) + retention_days: int = Field(default=7, ge=1, le=90) + concurrency: int = Field(default=5, ge=1, le=20) + targets: list[PortTrafficTargetIn] = Field(default_factory=list) + start_now: bool = False + + +class PortTrafficTaskUpdate(BaseModel): + title: str | None = Field(default=None, min_length=1, max_length=256) + interval_sec: int | None = Field(default=None, ge=15, le=3600) + retention_days: int | None = Field(default=None, ge=1, le=90) + concurrency: int | None = Field(default=None, ge=1, le=20) + + +class PortTrafficTaskOut(BaseModel): + id: str + title: str + status: str + interval_sec: int + retention_days: int + concurrency: int + collect_running: bool = False + target_count: int = 0 + active_target_count: int = 0 + last_collect_started_at: datetime | None = None + last_collect_ended_at: datetime | None = None + last_error: str = "" + created_at: datetime | None = None + updated_at: datetime | None = None + + +class PortTrafficTargetOut(BaseModel): + id: str + task_id: str + source: str + target_id: str + ne_name: str + ne_ip: str + vendor: str + ifname: str + if_description: str + bw_bps: int + status: str + last_error: str = "" + last_sample_at: datetime | None = None + created_at: datetime | None = None + + +class PortTrafficTargetsPut(BaseModel): + targets: list[PortTrafficTargetIn] + + +class DiscoverPortsRequest(BaseModel): + source: Literal["managed", "ume"] + id: str + + +class DiscoverPortItem(BaseModel): + ifname: str + attribute: str = "" + mode: str = "" + bw_raw: str = "" + bw_bps: int = 0 + admin: str = "" + phy: str = "" + prot: str = "" + description: str = "" + + +class DiscoverPortsResponse(BaseModel): + source: str + id: str + ne_name: str = "" + ne_ip: str = "" + vendor: str = "" + vendor_key: str = "" + ports: list[DiscoverPortItem] = Field(default_factory=list) + + +class PortTrafficSamplePoint(BaseModel): + ts: datetime + in_bps: float + out_bps: float + in_util_pct: float + out_util_pct: float + bw_bps: int + rate_period_sec: int = 0 + + +class PortTrafficSamplesOut(BaseModel): + target: PortTrafficTargetOut + points: list[PortTrafficSamplePoint] = Field(default_factory=list) + + +class PortTrafficDashboardOut(BaseModel): + task_count: int = 0 + running_task_count: int = 0 + active_target_count: int = 0 + sample_count_24h: int = 0 + last_sample_at: datetime | None = None diff --git a/netx_api/port_traffic_service.py b/netx_api/port_traffic_service.py new file mode 100644 index 0000000..f676a2a --- /dev/null +++ b/netx_api/port_traffic_service.py @@ -0,0 +1,398 @@ +"""Port traffic monitoring service: CRUD, discover, samples, dashboard.""" + +from __future__ import annotations + +import logging +from datetime import datetime, timedelta +from typing import Any +from uuid import uuid4 + +from fastapi import HTTPException +from sqlalchemy import func +from sqlalchemy.orm import Session + +from .cli_resolve import resolve_cli_target +from .config import settings +from .models import PortTrafficSample, PortTrafficTarget, PortTrafficTask +from .ne_session_factory import close_netmiko_connection, open_netmiko_connection +from .port_traffic_commands import commands_for_vendor +from .port_traffic_parsers import brief_port_to_dict, parse_zte_interface_brief +from .port_traffic_schemas import ( + DiscoverPortItem, + DiscoverPortsRequest, + DiscoverPortsResponse, + PortTrafficDashboardOut, + PortTrafficSamplePoint, + PortTrafficSamplesOut, + PortTrafficTargetIn, + PortTrafficTargetOut, + PortTrafficTargetsPut, + PortTrafficTaskCreate, + PortTrafficTaskOut, + PortTrafficTaskUpdate, +) + +_log = logging.getLogger("netx.port_traffic.service") + + +def _utcnow() -> datetime: + return datetime.utcnow() + + +def _task_out(db: Session, task: PortTrafficTask) -> PortTrafficTaskOut: + tid = str(task.id) + total = db.query(PortTrafficTarget).filter(PortTrafficTarget.task_id == tid).count() + active = ( + db.query(PortTrafficTarget) + .filter(PortTrafficTarget.task_id == tid, PortTrafficTarget.status == "active") + .count() + ) + return PortTrafficTaskOut( + id=tid, + title=str(task.title or ""), + status=str(task.status or ""), + interval_sec=int(task.interval_sec or 60), + retention_days=int(task.retention_days or 7), + concurrency=int(task.concurrency or 5), + collect_running=bool(task.collect_running), + target_count=int(total), + active_target_count=int(active), + last_collect_started_at=task.last_collect_started_at, + last_collect_ended_at=task.last_collect_ended_at, + last_error=str(task.last_error or ""), + created_at=task.created_at, + updated_at=task.updated_at, + ) + + +def _target_out(row: PortTrafficTarget) -> PortTrafficTargetOut: + return PortTrafficTargetOut( + id=str(row.id), + task_id=str(row.task_id), + source=str(row.source or ""), + target_id=str(row.target_id or ""), + ne_name=str(row.ne_name or ""), + ne_ip=str(row.ne_ip or ""), + vendor=str(row.vendor or ""), + ifname=str(row.ifname or ""), + if_description=str(row.if_description or ""), + bw_bps=int(row.bw_bps or 0), + status=str(row.status or ""), + last_error=str(row.last_error or ""), + last_sample_at=row.last_sample_at, + created_at=row.created_at, + ) + + +def _assert_zte_targets(targets: list[PortTrafficTargetIn]) -> None: + for t in targets: + cmds = commands_for_vendor(t.vendor or "", "") + if cmds is None: + raise HTTPException( + status_code=400, + detail=f"vendor_not_supported_for_port_traffic: {t.vendor or 'unknown'} ({t.ne_name or t.target_id})", + ) + + +def list_tasks(db: Session, *, page: int = 1, page_size: int = 20) -> dict[str, Any]: + q = db.query(PortTrafficTask).order_by(PortTrafficTask.created_at.desc()) + total = q.count() + rows = q.offset((page - 1) * page_size).limit(page_size).all() + return { + "total": total, + "page": page, + "page_size": page_size, + "items": [_task_out(db, r).model_dump() for r in rows], + } + + +def get_task(db: Session, task_id: str) -> PortTrafficTaskOut: + task = db.get(PortTrafficTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + return _task_out(db, task) + + +def create_task(db: Session, body: PortTrafficTaskCreate) -> PortTrafficTaskOut: + _assert_zte_targets(body.targets) + now = _utcnow() + status = "running" if body.start_now and body.targets else "draft" + if body.start_now and not body.targets: + status = "draft" + task = PortTrafficTask( + id=uuid4().hex, + title=body.title.strip(), + status=status, + interval_sec=int(body.interval_sec), + retention_days=int(body.retention_days), + concurrency=int(body.concurrency), + created_at=now, + updated_at=now, + ) + db.add(task) + db.flush() + for t in body.targets: + db.add( + PortTrafficTarget( + id=uuid4().hex, + task_id=task.id, + source=t.source, + target_id=t.target_id, + ne_name=t.ne_name or "", + ne_ip=t.ne_ip or "", + vendor=t.vendor or "", + ifname=t.ifname.strip(), + if_description=t.if_description or "", + bw_bps=int(t.bw_bps or 0), + status="active", + created_at=now, + ) + ) + db.commit() + db.refresh(task) + return _task_out(db, task) + + +def update_task(db: Session, task_id: str, body: PortTrafficTaskUpdate) -> PortTrafficTaskOut: + task = db.get(PortTrafficTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + data = body.model_dump(exclude_unset=True) + if "title" in data and data["title"] is not None: + task.title = str(data["title"]).strip() + if "interval_sec" in data and data["interval_sec"] is not None: + task.interval_sec = int(data["interval_sec"]) + if "retention_days" in data and data["retention_days"] is not None: + task.retention_days = int(data["retention_days"]) + if "concurrency" in data and data["concurrency"] is not None: + task.concurrency = int(data["concurrency"]) + task.updated_at = _utcnow() + db.commit() + db.refresh(task) + return _task_out(db, task) + + +def delete_task(db: Session, task_id: str) -> dict[str, Any]: + task = db.get(PortTrafficTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + if bool(task.collect_running): + raise HTTPException(status_code=409, detail="collect_running") + targets = db.query(PortTrafficTarget).filter(PortTrafficTarget.task_id == task_id).all() + ids = [str(t.id) for t in targets] + if ids: + db.query(PortTrafficSample).filter(PortTrafficSample.target_row_id.in_(ids)).delete( + synchronize_session=False + ) + db.query(PortTrafficTarget).filter(PortTrafficTarget.task_id == task_id).delete( + synchronize_session=False + ) + db.delete(task) + db.commit() + return {"ok": True, "id": task_id} + + +def set_task_status(db: Session, task_id: str, status: str) -> PortTrafficTaskOut: + task = db.get(PortTrafficTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + if status == "running": + active = ( + db.query(PortTrafficTarget) + .filter(PortTrafficTarget.task_id == task_id, PortTrafficTarget.status == "active") + .count() + ) + if active <= 0: + raise HTTPException(status_code=400, detail="no_active_targets") + task.status = status + task.updated_at = _utcnow() + if status in ("stopped", "paused"): + # leave collect_running for runner to finish; recovery clears stuck flags + pass + db.commit() + db.refresh(task) + return _task_out(db, task) + + +def put_targets(db: Session, task_id: str, body: PortTrafficTargetsPut) -> list[PortTrafficTargetOut]: + task = db.get(PortTrafficTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + if bool(task.collect_running): + raise HTTPException(status_code=409, detail="collect_running") + _assert_zte_targets(body.targets) + old = db.query(PortTrafficTarget).filter(PortTrafficTarget.task_id == task_id).all() + old_ids = [str(t.id) for t in old] + if old_ids: + db.query(PortTrafficSample).filter(PortTrafficSample.target_row_id.in_(old_ids)).delete( + synchronize_session=False + ) + db.query(PortTrafficTarget).filter(PortTrafficTarget.task_id == task_id).delete( + synchronize_session=False + ) + now = _utcnow() + rows: list[PortTrafficTarget] = [] + for t in body.targets: + row = PortTrafficTarget( + id=uuid4().hex, + task_id=task_id, + source=t.source, + target_id=t.target_id, + ne_name=t.ne_name or "", + ne_ip=t.ne_ip or "", + vendor=t.vendor or "", + ifname=t.ifname.strip(), + if_description=t.if_description or "", + bw_bps=int(t.bw_bps or 0), + status="active", + created_at=now, + ) + db.add(row) + rows.append(row) + task.updated_at = now + db.commit() + return [_target_out(r) for r in rows] + + +def list_targets(db: Session, task_id: str) -> list[PortTrafficTargetOut]: + task = db.get(PortTrafficTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task_not_found") + rows = ( + db.query(PortTrafficTarget) + .filter(PortTrafficTarget.task_id == task_id) + .order_by(PortTrafficTarget.ne_name, PortTrafficTarget.ifname) + .all() + ) + return [_target_out(r) for r in rows] + + +def discover_ports(db: Session, body: DiscoverPortsRequest) -> DiscoverPortsResponse: + try: + if body.source == "managed": + creds, device = resolve_cli_target(db, managed_ne_id=body.id) + else: + creds, device = resolve_cli_target(db, ume_ne_id=body.id) + except HTTPException: + raise + except Exception as exc: + raise HTTPException(status_code=400, detail=f"resolve_failed: {exc}") from exc + + vendor = str(device.get("vendor") or creds.get("vendor") or "") + device_type = str(device.get("device_type") or creds.get("device_type") or "") + ne_name = str(device.get("name") or creds.get("host") or "") + ne_ip = str(device.get("ip_address") or creds.get("host") or "") + cmds = commands_for_vendor(vendor, device_type) + if cmds is None: + raise HTTPException( + status_code=400, + detail=f"vendor_not_supported_for_port_traffic: {vendor or 'unknown'}", + ) + + per_cmd = int(settings.ne_collect_read_timeout_sec or 120) + conn = open_netmiko_connection(creds, session_timeout=per_cmd + 60) + try: + raw = str(conn.send_command(command_string=cmds.brief, read_timeout=per_cmd) or "") + finally: + close_netmiko_connection(conn) + + ports = [DiscoverPortItem(**brief_port_to_dict(p)) for p in parse_zte_interface_brief(raw)] + return DiscoverPortsResponse( + source=body.source, + id=body.id, + ne_name=ne_name, + ne_ip=ne_ip, + vendor=vendor, + vendor_key=cmds.vendor_key, + ports=ports, + ) + + +def get_samples( + db: Session, + *, + target_row_id: str, + from_ts: datetime | None = None, + to_ts: datetime | None = None, +) -> PortTrafficSamplesOut: + target = db.get(PortTrafficTarget, target_row_id) + if not target: + raise HTTPException(status_code=404, detail="target_not_found") + now = _utcnow() + if to_ts is None: + to_ts = now + if from_ts is None: + from_ts = to_ts - timedelta(hours=1) + rows = ( + db.query(PortTrafficSample) + .filter( + PortTrafficSample.target_row_id == target_row_id, + PortTrafficSample.ts >= from_ts, + PortTrafficSample.ts <= to_ts, + PortTrafficSample.raw_ok.is_(True), + ) + .order_by(PortTrafficSample.ts.asc()) + .all() + ) + points = [ + PortTrafficSamplePoint( + ts=r.ts, + in_bps=float(r.in_bps or 0), + out_bps=float(r.out_bps or 0), + in_util_pct=float(r.in_util_pct or 0), + out_util_pct=float(r.out_util_pct or 0), + bw_bps=int(r.bw_bps or 0), + rate_period_sec=int(r.rate_period_sec or 0), + ) + for r in rows + ] + return PortTrafficSamplesOut(target=_target_out(target), points=points) + + +def dashboard(db: Session) -> PortTrafficDashboardOut: + task_count = db.query(PortTrafficTask).count() + running = db.query(PortTrafficTask).filter(PortTrafficTask.status == "running").count() + active_targets = db.query(PortTrafficTarget).filter(PortTrafficTarget.status == "active").count() + since = _utcnow() - timedelta(hours=24) + sample_count = ( + db.query(PortTrafficSample) + .filter(PortTrafficSample.ts >= since, PortTrafficSample.raw_ok.is_(True)) + .count() + ) + last = db.query(func.max(PortTrafficSample.ts)).scalar() + return PortTrafficDashboardOut( + task_count=int(task_count), + running_task_count=int(running), + active_target_count=int(active_targets), + sample_count_24h=int(sample_count), + last_sample_at=last, + ) + + +def purge_expired_samples(db: Session) -> int: + """Delete samples older than each task's retention_days.""" + tasks = db.query(PortTrafficTask).all() + deleted = 0 + now = _utcnow() + for task in tasks: + days = max(1, int(task.retention_days or 7)) + cutoff = now - timedelta(days=days) + target_ids = [ + str(t.id) + for t in db.query(PortTrafficTarget.id).filter(PortTrafficTarget.task_id == task.id).all() + ] + if not target_ids: + continue + n = ( + db.query(PortTrafficSample) + .filter( + PortTrafficSample.target_row_id.in_(target_ids), + PortTrafficSample.ts < cutoff, + ) + .delete(synchronize_session=False) + ) + deleted += int(n or 0) + if deleted: + db.commit() + _log.info("port_traffic retention purged samples=%s", deleted) + return deleted diff --git a/tests/test_port_traffic_parsers.py b/tests/test_port_traffic_parsers.py new file mode 100644 index 0000000..283bac9 --- /dev/null +++ b/tests/test_port_traffic_parsers.py @@ -0,0 +1,116 @@ +"""Unit tests for ZTE port traffic brief/detail parsers (sample from oclaw Untitled-1.ps1).""" + +from __future__ import annotations + +import unittest + +from netx_api.port_traffic_commands import commands_for_vendor, detail_command +from netx_api.port_traffic_parsers import ( + parse_bw_to_bps, + parse_zte_interface_brief, + parse_zte_interface_detail, +) + +SAMPLE_DETAIL = """\ +AL5458-ACC-6120HS#show interface xgei-1/1/0/1 +13:47:08 Africa/Algiers Fri Jul 31 2026 +xgei-1/1/0/1 is up, ifindex: 8194 + Description: C2930L100-EQ2 + Line protocol is up, IPv4 protocol is up, IPv6 protocol is down, + detected status is RX-OK/TX-OK + Last line protocol up time : 2026-05-26 15:14:49 + Hardware is XGigabit Ethernet, address is 00d0.0000.088f + Internet address is 192.168.1.1/30 + BW 1 Gbit/s + IP MTU 1500 bytes + MTU 1562 bytes + MPLS MTU 1548 bytes + + Fec-eth : N/A + Fec-bypass : N/A + ARP type ARP + ARP Timeout 04:00:00 + Last Clear Time : 2026-05-26 15:13:48 Last Refresh Time: 2026-07-31 13:47:00 + Rate period : 30 s + Input : 824 bit/s 1 packet/s + Output : 824 bit/s 1 packet/s + Peak rate: + Input : 3536 bit/s peak time 2026-05-26 15:17:10 + Output : 6896 bit/s peak time 2026-05-26 15:17:10 + Intf utilization: input 0% output 0% + HardwareCounters: + In_Bytes 462427066 In_Packets 6529765 +""" + +SAMPLE_BRIEF = """\ +AL5458-ACC-6120HS#show interface brief +13:43:00 Africa/Algiers Fri Jul 31 2026 +Interface Attribute Mode BW Admin Phy Prot Description +xgei-1/1/0/1 optical Duplex/full 1G up up up C2930L100-EQ2 +xgei-1/1/0/2 optical Duplex/full 1G up down down +xgei-1/1/0/4 optical Duplex/full 10G up down down 123456test11255 +cgei-1/1/0/33 optical Duplex/full 100G up down down +xxvgei-1/1/0/11 optical Duplex/full 25G up down down +smartgroup1 N/A N/A up up down +smartgroup10 N/A N/A 1G up up up +bvi2 N/A N/A up up up +""" + + +class BwParseTests(unittest.TestCase): + def test_compact(self): + self.assertEqual(parse_bw_to_bps("1G"), 1_000_000_000) + self.assertEqual(parse_bw_to_bps("10G"), 10_000_000_000) + self.assertEqual(parse_bw_to_bps("100G"), 100_000_000_000) + self.assertEqual(parse_bw_to_bps("25G"), 25_000_000_000) + self.assertEqual(parse_bw_to_bps("1M"), 1_000_000) + + def test_detail_line(self): + self.assertEqual(parse_bw_to_bps("BW 1 Gbit/s"), 1_000_000_000) + + +class BriefParserTests(unittest.TestCase): + def test_sample_brief(self): + rows = parse_zte_interface_brief(SAMPLE_BRIEF) + by_name = {r.ifname: r for r in rows} + self.assertIn("xgei-1/1/0/1", by_name) + r1 = by_name["xgei-1/1/0/1"] + self.assertEqual(r1.bw_bps, 1_000_000_000) + self.assertEqual(r1.admin, "up") + self.assertEqual(r1.phy, "up") + self.assertEqual(r1.prot, "up") + self.assertEqual(r1.description, "C2930L100-EQ2") + self.assertEqual(by_name["xgei-1/1/0/4"].bw_bps, 10_000_000_000) + self.assertEqual(by_name["cgei-1/1/0/33"].bw_bps, 100_000_000_000) + self.assertEqual(by_name["xxvgei-1/1/0/11"].bw_bps, 25_000_000_000) + self.assertEqual(by_name["smartgroup1"].bw_bps, 0) + self.assertEqual(by_name["smartgroup10"].bw_bps, 1_000_000_000) + self.assertEqual(by_name["xgei-1/1/0/2"].phy, "down") + + +class DetailParserTests(unittest.TestCase): + def test_sample_detail(self): + d = parse_zte_interface_detail(SAMPLE_DETAIL) + self.assertEqual(d.ifname, "xgei-1/1/0/1") + self.assertEqual(d.bw_bps, 1_000_000_000) + self.assertEqual(d.rate_period_sec, 30) + self.assertEqual(d.in_bps, 824.0) + self.assertEqual(d.out_bps, 824.0) + self.assertEqual(d.in_util_pct, 0.0) + self.assertEqual(d.out_util_pct, 0.0) + self.assertEqual(d.description, "C2930L100-EQ2") + # Peak rates must not override Rate period values + self.assertNotEqual(d.in_bps, 3536.0) + + +class CommandsTests(unittest.TestCase): + def test_zte_matrix(self): + cmds = commands_for_vendor("ZTE", "zxros") + assert cmds is not None + self.assertEqual(cmds.brief, "show interface brief") + self.assertEqual(detail_command(cmds, "xgei-1/1/0/1"), "show interface xgei-1/1/0/1") + self.assertIsNone(commands_for_vendor("Cisco", "ios")) + + +if __name__ == "__main__": + unittest.main() diff --git a/web/WEB.md b/web/WEB.md index 9dca917..b69933a 100644 --- a/web/WEB.md +++ b/web/WEB.md @@ -42,7 +42,7 @@ src/ | `/network/topology` | redirect → `/topology` | — | | `/network/tasks/collect` | 采集任务 | `network` | | `/network/tasks/config-sync` | 配置同步 | `network` | -| `/network/tasks/port-traffic` | 端口流量(占位) | `network` | +| `/network/tasks/port-traffic` | 端口流量监控(任务向导 + 大屏) | `network` | | `/topology` | 拓扑管理 | `topology` | | `/webcrt` | WebCRT 终端 | `webcrt` | | `/collect` | redirect → `/network/tasks/collect` | — | @@ -105,6 +105,15 @@ src/ - 进程启动宽限:`NETX_CONFIG_SYNC_STARTUP_GRACE_SEC`(默认 3600)仅约束**新建**自动周期,不影响续跑 - 前端:`/network/tasks/config-sync`(看板)+ `/network/configs`(查看) +## 端口流量监控 + +- API:`/v1/port-traffic/*`(任务 CRUD、discover/ports、samples、dashboard) +- 厂商:本期仅 ZTE(`show interface brief` / `show interface {if}`),解析 Rate period **bit/s** +- 调度:`NETX_PORT_TRAFFIC_SCHEDULER_ENABLED`(默认开),tick `NETX_PORT_TRAFFIC_SCHEDULER_TICK_SEC`(默认 15) +- 单飞:同一任务同时只允许一轮采集;崩溃启动清除 `collect_running` +- 保留:按任务 `retention_days` 清理过期 sample +- 前端:`/network/tasks/port-traffic`(四步向导 + 任务启停 + uPlot 监控大屏) + ## WebCRT - API:`POST /v1/webcrt/sessions`(`ne_id` 或 `ume_ne_id`)、`WS /v1/webcrt/sessions/{id}/ws`、`DELETE /v1/webcrt/sessions/{id}` diff --git a/web/package-lock.json b/web/package-lock.json index 83f76a5..91167fd 100644 --- a/web/package-lock.json +++ b/web/package-lock.json @@ -14,7 +14,8 @@ "@xyflow/react": "^12.11.2", "react": "^19.2.5", "react-dom": "^19.2.5", - "react-router-dom": "^7.14.2" + "react-router-dom": "^7.14.2", + "uplot": "^1.6.32" }, "devDependencies": { "@eslint/js": "^10.0.1", @@ -2889,6 +2890,12 @@ "browserslist": ">= 4.21.0" } }, + "node_modules/uplot": { + "version": "1.6.32", + "resolved": "https://registry.npmjs.org/uplot/-/uplot-1.6.32.tgz", + "integrity": "sha512-KIMVnG68zvu5XXUbC4LQEPnhwOxBuLyW1AHtpm6IKTXImkbLgkMy+jabjLgSLMasNuGGzQm/ep3tOkyTxpiQIw==", + "license": "MIT" + }, "node_modules/uri-js": { "version": "4.4.1", "resolved": "https://registry.npmjs.org/uri-js/-/uri-js-4.4.1.tgz", diff --git a/web/package.json b/web/package.json index 560323c..8567dd9 100644 --- a/web/package.json +++ b/web/package.json @@ -19,7 +19,8 @@ "@xyflow/react": "^12.11.2", "react": "^19.2.5", "react-dom": "^19.2.5", - "react-router-dom": "^7.14.2" + "react-router-dom": "^7.14.2", + "uplot": "^1.6.32" }, "devDependencies": { "@eslint/js": "^10.0.1", diff --git a/web/src/App.tsx b/web/src/App.tsx index 3c91fa7..2b79f86 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -13,7 +13,7 @@ import { NetworkLayout } from "./pages/network/NetworkLayout"; import { NetworkDevicesPage } from "./pages/network/NetworkDevicesPage"; import { NetworkAlarmsPage } from "./pages/network/NetworkAlarmsPage"; import { NetworkConfigsPage } from "./pages/network/NetworkConfigsPage"; -import { NetworkPlaceholderPage } from "./pages/network/NetworkPlaceholderPage"; +import { PortTrafficPage } from "./pages/network/PortTrafficPage"; import { LoginPage } from "./pages/LoginPage"; import { UsersPage } from "./pages/UsersPage"; import { AuditPage } from "./pages/AuditPage"; @@ -95,7 +95,7 @@ function ProtectedApp() { } /> } /> } /> - } /> + } /> } /> } /> diff --git a/web/src/constants/queryKeys.ts b/web/src/constants/queryKeys.ts index fb0a30e..601b6b5 100644 --- a/web/src/constants/queryKeys.ts +++ b/web/src/constants/queryKeys.ts @@ -59,4 +59,10 @@ export const queryKeys = { networkConfigs: (page: number, keyword: string, source: string, vendor: string) => ["networkConfigs", page, keyword, source, vendor] as const, networkConfigDetail: (source: string, id: string) => ["networkConfigDetail", source, id] as const, + portTrafficDashboard: ["portTrafficDashboard"] as const, + portTrafficTasksAll: ["portTrafficTasks"] as const, + portTrafficTasks: (page: number) => ["portTrafficTasks", page] as const, + portTrafficTargets: (taskId: string) => ["portTrafficTargets", taskId] as const, + portTrafficSamples: (targetId: string, rangeHours: number) => + ["portTrafficSamples", targetId, rangeHours] as const, }; diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index b0e2032..1a670be 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -60,12 +60,65 @@ const en = { placeholder: { title: "Coming soon", body: "This capability is not available yet.", - portTrafficTitle: "Port traffic monitoring", - portTrafficBody: "Planned: collect per-port traffic on a schedule and show trends. Suggested path:", - portTrafficStep1: "Task definition: NE/port scope, interval, retention", - portTrafficStep2: "Collection: prefer CLI (display/show interface); SNMP later at scale", - portTrafficStep3: "Storage: time series or hourly rollups (in/out bps, errors)", - portTrafficStep4: "UI: task list + per-port trend charts", + }, + }, + portTraffic: { + title: "Port traffic", + create: "New task", + backList: "Back to list", + created: "Task created", + createFailed: "Create failed", + started: "Started", + paused: "Paused", + stopped: "Stopped", + deleted: "Deleted", + confirmDelete: "Delete this task and its samples?", + empty: "No monitoring tasks yet. Create one to start.", + wall: "Ops wall", + start: "Start", + pause: "Pause", + stop: "Stop", + delete: "Delete", + step1: "Params", + step2: "NEs", + step3: "Ports", + step4: "Confirm", + prev: "Back", + next: "Next", + fieldTitle: "Task name", + titlePh: "e.g. Access uplinks", + titleRequired: "Task name is required", + portsRequired: "Select at least one port", + interval: "Interval (sec)", + retention: "Retention (days)", + concurrency: "Concurrency", + startNow: "Start after create", + neKeywordPh: "Name / IP / vendor", + selectedNe: "{{count}} NE(s) selected", + selectedPorts: "{{count}} port(s) selected", + pickNeFirst: "Select NEs in the previous step first", + discover: "Discover ports", + confirmCreate: "Create task", + wallTask: "Task", + wallPort: "Interface", + pickPort: "Select interface", + range: "Range", + kpi: { + tasks: "Tasks", + running: "Running", + ports: "Ports", + samples24h: "Samples 24h", + }, + col: { + title: "Task", + status: "Status", + ports: "Ports", + interval: "Interval", + last: "Last collect", + actions: "Actions", + source: "Source", + name: "Device", + vendor: "Vendor", }, }, configSync: { diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index dc35aa8..2f7d474 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -60,12 +60,65 @@ const zh = { placeholder: { title: "功能规划中", body: "该能力尚未开放。", - portTrafficTitle: "端口流量监控", - portTrafficBody: "规划中:按任务采集设备端口流量并展示趋势。建议路径如下。", - portTrafficStep1: "任务定义:选择网元/端口范围、采集周期与保留策略", - portTrafficStep2: "采集:优先 CLI(display/show interface),规模上来后再 SNMP", - portTrafficStep3: "存储:端口 in/out bps、error 等时序或按小时汇总", - portTrafficStep4: "界面:任务列表 + 单端口趋势图", + }, + }, + portTraffic: { + title: "端口流量监控", + create: "新建任务", + backList: "返回列表", + created: "任务已创建", + createFailed: "创建失败", + started: "已启动", + paused: "已暂停", + stopped: "已停止", + deleted: "已删除", + confirmDelete: "删除该任务及其采样数据?", + empty: "暂无监控任务。请先创建任务。", + wall: "监控大屏", + start: "启动", + pause: "暂停", + stop: "停止", + delete: "删除", + step1: "参数", + step2: "选网元", + step3: "选端口", + step4: "确认", + prev: "上一步", + next: "下一步", + fieldTitle: "任务名称", + titlePh: "例如:接入侧上联", + titleRequired: "请填写任务名称", + portsRequired: "请至少选择一个端口", + interval: "周期(秒)", + retention: "保留(天)", + concurrency: "并发", + startNow: "创建后立即启动", + neKeywordPh: "名称 / IP / 厂商", + selectedNe: "已选 {{count}} 台网元", + selectedPorts: "已选 {{count}} 个端口", + pickNeFirst: "请先在上一步选择网元", + discover: "拉取端口", + confirmCreate: "创建任务", + wallTask: "任务", + wallPort: "接口", + pickPort: "选择接口", + range: "时间范围", + kpi: { + tasks: "任务数", + running: "运行中", + ports: "监控端口", + samples24h: "24h 样点", + }, + col: { + title: "任务", + status: "状态", + ports: "端口", + interval: "周期", + last: "最近采集", + actions: "操作", + source: "来源", + name: "设备", + vendor: "厂商", }, }, configSync: { diff --git a/web/src/index.css b/web/src/index.css index ded196a..9ca3d34 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -2918,3 +2918,196 @@ pre { color: #475569; line-height: 1.6; } + +/* Port traffic monitoring */ +.pt-wizard__steps { + display: flex; + flex-wrap: wrap; + gap: 8px; + margin-bottom: 16px; +} + +.pt-wizard__step { + border: 1px solid #cbd5e1; + background: #f8fafc; + color: #334155; + border-radius: 6px; + padding: 6px 10px; + font-size: 13px; + cursor: pointer; +} + +.pt-wizard__step.is-active { + border-color: #0f766e; + background: #ccfbf1; + color: #115e59; +} + +.pt-wizard__ne-block { + margin-bottom: 16px; + padding-bottom: 12px; + border-bottom: 1px solid #e2e8f0; +} + +.pt-wizard__confirm-list { + margin: 8px 0 16px; + padding-left: 1.2rem; + max-height: 220px; + overflow: auto; + color: #334155; +} + +.pt-wall-page__filters { + margin-bottom: 12px; + gap: 12px 18px; + align-items: flex-end; +} + +.pt-wall-page__filters label { + display: inline-flex; + flex-direction: column; + gap: 4px; + font-size: 12px; + color: #64748b; +} + +.pt-wall-page__filters select { + min-width: 160px; +} + +.pt-wall { + --pt-bg: #0b1220; + --pt-panel: #111827; + --pt-line: rgba(148, 163, 184, 0.22); + --pt-text: #e2e8f0; + --pt-muted: #94a3b8; + --pt-in: #2dd4bf; + --pt-out: #f59e0b; + background: + linear-gradient(180deg, rgba(15, 23, 42, 0.92), rgba(11, 18, 32, 0.98)), + repeating-linear-gradient( + 0deg, + transparent, + transparent 23px, + rgba(148, 163, 184, 0.05) 24px + ), + repeating-linear-gradient( + 90deg, + transparent, + transparent 23px, + rgba(148, 163, 184, 0.05) 24px + ); + border: 1px solid #1f2937; + border-radius: 10px; + padding: 16px 18px 12px; + color: var(--pt-text); +} + +.pt-wall__head { + display: flex; + justify-content: space-between; + align-items: baseline; + gap: 12px; + margin-bottom: 14px; +} + +.pt-wall__title { + display: flex; + flex-direction: column; + gap: 2px; + min-width: 0; +} + +.pt-wall__ne { + font-size: 13px; + color: var(--pt-muted); +} + +.pt-wall__if { + font-size: 20px; + font-weight: 650; + letter-spacing: 0.02em; + font-family: "IBM Plex Mono", "Cascadia Mono", Consolas, monospace; +} + +.pt-wall__range { + font-size: 12px; + color: var(--pt-muted); + white-space: nowrap; +} + +.pt-wall__kpis { + display: grid; + grid-template-columns: repeat(4, minmax(0, 1fr)); + gap: 10px; + margin-bottom: 14px; +} + +.pt-wall__kpi { + background: rgba(17, 24, 39, 0.72); + border: 1px solid var(--pt-line); + border-radius: 8px; + padding: 10px 12px; +} + +.pt-wall__kpi-label { + font-size: 11px; + text-transform: uppercase; + letter-spacing: 0.06em; + color: var(--pt-muted); + margin-bottom: 4px; +} + +.pt-wall__kpi-value { + font-size: 22px; + font-weight: 650; + font-variant-numeric: tabular-nums; + font-family: "IBM Plex Mono", "Cascadia Mono", Consolas, monospace; + line-height: 1.2; +} + +.pt-wall__kpi--in .pt-wall__kpi-value { + color: var(--pt-in); +} + +.pt-wall__kpi--out .pt-wall__kpi-value { + color: var(--pt-out); +} + +.pt-wall__kpi-unit, +.pt-wall__kpi-sep { + margin-left: 4px; + font-size: 12px; + color: var(--pt-muted); + font-weight: 500; +} + +.pt-wall__chart { + min-height: 240px; + width: 100%; +} + +.pt-wall__chart .uplot, +.pt-wall__chart .u-wrap { + font-family: "IBM Plex Mono", "Cascadia Mono", Consolas, monospace; +} + +.pt-wall__chart .u-legend { + color: var(--pt-muted); +} + +.pt-wall__empty { + display: flex; + align-items: center; + justify-content: center; + min-height: 240px; + color: var(--pt-muted); + border: 1px dashed rgba(148, 163, 184, 0.28); + border-radius: 8px; +} + +@media (max-width: 900px) { + .pt-wall__kpis { + grid-template-columns: repeat(2, minmax(0, 1fr)); + } +} diff --git a/web/src/pages/network/NetworkPlaceholderPage.tsx b/web/src/pages/network/NetworkPlaceholderPage.tsx index f348541..a1f088a 100644 --- a/web/src/pages/network/NetworkPlaceholderPage.tsx +++ b/web/src/pages/network/NetworkPlaceholderPage.tsx @@ -1,26 +1,13 @@ -import { useI18n } from "../../i18n"; - type Props = { - kind: "port-traffic"; + kind?: string; }; -export function NetworkPlaceholderPage({ kind }: Props) { - const { t } = useI18n(); - const titleKey = kind === "port-traffic" ? "network.placeholder.portTrafficTitle" : "network.placeholder.title"; - const bodyKey = kind === "port-traffic" ? "network.placeholder.portTrafficBody" : "network.placeholder.body"; - +/** Generic placeholder for unfinished network sub-pages. */ +export function NetworkPlaceholderPage(_props: Props) { return (
-

{t(titleKey)}

-

{t(bodyKey)}

- {kind === "port-traffic" ? ( -
    -
  • {t("network.placeholder.portTrafficStep1")}
  • -
  • {t("network.placeholder.portTrafficStep2")}
  • -
  • {t("network.placeholder.portTrafficStep3")}
  • -
  • {t("network.placeholder.portTrafficStep4")}
  • -
- ) : null} +

Coming soon

+

This capability is not available yet.

); } diff --git a/web/src/pages/network/PortTrafficPage.tsx b/web/src/pages/network/PortTrafficPage.tsx new file mode 100644 index 0000000..98a3923 --- /dev/null +++ b/web/src/pages/network/PortTrafficPage.tsx @@ -0,0 +1,645 @@ +import { useEffect, useMemo, useState } from "react"; +import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query"; +import { + createPortTrafficTask, + deletePortTrafficTask, + discoverPortTrafficPorts, + fetchCliTargets, + fetchPortTrafficDashboard, + fetchPortTrafficSamples, + fetchPortTrafficTargets, + fetchPortTrafficTasks, + pausePortTrafficTask, + startPortTrafficTask, + stopPortTrafficTask, +} from "../../services/api"; +import { queryKeys } from "../../constants/queryKeys"; +import { useI18n } from "../../i18n"; +import { useToast } from "../../hooks/useToast"; +import type { CliTargetItem, PortTrafficDiscoverPort, PortTrafficTargetIn } from "../../types"; +import { pageCount } from "../../utils/display"; +import { formatSystemTime } from "../../utils/time"; +import { PortTrafficWall } from "./PortTrafficWall"; + +const POLL_MS = 5000; +const TARGET_PAGE_SIZE = 20; + +type PortPick = PortTrafficTargetIn & { key: string }; + +function neKey(source: string, id: string) { + return `${source}:${id}`; +} + +function portKey(source: string, id: string, ifname: string) { + return `${source}:${id}:${ifname}`; +} + +function formatBw(bps: number) { + if (!bps) return "—"; + if (bps >= 1e9) return `${bps / 1e9}G`; + if (bps >= 1e6) return `${bps / 1e6}M`; + return `${bps}`; +} + +export function PortTrafficPage() { + const { t } = useI18n(); + const { showOk, showError } = useToast(); + const queryClient = useQueryClient(); + + const [view, setView] = useState<"list" | "wizard" | "wall">("list"); + const [taskPage, setTaskPage] = useState(1); + const [wallTaskId, setWallTaskId] = useState(""); + const [wallTargetId, setWallTargetId] = useState(""); + const [rangeHours, setRangeHours] = useState(1); + + // Wizard state + const [step, setStep] = useState(1); + const [title, setTitle] = useState(""); + const [intervalSec, setIntervalSec] = useState(60); + const [retentionDays, setRetentionDays] = useState(7); + const [concurrency, setConcurrency] = useState(5); + const [startNow, setStartNow] = useState(true); + const [neKeyword, setNeKeyword] = useState(""); + const [nePage, setNePage] = useState(1); + const [selectedNes, setSelectedNes] = useState>({}); + const [portsByNe, setPortsByNe] = useState< + Record + >({}); + const [pickedPorts, setPickedPorts] = useState>({}); + + const dashQuery = useQuery({ + queryKey: queryKeys.portTrafficDashboard, + queryFn: fetchPortTrafficDashboard, + staleTime: 2000, + refetchInterval: POLL_MS, + }); + + const tasksQuery = useQuery({ + queryKey: queryKeys.portTrafficTasks(taskPage), + queryFn: () => fetchPortTrafficTasks({ page: taskPage, pageSize: 10 }), + staleTime: 2000, + refetchInterval: POLL_MS, + }); + + const wallTasksQuery = useQuery({ + queryKey: queryKeys.portTrafficTasks(0), + queryFn: () => fetchPortTrafficTasks({ page: 1, pageSize: 100 }), + enabled: view === "wall", + staleTime: 5000, + }); + + const nesQuery = useQuery({ + queryKey: queryKeys.cliTargets(neKeyword, nePage, TARGET_PAGE_SIZE), + queryFn: () => + fetchCliTargets({ source: "all", keyword: neKeyword, page: nePage, pageSize: TARGET_PAGE_SIZE }), + enabled: view === "wizard" && step === 2, + staleTime: 5000, + }); + + const wallTargetsQuery = useQuery({ + queryKey: queryKeys.portTrafficTargets(wallTaskId), + queryFn: () => fetchPortTrafficTargets(wallTaskId), + enabled: view === "wall" && Boolean(wallTaskId), + staleTime: 3000, + }); + + const samplesQuery = useQuery({ + queryKey: queryKeys.portTrafficSamples(wallTargetId, rangeHours), + queryFn: () => { + const to = new Date(); + const from = new Date(to.getTime() - rangeHours * 3600 * 1000); + return fetchPortTrafficSamples({ + targetId: wallTargetId, + from: from.toISOString(), + to: to.toISOString(), + }); + }, + enabled: view === "wall" && Boolean(wallTargetId), + staleTime: 2000, + refetchInterval: POLL_MS, + }); + + const invalidateAll = () => { + void queryClient.invalidateQueries({ queryKey: queryKeys.portTrafficTasksAll }); + void queryClient.invalidateQueries({ queryKey: queryKeys.portTrafficDashboard }); + }; + + const createMut = useMutation({ + mutationFn: createPortTrafficTask, + onSuccess: () => { + showOk(t("portTraffic.created")); + invalidateAll(); + resetWizard(); + setView("list"); + }, + onError: (e: Error) => showError(e.message || t("portTraffic.createFailed")), + }); + + const startMut = useMutation({ + mutationFn: startPortTrafficTask, + onSuccess: () => { + showOk(t("portTraffic.started")); + invalidateAll(); + }, + onError: (e: Error) => showError(e.message), + }); + const pauseMut = useMutation({ + mutationFn: pausePortTrafficTask, + onSuccess: () => { + showOk(t("portTraffic.paused")); + invalidateAll(); + }, + onError: (e: Error) => showError(e.message), + }); + const stopMut = useMutation({ + mutationFn: stopPortTrafficTask, + onSuccess: () => { + showOk(t("portTraffic.stopped")); + invalidateAll(); + }, + onError: (e: Error) => showError(e.message), + }); + const deleteMut = useMutation({ + mutationFn: deletePortTrafficTask, + onSuccess: () => { + showOk(t("portTraffic.deleted")); + invalidateAll(); + }, + onError: (e: Error) => showError(e.message), + }); + + const resetWizard = () => { + setStep(1); + setTitle(""); + setIntervalSec(60); + setRetentionDays(7); + setConcurrency(5); + setStartNow(true); + setSelectedNes({}); + setPortsByNe({}); + setPickedPorts({}); + setNeKeyword(""); + setNePage(1); + }; + + const openWizard = () => { + resetWizard(); + setView("wizard"); + }; + + const openWall = (taskId: string) => { + setWallTaskId(taskId); + setWallTargetId(""); + setView("wall"); + }; + + const toggleNe = (row: CliTargetItem) => { + const k = neKey(row.source, row.id); + setSelectedNes((prev) => { + const next = { ...prev }; + if (next[k]) delete next[k]; + else next[k] = row; + return next; + }); + }; + + const discoverNe = async (row: CliTargetItem) => { + const k = neKey(row.source, row.id); + setPortsByNe((prev) => ({ ...prev, [k]: { loading: true, ports: prev[k]?.ports || [] } })); + try { + const source = row.source === "ume" ? "ume" : "managed"; + const res = await discoverPortTrafficPorts({ source, id: row.id }); + setPortsByNe((prev) => ({ + ...prev, + [k]: { + ports: res.ports, + meta: { ne_name: res.ne_name || row.name, ne_ip: res.ne_ip || row.ip_address, vendor: res.vendor || row.vendor || "" }, + }, + })); + } catch (e) { + const msg = e instanceof Error ? e.message : "discover_failed"; + setPortsByNe((prev) => ({ ...prev, [k]: { ports: [], error: msg } })); + showError(msg); + } + }; + + const togglePort = (ne: CliTargetItem, port: PortTrafficDiscoverPort, meta?: { ne_name: string; ne_ip: string; vendor: string }) => { + const source = (ne.source === "ume" ? "ume" : "managed") as "managed" | "ume"; + const k = portKey(source, ne.id, port.ifname); + setPickedPorts((prev) => { + const next = { ...prev }; + if (next[k]) { + delete next[k]; + return next; + } + next[k] = { + key: k, + source, + target_id: ne.id, + ne_name: meta?.ne_name || ne.name, + ne_ip: meta?.ne_ip || ne.ip_address, + vendor: meta?.vendor || ne.vendor || "", + ifname: port.ifname, + if_description: port.description, + bw_bps: port.bw_bps, + }; + return next; + }); + }; + + const selectedNeList = useMemo(() => Object.values(selectedNes), [selectedNes]); + const pickedList = useMemo(() => Object.values(pickedPorts), [pickedPorts]); + const wallTargets = wallTargetsQuery.data?.items || []; + const selectedWallTarget = wallTargets.find((x) => x.id === wallTargetId) || null; + + useEffect(() => { + if (view !== "wall" || wallTargetId || !wallTargets.length) return; + setWallTargetId(wallTargets[0].id); + }, [view, wallTargetId, wallTargets]); + + const submitWizard = () => { + if (!title.trim()) { + showError(t("portTraffic.titleRequired")); + return; + } + if (!pickedList.length) { + showError(t("portTraffic.portsRequired")); + return; + } + createMut.mutate({ + title: title.trim(), + interval_sec: intervalSec, + retention_days: retentionDays, + concurrency, + start_now: startNow, + targets: pickedList.map(({ key: _k, ...rest }) => rest), + }); + }; + + const dash = dashQuery.data; + const tasks = tasksQuery.data?.items || []; + const taskPages = pageCount(tasksQuery.data?.total || 0, 10); + + return ( +
+
+

{t("portTraffic.title")}

+
+ {view !== "list" ? ( + + ) : ( + + )} +
+
+ + {view === "list" ? ( + <> +
+
+
{t("portTraffic.kpi.tasks")}
+
{dash?.task_count ?? "—"}
+
+
+
{t("portTraffic.kpi.running")}
+
{dash?.running_task_count ?? "—"}
+
+
+
{t("portTraffic.kpi.ports")}
+
{dash?.active_target_count ?? "—"}
+
+
+
{t("portTraffic.kpi.samples24h")}
+
{dash?.sample_count_24h ?? "—"}
+
+
+ + + + + + + + + + + + + + {!tasks.length ? ( + + + + ) : ( + tasks.map((row) => ( + + + + + + + + + )) + )} + +
{t("portTraffic.col.title")}{t("portTraffic.col.status")}{t("portTraffic.col.ports")}{t("portTraffic.col.interval")}{t("portTraffic.col.last")}{t("portTraffic.col.actions")}
+ {t("portTraffic.empty")} +
{row.title} + {row.status} + {row.collect_running ? " · collecting" : ""} + + {row.active_target_count}/{row.target_count} + {row.interval_sec}s{formatSystemTime(row.last_collect_ended_at) || "—"} +
+ + {row.status !== "running" ? ( + + ) : ( + + )} + {row.status !== "stopped" ? ( + + ) : null} + +
+
+
+ + + {taskPage}/{Math.max(1, taskPages)} + + +
+ + ) : null} + + {view === "wizard" ? ( +
+
+ {[1, 2, 3, 4].map((n) => ( + + ))} +
+ + {step === 1 ? ( +
+ + + + + +
+ ) : null} + + {step === 2 ? ( + <> +
+ { + setNeKeyword(e.target.value); + setNePage(1); + }} + placeholder={t("portTraffic.neKeywordPh")} + /> + {t("portTraffic.selectedNe", { count: selectedNeList.length })} +
+ + + + + + + + + + + {(nesQuery.data?.items || []).map((row) => { + const k = neKey(row.source, row.id); + return ( + + + + + + + + ); + })} + +
+ {t("portTraffic.col.source")}{t("portTraffic.col.name")}IP{t("portTraffic.col.vendor")}
+ toggleNe(row)} /> + {row.source}{row.name}{row.ip_address}{row.vendor || "—"}
+ + ) : null} + + {step === 3 ? ( +
+ {!selectedNeList.length ?

{t("portTraffic.pickNeFirst")}

: null} + {selectedNeList.map((ne) => { + const k = neKey(ne.source, ne.id); + const bucket = portsByNe[k]; + return ( +
+
+ + {ne.name} ({ne.ip_address}) + + + {bucket?.error ? {bucket.error} : null} +
+ {bucket?.ports?.length ? ( + + + + + + + + + + + {bucket.ports.map((p) => { + const pk = portKey(ne.source === "ume" ? "ume" : "managed", ne.id, p.ifname); + return ( + + + + + + + + ); + })} + +
+ InterfaceBWAdmin/Phy/ProtDescription
+ togglePort(ne, p, bucket.meta)} + /> + {p.ifname}{p.bw_raw || formatBw(p.bw_bps)} + {p.admin}/{p.phy}/{p.prot} + {p.description || "—"}
+ ) : null} +
+ ); + })} +

{t("portTraffic.selectedPorts", { count: String(pickedList.length) })}

+
+ ) : null} + + {step === 4 ? ( +
+

+ {title || "—"} · {intervalSec}s · {retentionDays}d · ×{concurrency} +

+

{t("portTraffic.selectedPorts", { count: String(pickedList.length) })}

+
    + {pickedList.slice(0, 40).map((p) => ( +
  • + {p.ne_name} / {p.ifname} ({formatBw(p.bw_bps || 0)}) +
  • + ))} + {pickedList.length > 40 ?
  • … +{pickedList.length - 40}
  • : null} +
+ +
+ ) : null} + +
+ + +
+
+ ) : null} + + {view === "wall" ? ( +
+
+ + + +
+ +
+ ) : null} +
+ ); +} diff --git a/web/src/pages/network/PortTrafficWall.tsx b/web/src/pages/network/PortTrafficWall.tsx new file mode 100644 index 0000000..2f3b061 --- /dev/null +++ b/web/src/pages/network/PortTrafficWall.tsx @@ -0,0 +1,157 @@ +import { useEffect, useMemo, useRef } from "react"; +import uPlot from "uplot"; +import "uplot/dist/uPlot.min.css"; +import type { PortTrafficSamplePoint, PortTrafficTarget } from "../../types"; + +function formatBps(n: number): string { + const v = Math.abs(n); + if (v >= 1e9) return `${(n / 1e9).toFixed(2)} G`; + if (v >= 1e6) return `${(n / 1e6).toFixed(2)} M`; + if (v >= 1e3) return `${(n / 1e3).toFixed(1)} K`; + return `${n.toFixed(0)}`; +} + +function formatBwLabel(bps: number): string { + if (!bps) return "—"; + return `${formatBps(bps)}bit/s`; +} + +type Props = { + target: PortTrafficTarget | null; + points: PortTrafficSamplePoint[]; + rangeLabel: string; +}; + +export function PortTrafficWall({ target, points, rangeLabel }: Props) { + const rootRef = useRef(null); + const plotRef = useRef(null); + + const latest = points.length ? points[points.length - 1] : null; + const kpi = useMemo(() => { + const bw = latest?.bw_bps || target?.bw_bps || 0; + return { + bw, + inBps: latest?.in_bps ?? 0, + outBps: latest?.out_bps ?? 0, + inUtil: latest?.in_util_pct ?? 0, + outUtil: latest?.out_util_pct ?? 0, + }; + }, [latest, target]); + + useEffect(() => { + const el = rootRef.current; + if (!el) return; + + const xs = points.map((p) => Math.floor(new Date(p.ts).getTime() / 1000)); + const inSeries = points.map((p) => p.in_bps); + const outSeries = points.map((p) => p.out_bps); + + const opts: uPlot.Options = { + width: Math.max(320, el.clientWidth || 800), + height: Math.max(220, Math.min(420, Math.floor((el.clientWidth || 800) * 0.32))), + series: [ + {}, + { + label: "In", + stroke: "#2dd4bf", + width: 2, + points: { show: false }, + }, + { + label: "Out", + stroke: "#f59e0b", + width: 2, + points: { show: false }, + }, + ], + axes: [ + { + stroke: "#94a3b8", + grid: { stroke: "rgba(148,163,184,0.18)", width: 1 }, + ticks: { stroke: "rgba(148,163,184,0.35)" }, + }, + { + stroke: "#94a3b8", + grid: { stroke: "rgba(148,163,184,0.12)", width: 1 }, + ticks: { stroke: "rgba(148,163,184,0.35)" }, + values: (_u, splits) => splits.map((v) => formatBps(v)), + size: 56, + }, + ], + scales: { + x: { time: true }, + }, + legend: { show: true }, + cursor: { + drag: { x: true, y: false }, + }, + }; + + plotRef.current?.destroy(); + plotRef.current = null; + + if (xs.length === 0) { + el.replaceChildren(); + return; + } + + plotRef.current = new uPlot(opts, [xs, inSeries, outSeries], el); + + const onResize = () => { + if (!plotRef.current || !rootRef.current) return; + plotRef.current.setSize({ + width: Math.max(320, rootRef.current.clientWidth), + height: Math.max(220, Math.min(420, Math.floor(rootRef.current.clientWidth * 0.32))), + }); + }; + window.addEventListener("resize", onResize); + return () => { + window.removeEventListener("resize", onResize); + plotRef.current?.destroy(); + plotRef.current = null; + }; + }, [points]); + + return ( +
+
+
+ {target ? target.ne_name || target.ne_ip || "—" : "—"} + {target?.ifname || "Select interface"} +
+
{rangeLabel}
+
+
+
+
Bandwidth
+
{formatBwLabel(kpi.bw)}
+
+
+
In
+
+ {formatBps(kpi.inBps)} + bit/s +
+
+
+
Out
+
+ {formatBps(kpi.outBps)} + bit/s +
+
+
+
Util In / Out
+
+ {kpi.inUtil.toFixed(1)}% + / + {kpi.outUtil.toFixed(1)}% +
+
+
+
+ {!points.length ?
No samples in range
: null} +
+
+ ); +} diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 6ad473c..323b8ef 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -29,6 +29,12 @@ import type { ConfigSyncTask, NeConfigSnapshotDetail, NeConfigSnapshotMeta, + PortTrafficDashboard, + PortTrafficDiscoverResponse, + PortTrafficSamples, + PortTrafficTarget, + PortTrafficTargetIn, + PortTrafficTask, } from "../types"; export const AUTH_TOKEN_KEY = "netx_access_token"; @@ -779,3 +785,56 @@ export const downloadNeConfigSnapshot = async ( URL.revokeObjectURL(url); } }; + +export const fetchPortTrafficDashboard = () => + apiGet("/v1/port-traffic/dashboard"); + +export const fetchPortTrafficTasks = (params: { page?: number; pageSize?: number }) => { + const p = new URLSearchParams(); + p.set("page", String(Math.max(1, Number(params.page || 1)))); + p.set("page_size", String(Math.max(1, Math.min(100, Number(params.pageSize || 20))))); + return apiGet<{ total: number; page: number; page_size: number; items: PortTrafficTask[] }>( + `/v1/port-traffic/tasks?${p.toString()}`, + ); +}; + +export const createPortTrafficTask = (body: { + title: string; + interval_sec?: number; + retention_days?: number; + concurrency?: number; + targets?: PortTrafficTargetIn[]; + start_now?: boolean; +}) => apiPost("/v1/port-traffic/tasks", body); + +export const startPortTrafficTask = (taskId: string) => + apiPost(`/v1/port-traffic/tasks/${encodeURIComponent(taskId)}/start`, {}); + +export const pausePortTrafficTask = (taskId: string) => + apiPost(`/v1/port-traffic/tasks/${encodeURIComponent(taskId)}/pause`, {}); + +export const stopPortTrafficTask = (taskId: string) => + apiPost(`/v1/port-traffic/tasks/${encodeURIComponent(taskId)}/stop`, {}); + +export const deletePortTrafficTask = (taskId: string) => + apiDelete<{ ok: boolean; id: string }>(`/v1/port-traffic/tasks/${encodeURIComponent(taskId)}`); + +export const fetchPortTrafficTargets = (taskId: string) => + apiGet<{ items: PortTrafficTarget[] }>( + `/v1/port-traffic/tasks/${encodeURIComponent(taskId)}/targets`, + ); + +export const discoverPortTrafficPorts = (body: { source: "managed" | "ume"; id: string }) => + apiPost("/v1/port-traffic/discover/ports", body); + +export const fetchPortTrafficSamples = (params: { + targetId: string; + from?: string; + to?: string; +}) => { + const p = new URLSearchParams(); + p.set("target_id", params.targetId); + if (params.from) p.set("from", params.from); + if (params.to) p.set("to", params.to); + return apiGet(`/v1/port-traffic/samples?${p.toString()}`); +}; diff --git a/web/src/types.ts b/web/src/types.ts index 0c869c3..f768dc0 100644 --- a/web/src/types.ts +++ b/web/src/types.ts @@ -507,3 +507,93 @@ export type NeConfigSnapshotDetail = NeConfigSnapshotMeta & { config_alt_text: string; }; +export type PortTrafficTask = { + id: string; + title: string; + status: string; + interval_sec: number; + retention_days: number; + concurrency: number; + collect_running: boolean; + target_count: number; + active_target_count: number; + last_collect_started_at?: string | null; + last_collect_ended_at?: string | null; + last_error: string; + created_at?: string | null; + updated_at?: string | null; +}; + +export type PortTrafficTarget = { + id: string; + task_id: string; + source: string; + target_id: string; + ne_name: string; + ne_ip: string; + vendor: string; + ifname: string; + if_description: string; + bw_bps: number; + status: string; + last_error: string; + last_sample_at?: string | null; + created_at?: string | null; +}; + +export type PortTrafficTargetIn = { + source: "managed" | "ume"; + target_id: string; + ne_name?: string; + ne_ip?: string; + vendor?: string; + ifname: string; + if_description?: string; + bw_bps?: number; +}; + +export type PortTrafficDiscoverPort = { + ifname: string; + attribute: string; + mode: string; + bw_raw: string; + bw_bps: number; + admin: string; + phy: string; + prot: string; + description: string; +}; + +export type PortTrafficDiscoverResponse = { + source: string; + id: string; + ne_name: string; + ne_ip: string; + vendor: string; + vendor_key: string; + ports: PortTrafficDiscoverPort[]; +}; + +export type PortTrafficSamplePoint = { + ts: string; + in_bps: number; + out_bps: number; + in_util_pct: number; + out_util_pct: number; + bw_bps: number; + rate_period_sec: number; +}; + +export type PortTrafficSamples = { + target: PortTrafficTarget; + points: PortTrafficSamplePoint[]; +}; + +export type PortTrafficDashboard = { + task_count: number; + running_task_count: number; + active_target_count: number; + sample_count_24h: number; + last_sample_at?: string | null; +}; +