diff --git a/netx_api/main.py b/netx_api/main.py
index 1ee43cf..fbd6c3d 100644
--- a/netx_api/main.py
+++ b/netx_api/main.py
@@ -32,6 +32,7 @@ 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
from .lldp_collect_router import router as lldp_collect_router
+from .ops_router import router as ops_router
from .importer import aggregate_alarms, import_alarm_excel, query_alarms
from .models import (
AiAnalyzeHistory,
@@ -141,6 +142,7 @@ app.include_router(port_traffic_router)
app.include_router(webcrt_router)
app.include_router(topology_router)
app.include_router(lldp_collect_router)
+app.include_router(ops_router)
parser_cfg = load_parser_config()
_UME_CLIENT_SINGLETON = UMEClient(
token_loader=lambda: load_shared_token(),
diff --git a/netx_api/models.py b/netx_api/models.py
index 37d8f4a..aa58f40 100644
--- a/netx_api/models.py
+++ b/netx_api/models.py
@@ -930,6 +930,7 @@ class PortTrafficPanel(Base):
range_hours: Mapped[int] = mapped_column(Integer, default=24)
baseline: Mapped[str] = mapped_column(String(16), default="off") # off|day|week|shift|custom
offset_hours: Mapped[int] = mapped_column(Integer, default=0)
+ ahead_hours: Mapped[int] = mapped_column(Integer, default=1) # extend chart past "now" for baseline peek
baseline_target_id: Mapped[str] = mapped_column(String(64), default="")
y_mode: Mapped[str] = mapped_column(String(16), default="auto") # auto|current|util
ord: Mapped[int] = mapped_column(Integer, default=0)
diff --git a/netx_api/ops_router.py b/netx_api/ops_router.py
new file mode 100644
index 0000000..72af5cb
--- /dev/null
+++ b/netx_api/ops_router.py
@@ -0,0 +1,23 @@
+"""Ops overview APIs (live task board)."""
+
+from __future__ import annotations
+
+from typing import Annotated, Any
+
+from fastapi import APIRouter, Depends
+from sqlalchemy.orm import Session
+
+from .auth_deps import AuthContext, require_user
+from .db import get_db
+from .ops_tasks_service import list_ops_tasks
+
+router = APIRouter(prefix="/v1/ops", tags=["ops"])
+
+
+@router.get("/tasks")
+def api_ops_tasks(
+ ctx: Annotated[AuthContext, Depends(require_user)],
+ db: Session = Depends(get_db),
+) -> dict[str, Any]:
+ _ = ctx
+ return list_ops_tasks(db)
diff --git a/netx_api/ops_tasks_service.py b/netx_api/ops_tasks_service.py
new file mode 100644
index 0000000..2d37739
--- /dev/null
+++ b/netx_api/ops_tasks_service.py
@@ -0,0 +1,472 @@
+"""Unified live task overview across NetX runners."""
+
+from __future__ import annotations
+
+from datetime import datetime, timedelta
+from typing import Any
+
+from sqlalchemy.orm import Session
+
+from .models import (
+ AuditLog,
+ ConfigSyncCycle,
+ ConfigSyncTask,
+ ManagedNE,
+ NeCollectionJob,
+ NeCollectionRun,
+ PortTrafficDevice,
+ TopoDiscoverJob,
+ UmeCliOverride,
+ UmeSyncJob,
+)
+
+
+def _iso(dt: datetime | None) -> str | None:
+ if dt is None:
+ return None
+ try:
+ return dt.isoformat()
+ except Exception:
+ return str(dt)
+
+
+def _item(
+ *,
+ kind: str,
+ id: str,
+ title: str,
+ status: str,
+ trigger: str = "",
+ actor: str = "",
+ started_at: datetime | None = None,
+ updated_at: datetime | None = None,
+ progress: str = "",
+ detail: str = "",
+ href: str = "",
+) -> dict[str, Any]:
+ return {
+ "kind": kind,
+ "id": id,
+ "title": title,
+ "status": status,
+ "trigger": trigger,
+ "actor": actor or "—",
+ "started_at": _iso(started_at),
+ "updated_at": _iso(updated_at),
+ "progress": progress,
+ "detail": detail,
+ "href": href,
+ }
+
+
+def _audit_actor_map(db: Session, *, hours: int = 72) -> dict[str, str]:
+ """Best-effort map of resource id → latest actor username from audit detail."""
+ since = datetime.utcnow() - timedelta(hours=max(1, hours))
+ rows = (
+ db.query(AuditLog)
+ .filter(AuditLog.ts >= since, AuditLog.status_code < 400)
+ .order_by(AuditLog.ts.desc())
+ .limit(800)
+ .all()
+ )
+ out: dict[str, str] = {}
+ for row in rows:
+ user = str(row.actor_username or "").strip()
+ if not user:
+ continue
+ detail = row.detail if isinstance(row.detail, dict) else {}
+ for key in ("id", "cycle_id", "job_id", "device_id", "session_id", "board_id"):
+ rid = str(detail.get(key) or "").strip()
+ if rid and rid not in out:
+ out[rid] = user
+ return out
+
+
+def _port_traffic_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]:
+ rows = (
+ db.query(PortTrafficDevice)
+ .filter(
+ (PortTrafficDevice.collect_running.is_(True))
+ | (PortTrafficDevice.status == "running")
+ )
+ .order_by(PortTrafficDevice.updated_at.desc())
+ .limit(200)
+ .all()
+ )
+ items: list[dict[str, Any]] = []
+ for row in rows:
+ did = str(row.id)
+ name = str(row.ne_name or row.ne_ip or did)
+ if bool(row.collect_running):
+ status = "collecting"
+ else:
+ status = str(row.status or "running")
+ actor = actors.get(did) or "scheduler"
+ items.append(
+ _item(
+ kind="port_traffic",
+ id=did,
+ title=f"端口流量 · {name}",
+ status=status,
+ trigger="schedule" if actor == "scheduler" else "manual",
+ actor=actor,
+ started_at=row.last_collect_started_at or row.updated_at,
+ updated_at=row.last_collect_ended_at or row.updated_at,
+ progress="",
+ detail=str(row.last_error or "")[:240],
+ href="/network/tasks/port-traffic",
+ )
+ )
+ return items
+
+
+def _config_sync_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]:
+ rows = (
+ db.query(ConfigSyncCycle)
+ .filter(ConfigSyncCycle.status.in_(("pending", "running", "paused")))
+ .order_by(ConfigSyncCycle.created_at.desc())
+ .limit(50)
+ .all()
+ )
+ items: list[dict[str, Any]] = []
+ for row in rows:
+ cid = str(row.id)
+ running = (
+ db.query(ConfigSyncTask)
+ .filter(ConfigSyncTask.cycle_id == cid, ConfigSyncTask.status == "running")
+ .count()
+ )
+ done = int(row.success_count or 0) + int(row.fail_count or 0) + int(row.skip_count or 0)
+ planned = int(row.planned_count or 0)
+ trigger = str(row.trigger_mode or "schedule")
+ actor = actors.get(cid) or ("scheduler" if trigger == "schedule" else "—")
+ items.append(
+ _item(
+ kind="config_sync",
+ id=cid,
+ title=f"配置同步 · {cid[:8]}",
+ status=str(row.status or "pending"),
+ trigger=trigger,
+ actor=actor,
+ started_at=row.started_at or row.created_at,
+ updated_at=row.ended_at or row.started_at or row.created_at,
+ progress=f"{done}/{planned}" + (f" · running {running}" if running else ""),
+ detail=str(row.error_message or "")[:240],
+ href="/network/tasks/config-sync",
+ )
+ )
+ return items
+
+
+def _collection_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]:
+ rows = (
+ db.query(NeCollectionJob)
+ .filter(NeCollectionJob.status.in_(("pending", "running", "paused")))
+ .order_by(NeCollectionJob.created_at.desc())
+ .limit(50)
+ .all()
+ )
+ items: list[dict[str, Any]] = []
+ for row in rows:
+ jid = str(row.id)
+ running = (
+ db.query(NeCollectionRun)
+ .filter(NeCollectionRun.job_id == jid, NeCollectionRun.status == "running")
+ .count()
+ )
+ title = str(row.title or "").strip() or f"采集任务 · {jid[:8]}"
+ actor = actors.get(jid) or "—"
+ items.append(
+ _item(
+ kind="ne_collect",
+ id=jid,
+ title=title,
+ status=str(row.status or "pending"),
+ trigger="manual",
+ actor=actor,
+ started_at=row.started_at or row.created_at,
+ updated_at=row.last_run_at or row.ended_at or row.started_at or row.created_at,
+ progress=(
+ f"ok {int(row.success_count or 0)} / fail {int(row.fail_count or 0)}"
+ f" / total {int(row.ne_count or 0)}"
+ + (f" · running {running}" if running else "")
+ ),
+ detail=str(row.error_message or "")[:240],
+ href="/network/tasks/collect",
+ )
+ )
+ return items
+
+
+def _lldp_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]:
+ rows = (
+ db.query(TopoDiscoverJob)
+ .filter(TopoDiscoverJob.status.in_(("pending", "running")))
+ .order_by(TopoDiscoverJob.created_at.desc())
+ .limit(20)
+ .all()
+ )
+ items: list[dict[str, Any]] = []
+ for row in rows:
+ jid = str(row.id)
+ trigger = str(row.trigger_mode or "manual")
+ actor = actors.get(jid) or (
+ "scheduler" if trigger == "schedule" else ("topology" if trigger == "topology" else "—")
+ )
+ done = int(row.done or 0)
+ total = int(row.total or 0)
+ items.append(
+ _item(
+ kind="lldp_discover",
+ id=jid,
+ title=f"LLDP 发现 · {trigger}",
+ status=str(row.status or "pending"),
+ trigger=trigger,
+ actor=actor,
+ started_at=row.started_at or row.created_at,
+ updated_at=row.updated_at or row.ended_at or row.started_at,
+ progress=f"{done}/{total}",
+ detail=str(row.error or "")[:240],
+ href="/network/topology/lldp",
+ )
+ )
+ return items
+
+
+def _ume_sync_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]:
+ """In-flight UME inventory/alarm sync jobs (manual or scheduled)."""
+ rows = (
+ db.query(UmeSyncJob)
+ .filter(UmeSyncJob.status == "running", UmeSyncJob.ended_at.is_(None))
+ .order_by(UmeSyncJob.started_at.desc())
+ .limit(20)
+ .all()
+ )
+ items: list[dict[str, Any]] = []
+ for row in rows:
+ jid = str(row.id)
+ domain = str(row.domain or "sync")
+ trigger = str(row.trigger_mode or "manual")
+ actor = actors.get(jid) or ("scheduler" if trigger == "schedule" else "—")
+ pulled = int(row.pulled_count or 0)
+ inserted = int(row.inserted_count or 0)
+ updated = int(row.updated_count or 0)
+ items.append(
+ _item(
+ kind="ume_sync",
+ id=jid,
+ title=f"UME 同步 · {domain}",
+ status="running",
+ trigger=trigger,
+ actor=actor,
+ started_at=row.started_at,
+ updated_at=row.started_at,
+ progress=f"pull {pulled} · +{inserted} ~{updated}",
+ detail=str(row.error_message or "")[:240],
+ href="/ume",
+ )
+ )
+ return items
+
+
+def _ne_connect_items(db: Session, actors: dict[str, str]) -> list[dict[str, Any]]:
+ """CLI connect-test pool work marked as testing on NE / UME override rows."""
+ items: list[dict[str, Any]] = []
+ managed = (
+ db.query(ManagedNE)
+ .filter(ManagedNE.connect_status == "testing")
+ .order_by(ManagedNE.updated_at.desc())
+ .limit(100)
+ .all()
+ )
+ for row in managed:
+ nid = str(row.id)
+ name = str(row.name or row.ip_address or nid[:8])
+ items.append(
+ _item(
+ kind="ne_connect",
+ id=f"managed:{nid}",
+ title=f"连通性测试 · {name}",
+ status="testing",
+ trigger="manual",
+ actor=actors.get(nid) or "—",
+ started_at=row.connect_tested_at or row.updated_at,
+ updated_at=row.updated_at,
+ progress=str(row.ip_address or ""),
+ detail=str(row.connect_message or "")[:240],
+ href="/network/devices",
+ )
+ )
+ overrides = (
+ db.query(UmeCliOverride)
+ .filter(UmeCliOverride.connect_status == "testing")
+ .order_by(UmeCliOverride.updated_at.desc())
+ .limit(100)
+ .all()
+ )
+ for row in overrides:
+ uid = str(row.ume_ne_id)
+ items.append(
+ _item(
+ kind="ne_connect",
+ id=f"ume:{uid}",
+ title=f"连通性测试 · UME {uid[:12]}",
+ status="testing",
+ trigger="manual",
+ actor=actors.get(uid) or "—",
+ started_at=row.connect_tested_at or row.updated_at,
+ updated_at=row.updated_at,
+ progress="",
+ detail=str(row.connect_message or "")[:240],
+ href="/ume",
+ )
+ )
+ return items
+
+
+def _ume_runtime_items() -> list[dict[str, Any]]:
+ try:
+ from .main import _list_runtime_tasks
+ except Exception:
+ return []
+ items: list[dict[str, Any]] = []
+ for row in _list_runtime_tasks():
+ task = str(row.get("task") or "")
+ status = str(row.get("status") or "unknown")
+ # Always show UME background loops so ops can see paused/idle too.
+ items.append(
+ _item(
+ kind="ume_runtime",
+ id=task,
+ title=f"UME · {task}",
+ status=status,
+ trigger="system",
+ actor="system",
+ started_at=None,
+ updated_at=None,
+ progress=str(row.get("interval_label") or ""),
+ detail=str(row.get("last_error") or "")[:240],
+ href="/ume",
+ )
+ )
+ # Fix updated_at from last_run_at string if present
+ if row.get("last_run_at") and items:
+ items[-1]["updated_at"] = row.get("last_run_at")
+ return items
+
+
+def _webcrt_items(actors: dict[str, str]) -> list[dict[str, Any]]:
+ try:
+ from .webcrt_service import list_sessions
+ except Exception:
+ return []
+ data = list_sessions()
+ items: list[dict[str, Any]] = []
+ now = datetime.utcnow().timestamp()
+ for row in data.get("items") or []:
+ sid = str(row.get("session_id") or "")
+ name = str(row.get("ne_name") or row.get("ne_ip") or sid[:8])
+ lifecycle = str(row.get("lifecycle") or row.get("state") or "unknown")
+ actor = actors.get(sid) or "—"
+ detail = ""
+ progress = ""
+ if lifecycle == "connecting":
+ elapsed_ms = row.get("elapsed_ms")
+ progress = f"{int(elapsed_ms)} ms" if isinstance(elapsed_ms, int) else "logging in"
+ detail = "authenticating"
+ elif lifecycle == "ready":
+ progress = "attached"
+ if row.get("connect_ms") is not None:
+ detail = f"connect {int(row['connect_ms'])} ms"
+ elif lifecycle == "detached":
+ progress = "detached"
+ deadline = row.get("detach_deadline")
+ if isinstance(deadline, (int, float)) and deadline > 0:
+ left = max(0, int(deadline - now))
+ detail = f"grace {left}s"
+ else:
+ detail = "awaiting reconnect / close"
+ elif lifecycle == "error":
+ progress = "error"
+ detail = str(row.get("connect_error") or "")[:240]
+ items.append(
+ _item(
+ kind="webcrt",
+ id=sid,
+ title=f"WebCRT · {name}",
+ status=lifecycle,
+ trigger="manual",
+ actor=actor,
+ started_at=None,
+ updated_at=None,
+ progress=progress,
+ detail=detail,
+ href="/webcrt",
+ )
+ )
+ if row.get("created_at"):
+ items[-1]["started_at"] = row.get("created_at")
+ if row.get("last_activity"):
+ items[-1]["updated_at"] = row.get("last_activity")
+ return items
+
+
+def list_ops_tasks(db: Session) -> dict[str, Any]:
+ actors = _audit_actor_map(db)
+ items: list[dict[str, Any]] = []
+ items.extend(_port_traffic_items(db, actors))
+ items.extend(_config_sync_items(db, actors))
+ items.extend(_collection_items(db, actors))
+ items.extend(_lldp_items(db, actors))
+ items.extend(_ume_sync_items(db, actors))
+ items.extend(_ne_connect_items(db, actors))
+ items.extend(_ume_runtime_items())
+ items.extend(_webcrt_items(actors))
+
+ by_kind: dict[str, int] = {}
+ by_status: dict[str, int] = {}
+ active = 0
+ for it in items:
+ kind = str(it.get("kind") or "")
+ status = str(it.get("status") or "")
+ by_kind[kind] = int(by_kind.get(kind, 0)) + 1
+ by_status[status] = int(by_status.get(status, 0)) + 1
+ if status in (
+ "running",
+ "collecting",
+ "pending",
+ "paused",
+ "connecting",
+ "ready",
+ "testing",
+ ):
+ active += 1
+
+ # Sort: active-ish first, then kind/title
+ rank = {
+ "collecting": 0,
+ "running": 1,
+ "testing": 2,
+ "connecting": 3,
+ "pending": 4,
+ "paused": 5,
+ "ready": 6,
+ "detached": 7,
+ "error": 8,
+ }
+ items.sort(
+ key=lambda x: (
+ rank.get(str(x.get("status") or ""), 50),
+ str(x.get("kind") or ""),
+ str(x.get("title") or ""),
+ )
+ )
+
+ return {
+ "generated_at": datetime.utcnow().isoformat() + "Z",
+ "total": len(items),
+ "active": active,
+ "by_kind": by_kind,
+ "by_status": by_status,
+ "items": items,
+ }
diff --git a/netx_api/port_traffic_board_service.py b/netx_api/port_traffic_board_service.py
index 4d961d4..0d910d0 100644
--- a/netx_api/port_traffic_board_service.py
+++ b/netx_api/port_traffic_board_service.py
@@ -38,6 +38,7 @@ def _panel_out(db: Session, row: PortTrafficPanel) -> PortTrafficPanelOut:
range_hours=int(row.range_hours or 24),
baseline=str(row.baseline or "off"),
offset_hours=int(row.offset_hours or 0),
+ ahead_hours=max(0, int(row.ahead_hours if getattr(row, "ahead_hours", None) is not None else 1)),
baseline_target_id=str(row.baseline_target_id or ""),
y_mode=str(row.y_mode or "auto"),
ord=int(row.ord or 0),
@@ -120,6 +121,7 @@ def _apply_panels(
range_hours=int(item.range_hours),
baseline=str(item.baseline or "off"),
offset_hours=int(item.offset_hours or 0),
+ ahead_hours=max(0, min(24, int(item.ahead_hours if item.ahead_hours is not None else 1))),
baseline_target_id=str(item.baseline_target_id or "").strip(),
y_mode=str(item.y_mode or "auto"),
ord=int(item.ord if item.ord is not None else i),
diff --git a/netx_api/port_traffic_migrate.py b/netx_api/port_traffic_migrate.py
index e500030..c2fe92a 100644
--- a/netx_api/port_traffic_migrate.py
+++ b/netx_api/port_traffic_migrate.py
@@ -154,6 +154,7 @@ def ensure_port_traffic_series_schema(conn) -> None:
range_hours INTEGER DEFAULT 24,
baseline VARCHAR(16) DEFAULT 'off',
offset_hours INTEGER DEFAULT 0,
+ ahead_hours INTEGER DEFAULT 1,
baseline_target_id VARCHAR(64) DEFAULT '',
y_mode VARCHAR(16) DEFAULT 'auto',
ord INTEGER DEFAULT 0,
@@ -164,6 +165,9 @@ def ensure_port_traffic_series_schema(conn) -> None:
)
"""
)
+ conn.exec_driver_sql(
+ "ALTER TABLE port_traffic_panel ADD COLUMN IF NOT EXISTS ahead_hours INTEGER DEFAULT 1"
+ )
conn.exec_driver_sql(
"CREATE INDEX IF NOT EXISTS ix_port_traffic_panel_board_id ON port_traffic_panel (board_id)"
)
diff --git a/netx_api/port_traffic_router.py b/netx_api/port_traffic_router.py
index cfea399..ac853cb 100644
--- a/netx_api/port_traffic_router.py
+++ b/netx_api/port_traffic_router.py
@@ -440,6 +440,12 @@ def api_compare(
range_hours: float = Query(default=24, ge=0.25, le=24 * 90),
baseline: str = Query(default="off", description="off|shift|day|week|custom"),
offset_hours: float | None = Query(default=None, ge=0.25, le=24 * 90),
+ ahead_hours: float = Query(
+ default=0,
+ ge=0,
+ le=24,
+ description="extend chart end past now so baseline can show upcoming trend",
+ ),
baseline_target_id: str | None = Query(
default=None,
description="optional mapped interface for baseline overlay",
@@ -454,5 +460,6 @@ def api_compare(
baseline=baseline,
offset_hours=offset_hours,
baseline_target_id=baseline_target_id,
+ ahead_hours=ahead_hours,
to_ts=to_ts,
).model_dump()
diff --git a/netx_api/port_traffic_schemas.py b/netx_api/port_traffic_schemas.py
index 95fbd0f..20e1568 100644
--- a/netx_api/port_traffic_schemas.py
+++ b/netx_api/port_traffic_schemas.py
@@ -160,21 +160,22 @@ class PortTrafficSamplePoint(BaseModel):
ts_raw: datetime | None = None
-class PortTrafficSamplesOut(BaseModel):
- target: PortTrafficTargetOut
- points: list[PortTrafficSamplePoint] = Field(default_factory=list)
-
-
class PortTrafficCompareMeta(BaseModel):
target_id: str
baseline: str = "off"
offset_hours: float = 0
range_hours: float = 24
+ ahead_hours: float = 0
current_target: PortTrafficTargetOut | None = None
baseline_target: PortTrafficTargetOut | None = None
baseline_target_id: str = ""
+class PortTrafficSamplesOut(BaseModel):
+ target: PortTrafficTargetOut
+ points: list[PortTrafficSamplePoint] = Field(default_factory=list)
+
+
class PortTrafficCompareOut(BaseModel):
meta: PortTrafficCompareMeta
current: list[PortTrafficSamplePoint] = Field(default_factory=list)
@@ -214,6 +215,7 @@ class PortTrafficPanelIn(BaseModel):
range_hours: int = Field(default=24, ge=1, le=24 * 90)
baseline: Literal["off", "day", "week", "shift", "custom"] = "off"
offset_hours: int = Field(default=0, ge=0, le=24 * 90)
+ ahead_hours: int = Field(default=1, ge=0, le=24)
baseline_target_id: str = ""
y_mode: Literal["auto", "current", "util"] = "auto"
ord: int = 0
@@ -229,6 +231,7 @@ class PortTrafficPanelOut(BaseModel):
range_hours: int = 24
baseline: str = "off"
offset_hours: int = 0
+ ahead_hours: int = 1
baseline_target_id: str = ""
y_mode: str = "auto"
ord: int = 0
diff --git a/netx_api/port_traffic_service.py b/netx_api/port_traffic_service.py
index 75a2c14..1f329a9 100644
--- a/netx_api/port_traffic_service.py
+++ b/netx_api/port_traffic_service.py
@@ -720,6 +720,7 @@ def compare_targets(
baseline: str = "off",
offset_hours: float | None = None,
baseline_target_id: str | None = None,
+ ahead_hours: float = 0,
to_ts: datetime | None = None,
) -> PortTrafficCompareOut:
target = db.get(PortTrafficTarget, target_row_id)
@@ -736,10 +737,14 @@ def compare_targets(
raise HTTPException(status_code=404, detail="baseline_target_not_found")
now = _utcnow()
- to_ts = _as_naive_utc(to_ts) or now
+ # Anchor = "now" (or explicit to). Lookback is from anchor; ahead extends past it
+ # so period compare can show baseline trend after the current clock time.
+ anchor = _as_naive_utc(to_ts) or now
range_h = max(0.25, float(range_hours or 24))
- from_ts = to_ts - timedelta(hours=range_h)
- to_q = to_ts + timedelta(seconds=5)
+ ahead_h = max(0.0, min(24.0, float(ahead_hours or 0)))
+ from_ts = anchor - timedelta(hours=range_h)
+ to_end = anchor + timedelta(hours=ahead_h)
+ to_q = to_end + timedelta(seconds=5)
current = _sample_points(
_query_target_samples(db, target_row_id=str(target.id), from_ts=from_ts, to_ts=to_q)
)
@@ -770,6 +775,7 @@ def compare_targets(
baseline=str(baseline or "off"),
offset_hours=float(off_h or 0),
range_hours=range_h,
+ ahead_hours=ahead_h,
current_target=_target_out(target),
baseline_target=_target_out(mapped) if mapped is not None else None,
baseline_target_id=str(mapped.id) if mapped is not None else "",
diff --git a/netx_api/topology_lldp.py b/netx_api/topology_lldp.py
index ff4cb06..5a52384 100644
--- a/netx_api/topology_lldp.py
+++ b/netx_api/topology_lldp.py
@@ -77,7 +77,7 @@ VENDOR_LLDP_PROFILES: dict[str, VendorLldpProfile] = {
key="cisco",
lldp_command="show lldp neighbors detail",
cdp_command="show cdp neighbors detail",
- notes="device_type cisco_*; use detail form for System Name / Port id.",
+ notes="device_type cisco_*; prefer detail for System Name / Port id / mgmt IP.",
),
"huawei": VendorLldpProfile(
key="huawei",
diff --git a/netx_api/webcrt_service.py b/netx_api/webcrt_service.py
index d2d89f4..6474c7a 100644
--- a/netx_api/webcrt_service.py
+++ b/netx_api/webcrt_service.py
@@ -471,8 +471,7 @@ class WebcrtSession:
detach_deadline: float | None = None
closed: bool = False
close_reason: str = ""
- # connecting | ready | error | closed
- state: str = "ready"
+ state: str = "connecting"
connect_error: str = ""
connect_started_at: float = field(default_factory=time.time)
connect_finished_at: float | None = None
@@ -1444,32 +1443,54 @@ def close_session(session_id: str, *, reason: str = "closed", client: str = "")
def list_sessions() -> dict[str, Any]:
with _sessions_lock:
- items = [
- {
- "session_id": s.session_id,
- "ne_id": s.ne_id,
- "ne_name": s.ne_name,
- "ne_ip": s.ne_ip,
- "protocol": s.protocol,
- "encoding": s.encoding,
- "keepalive_sec": int(s.keepalive_sec or 0),
- "state": s.state,
- "attached": s.attached,
- "created_at": datetime.fromtimestamp(s.created_at, tz=timezone.utc).isoformat(),
- "last_activity": datetime.fromtimestamp(s.last_activity, tz=timezone.utc).isoformat(),
- "bytes_in": s.bytes_in,
- "bytes_out": s.bytes_out,
- "queue_depth": s.out_queue.qsize(),
- "queue_dropped": getattr(s.out_queue, "dropped", 0),
- "connect_ms": (
- int((s.connect_finished_at - s.connect_started_at) * 1000)
- if s.connect_finished_at
- else None
- ),
- }
- for s in _sessions.values()
- if not s.closed
- ]
+ items = []
+ for s in _sessions.values():
+ if s.closed:
+ continue
+ state = str(s.state or "unknown")
+ attached = bool(s.attached)
+ # Lifecycle for ops UI: distinguish login vs live vs grace-period detach.
+ if state == "connecting":
+ lifecycle = "connecting"
+ elif state == "error":
+ lifecycle = "error"
+ elif state == "ready" and attached:
+ lifecycle = "ready"
+ elif state == "ready" and not attached:
+ lifecycle = "detached"
+ else:
+ lifecycle = state
+ elapsed_ms = None
+ if state == "connecting":
+ elapsed_ms = int(max(0.0, time.time() - float(s.connect_started_at or time.time())) * 1000)
+ items.append(
+ {
+ "session_id": s.session_id,
+ "ne_id": s.ne_id,
+ "ne_name": s.ne_name,
+ "ne_ip": s.ne_ip,
+ "protocol": s.protocol,
+ "encoding": s.encoding,
+ "keepalive_sec": int(s.keepalive_sec or 0),
+ "state": state,
+ "lifecycle": lifecycle,
+ "attached": attached,
+ "detach_deadline": s.detach_deadline,
+ "connect_error": str(s.connect_error or "")[:500],
+ "elapsed_ms": elapsed_ms,
+ "created_at": datetime.fromtimestamp(s.created_at, tz=timezone.utc).isoformat(),
+ "last_activity": datetime.fromtimestamp(s.last_activity, tz=timezone.utc).isoformat(),
+ "bytes_in": s.bytes_in,
+ "bytes_out": s.bytes_out,
+ "queue_depth": s.out_queue.qsize(),
+ "queue_dropped": getattr(s.out_queue, "dropped", 0),
+ "connect_ms": (
+ int((s.connect_finished_at - s.connect_started_at) * 1000)
+ if s.connect_finished_at
+ else None
+ ),
+ }
+ )
return {
"total": len(items),
"max_sessions": max(1, int(settings.webcrt_max_sessions or 20)),
diff --git a/tests/test_port_traffic_compare.py b/tests/test_port_traffic_compare.py
index 04fbac6..156eab9 100644
--- a/tests/test_port_traffic_compare.py
+++ b/tests/test_port_traffic_compare.py
@@ -236,6 +236,46 @@ class CompareTargetsTests(unittest.TestCase):
self.assertGreaterEqual(len(out.baseline), 2)
self.assertEqual(out.meta.offset_hours, 2.0)
+ def test_ahead_hours_extends_baseline_past_now(self):
+ # Yesterday sample 30min after "now" clock — only visible when ahead_hours > 0.
+ future_raw = self.now - timedelta(days=1) + timedelta(minutes=30)
+ self.db.add(
+ PortTrafficSample(
+ id=uuid4().hex,
+ target_row_id=self.target_id,
+ series_id=self.series_id,
+ ts=future_raw,
+ in_bps=777.0,
+ out_bps=888.0,
+ raw_ok=True,
+ )
+ )
+ self.db.commit()
+
+ without = compare_targets(
+ self.db,
+ target_row_id=self.target_id,
+ range_hours=2,
+ baseline="day",
+ ahead_hours=0,
+ to_ts=self.now,
+ )
+ self.assertFalse(any(abs((p.ts_raw or p.ts) - future_raw).total_seconds() < 1 for p in without.baseline))
+
+ with_ahead = compare_targets(
+ self.db,
+ target_row_id=self.target_id,
+ range_hours=2,
+ baseline="day",
+ ahead_hours=1,
+ to_ts=self.now,
+ )
+ self.assertEqual(with_ahead.meta.ahead_hours, 1.0)
+ hit = [p for p in with_ahead.baseline if p.in_bps == 777.0]
+ self.assertEqual(len(hit), 1)
+ self.assertGreater(hit[0].ts, self.now)
+ self.assertLessEqual(hit[0].ts, self.now + timedelta(hours=1, seconds=5))
+
if __name__ == "__main__":
unittest.main()
diff --git a/tests/test_topology.py b/tests/test_topology.py
index 96b62eb..2d3736a 100644
--- a/tests/test_topology.py
+++ b/tests/test_topology.py
@@ -34,12 +34,15 @@ from netx_api.topology_schemas import (
CISCO_LLDP_BRIEF = """
+R2#show lldp neighbors
Capability codes:
- (R) Router, (B) Bridge
+ (R) Router, (B) Bridge, (T) Telephone, (C) DOCSIS Cable Device
+ (W) WLAN Access Point, (P) Repeater, (S) Station, (O) Other
Device ID Local Intf Hold-time Capability Port ID
-R1 Gi0/0 120 R Gi0/1
-R3 Gi0/1 120 R Gi0/0
+r1 Gi0/1 120 B,R Ethernet1/0/1
+
+Total entries displayed: 1
"""
CISCO_LLDP_DETAIL = """
@@ -72,8 +75,16 @@ gei-0/1/0/1 0011.2233.4455 gei-0/1/0/2 R1
class LldpParserTests(unittest.TestCase):
def test_cisco_brief(self) -> None:
hits = lldp.parse_cisco_lldp(CISCO_LLDP_BRIEF)
- self.assertGreaterEqual(len(hits), 1)
- self.assertTrue(any(h.remote_name.upper().startswith("R1") for h in hits))
+ self.assertEqual(len(hits), 1)
+ self.assertEqual(hits[0].remote_name.lower(), "r1")
+ self.assertEqual(hits[0].local_port, "Gi0/1")
+ self.assertEqual(hits[0].remote_port, "Ethernet1/0/1")
+
+ def test_cisco_detail_fallback(self) -> None:
+ hits = lldp.parse_cisco_lldp(CISCO_LLDP_DETAIL)
+ self.assertEqual(len(hits), 1)
+ self.assertEqual(hits[0].remote_name.lower(), "r1")
+ self.assertEqual(hits[0].remote_ip, "192.168.0.1")
def test_huawei(self) -> None:
hits = lldp.parse_huawei_lldp(HUAWEI_LLDP)
@@ -83,7 +94,7 @@ class LldpParserTests(unittest.TestCase):
def test_pick_command_lldp_only(self) -> None:
cmd, tag = lldp.pick_neighbor_command(protocol="cdp", vendor="Cisco", device_type="cisco_ios")
self.assertEqual(tag, "lldp")
- self.assertIn("lldp", cmd.lower())
+ self.assertEqual(cmd, "show lldp neighbors detail")
class FabricTopologyTests(unittest.TestCase):
diff --git a/tests/test_webcrt.py b/tests/test_webcrt.py
index 3266133..d6c02d9 100644
--- a/tests/test_webcrt.py
+++ b/tests/test_webcrt.py
@@ -503,6 +503,8 @@ class WebcrtServiceTests(unittest.TestCase):
rows=24,
conn=conn, # type: ignore[arg-type]
)
+ sess.state = "ready"
+ sess.connect_finished_at = time.time()
sess.attached = True
with svc._sessions_lock:
svc._sessions["grace"] = sess
@@ -522,6 +524,42 @@ class WebcrtServiceTests(unittest.TestCase):
svc._reap_sessions()
self.assertIsNone(svc.get_session("grace"))
+ @patch.object(svc, "_audit")
+ def test_list_sessions_lifecycle_ready_vs_detached(self, _mock_audit: MagicMock) -> None:
+ conn = _FakeConn()
+ sess = svc.WebcrtSession(
+ session_id="life",
+ ne_id="ne1",
+ ne_name="lab",
+ ne_ip="1.2.3.4",
+ protocol="ssh",
+ cols=80,
+ rows=24,
+ conn=conn, # type: ignore[arg-type]
+ )
+ sess.state = "ready"
+ sess.attached = True
+ with svc._sessions_lock:
+ svc._sessions["life"] = sess
+ ready_rows = {r["session_id"]: r for r in svc.list_sessions()["items"]}
+ self.assertEqual(ready_rows["life"]["lifecycle"], "ready")
+ self.assertTrue(ready_rows["life"]["attached"])
+
+ svc.detach_session("life", grace_sec=120.0, attach_gen=0)
+ detached_rows = {r["session_id"]: r for r in svc.list_sessions()["items"]}
+ self.assertEqual(detached_rows["life"]["lifecycle"], "detached")
+ self.assertFalse(detached_rows["life"]["attached"])
+ self.assertIsNotNone(detached_rows["life"]["detach_deadline"])
+
+ sess.state = "connecting"
+ sess.attached = False
+ connecting_rows = {r["session_id"]: r for r in svc.list_sessions()["items"]}
+ self.assertEqual(connecting_rows["life"]["lifecycle"], "connecting")
+ self.assertIsInstance(connecting_rows["life"]["elapsed_ms"], int)
+
+ svc.close_session("life", reason="test")
+ self.assertEqual(svc.list_sessions()["total"], 0)
+
@patch.object(svc, "_audit")
def test_attach_timeout_reaper(self, _mock_audit: MagicMock) -> None:
conn = _FakeConn()
diff --git a/web/src/App.tsx b/web/src/App.tsx
index 6ec612c..ea4d1b5 100644
--- a/web/src/App.tsx
+++ b/web/src/App.tsx
@@ -20,6 +20,8 @@ import { PortTrafficBoardListPage } from "./pages/network/PortTrafficBoardListPa
import { PortTrafficWallPage } from "./pages/network/PortTrafficWallPage";
import { LoginPage } from "./pages/LoginPage";
import { UsersPage } from "./pages/UsersPage";
+import { AuditLayout } from "./pages/audit/AuditLayout";
+import { TaskOverviewPage } from "./pages/audit/TaskOverviewPage";
import { AuditPage } from "./pages/AuditPage";
import { ApiTokensPage } from "./pages/ApiTokensPage";
import { ForceChangePasswordPage } from "./pages/ForceChangePasswordPage";
@@ -117,7 +119,11 @@ function ProtectedApp() {
{isAdmin ? t("auth.auditHintAdmin") : t("auth.auditHintUser")}
diff --git a/web/src/pages/audit/AuditLayout.tsx b/web/src/pages/audit/AuditLayout.tsx new file mode 100644 index 0000000..debac18 --- /dev/null +++ b/web/src/pages/audit/AuditLayout.tsx @@ -0,0 +1,32 @@ +import { NavLink, Outlet } from "react-router-dom"; +import { AUDIT_NAV } from "../../config/auditNav"; +import { useI18n } from "../../i18n"; + +export function AuditLayout() { + const { t } = useI18n(); + + return ( +{t("audit.tasks.hint")}
+ +{t("common.refreshing")}
: null} + + {!items.length && !query.isLoading ? ( +{t("audit.tasks.empty")}
+| {t("audit.tasks.colKind")} | +{t("audit.tasks.colTitle")} | +{t("audit.tasks.colStatus")} | +{t("audit.tasks.colActor")} | +{t("audit.tasks.colTrigger")} | +{t("audit.tasks.colProgress")} | +{t("audit.tasks.colStarted")} | +{t("audit.tasks.colUpdated")} | +{t("audit.tasks.colLink")} | +
|---|---|---|---|---|---|---|---|---|
| {kindLabel(row.kind, t)} | +
+ {row.title}
+ {row.detail ? (
+
+ {row.detail}
+
+ ) : null}
+ |
+ + + {statusLabel(row.status, t)} + + | +{row.actor || "—"} | +{row.trigger || "—"} | +{row.progress || "—"} | +{formatSystemTime(row.started_at) || "—"} | +{formatSystemTime(row.updated_at) || "—"} | ++ {row.href ? ( + + {t("audit.tasks.open")} + + ) : ( + "—" + )} + | +
{t("portTraffic.retentionHint")}
) : null} + {editPanel.baseline_target_id ? ( ++ {t("portTraffic.mapBaselineHint")} +
+ ) : null}