Improve port-traffic logs and add manual NE rebind.

Fix unreadable collect-log contrast, surface errors on the Logs button, and let operators re-link a monitor to a chosen inventory NE without wiping samples.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-08-02 12:52:26 +08:00
parent a6c7267b47
commit 3f0787da11
9 changed files with 537 additions and 60 deletions

View file

@ -16,6 +16,7 @@ from .port_traffic_schemas import (
PortTrafficBoardPanelsPut,
PortTrafficBoardUpdate,
PortTrafficDeviceCreate,
PortTrafficDeviceRebind,
PortTrafficDeviceUpdate,
PortTrafficInterfacesPut,
PortTrafficReplacePortRequest,
@ -41,6 +42,7 @@ from .port_traffic_service import (
list_series,
list_targets,
put_interfaces,
rebind_device,
replace_series_port,
set_device_status,
update_device,
@ -236,6 +238,28 @@ def api_delete_device(device_id: str, request: Request, db: Session = Depends(ge
return out
@router.post("/devices/{device_id}/rebind")
def api_rebind_device(
device_id: str,
body: PortTrafficDeviceRebind,
request: Request,
db: Session = Depends(get_db),
):
out = rebind_device(db, device_id, body)
uid, uname = _actor(request)
write_audit(
db,
action="port_traffic.device.rebind",
actor_user_id=uid,
actor_username=uname,
method="POST",
path=f"/v1/port-traffic/devices/{device_id}/rebind",
status_code=200,
detail={"id": device_id, "ne_id": out.ne_id},
)
return out.model_dump()
@router.post("/devices/{device_id}/start")
def api_start_device(
device_id: str,

View file

@ -217,8 +217,12 @@ def _save_sample(target_row_id: str, parsed: Any, vendor_hint_bw: int = 0) -> No
db.close()
def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> int:
"""One CLI login for the device; run show per interface. Returns error count."""
def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> tuple[int, str]:
"""One CLI login for the device; run show per interface.
Returns (error_count, device_error). device_error is set when the whole
round fails for one shared reason (e.g. managed_ne_not_found).
"""
per_cmd = int(settings.ne_collect_read_timeout_sec or 120)
cap = int(settings.ne_collect_run_timeout_cap_sec or 600)
errors = 0
@ -227,7 +231,7 @@ def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> int
try:
device = db.get(PortTrafficDevice, device_id)
if not device:
return len(target_ids)
return len(target_ids), "device_not_found"
source = str(device.source or "").strip().lower()
ne_id = str(device.ne_id or "").strip()
vendor_hint = str(device.vendor or "")
@ -239,17 +243,17 @@ def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> int
else:
for tid in target_ids:
_set_target_error(tid, "invalid_source")
return len(target_ids)
return len(target_ids), "invalid_source"
except HTTPException as exc:
msg = str(exc.detail or "resolve_failed")[:1020]
for tid in target_ids:
_set_target_error(tid, msg)
return len(target_ids)
return len(target_ids), msg
except Exception as exc:
msg = _format_error(exc)
for tid in target_ids:
_set_target_error(tid, msg)
return len(target_ids)
return len(target_ids), msg
vendor = str(info.get("vendor") or vendor_hint or "")
device_type = str(info.get("device_type") or "")
@ -257,7 +261,7 @@ def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> int
if cmds is None:
for tid in target_ids:
_set_target_error(tid, "unsupported_vendor")
return len(target_ids)
return len(target_ids), "unsupported_vendor"
vendor_key = cmds.vendor_key
targets = (
@ -270,7 +274,7 @@ def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> int
db.close()
if not ifaces:
return 0
return 0, ""
budget = min(cap, per_cmd * max(1, len(ifaces)) + 90)
conn = None
@ -292,10 +296,11 @@ def _sample_targets_shared_session(device_id: str, target_ids: list[str]) -> int
_set_target_error(tid, msg)
errors = len(ifaces)
_log.exception("port_traffic session failed device=%s", device_id)
return errors, msg
finally:
if conn is not None:
close_netmiko_connection(conn)
return errors
return errors, ""
def dispatch_collect(device_id: str) -> int:
@ -304,13 +309,18 @@ def dispatch_collect(device_id: str) -> int:
if not target_ids:
return 0
session_err = ""
try:
errors = _sample_targets_shared_session(device_id, target_ids)
errors, session_err = _sample_targets_shared_session(device_id, target_ids)
except Exception:
errors = len(target_ids)
session_err = "collect_failed"
_log.exception("port_traffic collect failed device=%s", device_id)
finally:
err_msg = f"{errors}_target_errors" if errors else ""
if session_err and errors:
err_msg = session_err[:1020]
else:
err_msg = f"{errors}_target_errors" if errors else ""
_finish_collect_round(device_id, error=err_msg)
return len(target_ids)

View file

@ -51,6 +51,12 @@ class PortTrafficDeviceUpdate(BaseModel):
vendor: str | None = None
class PortTrafficDeviceRebind(BaseModel):
"""Re-link a monitor device to an explicitly chosen managed/UME NE; keeps samples."""
ne_id: str = Field(min_length=1, max_length=128)
class PortTrafficDeviceOut(BaseModel):
id: str
source: str

View file

@ -13,7 +13,15 @@ from sqlalchemy.orm import Session
from .cli_resolve import resolve_cli_target
from .config import settings
from .models import PortTrafficDevice, PortTrafficEvent, PortTrafficSample, PortTrafficSeries, PortTrafficTarget
from .models import (
ManagedNE,
PortTrafficDevice,
PortTrafficEvent,
PortTrafficSample,
PortTrafficSeries,
PortTrafficTarget,
UmeInventoryNE,
)
from .ne_session_factory import close_netmiko_connection, open_netmiko_connection
from .ne_netmiko import send_show_command
from .port_traffic_commands import commands_for_vendor
@ -28,6 +36,7 @@ from .port_traffic_schemas import (
PortTrafficDashboardOut,
PortTrafficDeviceCreate,
PortTrafficDeviceOut,
PortTrafficDeviceRebind,
PortTrafficDeviceUpdate,
PortTrafficEventOut,
PortTrafficEventsOut,
@ -257,6 +266,101 @@ def update_device(db: Session, device_id: str, body: PortTrafficDeviceUpdate) ->
return _device_out(db, device)
def rebind_device(
db: Session,
device_id: str,
body: PortTrafficDeviceRebind,
) -> PortTrafficDeviceOut:
"""Point a monitor device at an explicitly chosen inventory NE; keeps samples."""
device = db.get(PortTrafficDevice, device_id)
if not device:
raise HTTPException(status_code=404, detail="device_not_found")
if bool(device.collect_running):
raise HTTPException(status_code=409, detail="collect_running")
source = str(device.source or "").strip().lower() or "managed"
want_id = str(body.ne_id or "").strip()
if not want_id:
raise HTTPException(status_code=400, detail="ne_id_required")
new_id = ""
new_name = ""
new_ip = ""
new_vendor = ""
if source == "managed":
row = db.get(ManagedNE, want_id)
if not row:
raise HTTPException(status_code=404, detail="managed_ne_not_found")
new_id = str(row.id)
new_name = str(row.name or "")
new_ip = str(row.ip_address or "")
new_vendor = str(row.vendor or "")
elif source == "ume":
inv = db.get(UmeInventoryNE, want_id)
if not inv:
raise HTTPException(status_code=404, detail="ume_ne_not_found")
new_id = str(inv.ne_id)
new_ip = str(inv.ip_address or "")
new_name = str(inv.user_label or inv.ne_name or inv.host_name or new_ip or "").strip()
new_vendor = str(inv.vendor or "")
else:
raise HTTPException(status_code=400, detail="invalid_source")
if new_id == str(device.ne_id or ""):
# Already bound; clear stale errors so collect can resume.
device.last_error = ""
device.updated_at = _utcnow()
for tgt in db.query(PortTrafficTarget).filter(PortTrafficTarget.device_id == device_id).all():
if "managed_ne_not_found" in str(tgt.last_error or "") or "ume_ne_not_found" in str(
tgt.last_error or ""
):
tgt.last_error = ""
db.commit()
db.refresh(device)
return _device_out(db, device)
clash = (
db.query(PortTrafficDevice)
.filter(
PortTrafficDevice.source == source,
PortTrafficDevice.ne_id == new_id,
PortTrafficDevice.id != device_id,
)
.first()
)
if clash:
raise HTTPException(status_code=409, detail="device_already_monitored")
device.ne_id = new_id
if new_name:
device.ne_name = new_name
if new_ip:
device.ne_ip = new_ip
if new_vendor:
device.vendor = new_vendor
device.last_error = ""
device.updated_at = _utcnow()
for tgt in db.query(PortTrafficTarget).filter(PortTrafficTarget.device_id == device_id).all():
tgt.target_id = new_id
if new_name:
tgt.ne_name = new_name
if new_ip:
tgt.ne_ip = new_ip
if new_vendor:
tgt.vendor = new_vendor
if "managed_ne_not_found" in str(tgt.last_error or "") or "ume_ne_not_found" in str(
tgt.last_error or ""
):
tgt.last_error = ""
db.commit()
db.refresh(device)
_log.info("port_traffic rebind device=%s -> ne_id=%s ip=%s", device_id, new_id, new_ip)
return _device_out(db, device)
def delete_device(db: Session, device_id: str) -> dict[str, Any]:
device = db.get(PortTrafficDevice, device_id)
if not device: