From eea404384845aa34d8aeb86e26d671656230e261 Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 31 Jul 2026 15:10:33 +0800 Subject: [PATCH] Fix port traffic wall chart refresh and first-sample visibility. Trigger collect on start, normalize sample time bounds, and keep the uPlot series updating live. Co-authored-by: Cursor --- netx_api/port_traffic_router.py | 57 +++++++++- netx_api/port_traffic_service.py | 24 +++-- web/src/i18n/en.ts | 6 ++ web/src/i18n/zh.ts | 6 ++ web/src/index.css | 16 ++- web/src/pages/network/PortTrafficPage.tsx | 43 +++++++- web/src/pages/network/PortTrafficWall.tsx | 120 +++++++++++++++------- web/src/services/api.ts | 6 ++ web/src/utils/time.ts | 16 ++- 9 files changed, 240 insertions(+), 54 deletions(-) diff --git a/netx_api/port_traffic_router.py b/netx_api/port_traffic_router.py index 480b30b..33764fb 100644 --- a/netx_api/port_traffic_router.py +++ b/netx_api/port_traffic_router.py @@ -4,7 +4,7 @@ from __future__ import annotations from datetime import datetime -from fastapi import APIRouter, Depends, Query, Request +from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Query, Request from sqlalchemy.orm import Session from .auth_service import write_audit @@ -56,10 +56,15 @@ def api_list_tasks( @router.post("/tasks") def api_create_task( body: PortTrafficTaskCreate, + background_tasks: BackgroundTasks, request: Request, db: Session = Depends(get_db), ): out = create_task(db, body) + if body.start_now and out.status == "running": + from .port_traffic_runner import dispatch_collect + + background_tasks.add_task(dispatch_collect, out.id) uid, uname = _actor(request) write_audit( db, @@ -119,8 +124,16 @@ def api_delete_task(task_id: str, request: Request, db: Session = Depends(get_db @router.post("/tasks/{task_id}/start") -def api_start_task(task_id: str, request: Request, db: Session = Depends(get_db)): +def api_start_task( + task_id: str, + background_tasks: BackgroundTasks, + request: Request, + db: Session = Depends(get_db), +): out = set_task_status(db, task_id, "running") + from .port_traffic_runner import dispatch_collect + + background_tasks.add_task(dispatch_collect, task_id) uid, uname = _actor(request) write_audit( db, @@ -135,6 +148,46 @@ def api_start_task(task_id: str, request: Request, db: Session = Depends(get_db) return out.model_dump() +@router.post("/tasks/{task_id}/collect-now") +def api_collect_now( + task_id: str, + background_tasks: BackgroundTasks, + request: Request, + db: Session = Depends(get_db), +): + task = get_task(db, task_id) + if task.status not in ("running", "paused", "draft", "stopped"): + raise HTTPException(status_code=400, detail="invalid_status") + # Force due: clear last end so claim accepts, ensure running for this round. + from .models import PortTrafficTask + from .port_traffic_runner import dispatch_collect + + row = db.get(PortTrafficTask, task_id) + if not row: + raise HTTPException(status_code=404, detail="task_not_found") + if bool(row.collect_running): + return {"ok": True, "started": False, "reason": "already_collecting", **task.model_dump()} + if str(row.status) != "running": + row.status = "running" + row.last_collect_ended_at = None + row.updated_at = datetime.utcnow() + db.commit() + background_tasks.add_task(dispatch_collect, task_id) + uid, uname = _actor(request) + write_audit( + db, + action="port_traffic.task.collect_now", + actor_user_id=uid, + actor_username=uname, + method="POST", + path=f"/v1/port-traffic/tasks/{task_id}/collect-now", + status_code=200, + detail={"id": task_id}, + ) + out = get_task(db, task_id) + return {"ok": True, "started": True, **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") diff --git a/netx_api/port_traffic_service.py b/netx_api/port_traffic_service.py index f676a2a..b991b7d 100644 --- a/netx_api/port_traffic_service.py +++ b/netx_api/port_traffic_service.py @@ -3,7 +3,7 @@ from __future__ import annotations import logging -from datetime import datetime, timedelta +from datetime import datetime, timedelta, timezone from typing import Any from uuid import uuid4 @@ -204,11 +204,10 @@ def set_task_status(db: Session, task_id: str, status: str) -> PortTrafficTaskOu ) if active <= 0: raise HTTPException(status_code=400, detail="no_active_targets") + # Allow scheduler/collect-now to fire immediately after start/resume. + task.last_collect_ended_at = None 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) @@ -308,6 +307,15 @@ def discover_ports(db: Session, body: DiscoverPortsRequest) -> DiscoverPortsResp ) +def _as_naive_utc(value: datetime | None) -> datetime | None: + """Normalize query bounds to naive UTC (DB columns are naive utcnow).""" + if value is None: + return None + if value.tzinfo is None: + return value + return value.astimezone(timezone.utc).replace(tzinfo=None) + + def get_samples( db: Session, *, @@ -319,10 +327,10 @@ def get_samples( 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) + to_ts = _as_naive_utc(to_ts) or now + from_ts = _as_naive_utc(from_ts) or (to_ts - timedelta(hours=1)) + # Slight skew so just-written samples are not clipped by client clock. + to_ts = to_ts + timedelta(seconds=5) rows = ( db.query(PortTrafficSample) .filter( diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 1a670be..fa2d443 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -103,6 +103,12 @@ const en = { wallPort: "Interface", pickPort: "Select interface", range: "Range", + collectNow: "Collect now", + collectStarted: "Collect triggered", + collectBusy: "Collect already running", + targetError: "Last interface error", + waitSamples: "Samples appear after the task runs; or click Collect now. Chart auto-refreshes.", + waitAfterError: "Last collect failed — fix connectivity, then Collect now.", kpi: { tasks: "Tasks", running: "Running", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index 2f7d474..46785a9 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -103,6 +103,12 @@ const zh = { wallPort: "接口", pickPort: "选择接口", range: "时间范围", + collectNow: "立即采集", + collectStarted: "已触发采集", + collectBusy: "采集进行中", + targetError: "接口最近错误", + waitSamples: "任务运行后会按周期写入样点;也可点「立即采集」。图会每几秒自动刷新。", + waitAfterError: "采集失败,请检查连通性后点「立即采集」重试。", kpi: { tasks: "任务数", running: "运行中", diff --git a/web/src/index.css b/web/src/index.css index 9ca3d34..dd5b303 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -3083,7 +3083,13 @@ pre { } .pt-wall__chart { - min-height: 240px; + min-height: 260px; + width: 100%; +} + +.pt-wall__chart-wrap { + position: relative; + min-height: 260px; width: 100%; } @@ -3097,13 +3103,19 @@ pre { } .pt-wall__empty { + position: absolute; + inset: 0; display: flex; align-items: center; justify-content: center; - min-height: 240px; + padding: 16px 24px; + text-align: center; + line-height: 1.5; color: var(--pt-muted); + background: rgba(11, 18, 32, 0.72); border: 1px dashed rgba(148, 163, 184, 0.28); border-radius: 8px; + pointer-events: none; } @media (max-width: 900px) { diff --git a/web/src/pages/network/PortTrafficPage.tsx b/web/src/pages/network/PortTrafficPage.tsx index 98a3923..02d53bf 100644 --- a/web/src/pages/network/PortTrafficPage.tsx +++ b/web/src/pages/network/PortTrafficPage.tsx @@ -1,6 +1,7 @@ import { useEffect, useMemo, useState } from "react"; import { useMutation, useQuery, useQueryClient } from "@tanstack/react-query"; import { + collectPortTrafficNow, createPortTrafficTask, deletePortTrafficTask, discoverPortTrafficPorts, @@ -100,7 +101,8 @@ export function PortTrafficPage() { queryKey: queryKeys.portTrafficTargets(wallTaskId), queryFn: () => fetchPortTrafficTargets(wallTaskId), enabled: view === "wall" && Boolean(wallTaskId), - staleTime: 3000, + staleTime: 2000, + refetchInterval: POLL_MS, }); const samplesQuery = useQuery({ @@ -115,8 +117,12 @@ export function PortTrafficPage() { }); }, enabled: view === "wall" && Boolean(wallTargetId), - staleTime: 2000, - refetchInterval: POLL_MS, + staleTime: 1000, + refetchInterval: (q) => { + const n = q.state.data?.points?.length ?? 0; + // Faster while waiting for the first sample; then every 5s for live wall. + return n === 0 ? 2500 : POLL_MS; + }, }); const invalidateAll = () => { @@ -159,6 +165,16 @@ export function PortTrafficPage() { }, onError: (e: Error) => showError(e.message), }); + const collectMut = useMutation({ + mutationFn: collectPortTrafficNow, + onSuccess: (res) => { + showOk(res.started ? t("portTraffic.collectStarted") : t("portTraffic.collectBusy")); + invalidateAll(); + void queryClient.invalidateQueries({ queryKey: queryKeys.portTrafficTargets(wallTaskId) }); + void queryClient.invalidateQueries({ queryKey: ["portTrafficSamples"] }); + }, + onError: (e: Error) => showError(e.message), + }); const deleteMut = useMutation({ mutationFn: deletePortTrafficTask, onSuccess: () => { @@ -632,11 +648,32 @@ export function PortTrafficPage() { + + {selectedWallTarget?.last_error ? ( +

+ {t("portTraffic.targetError")}: {selectedWallTarget.last_error} +

+ ) : null} ) : null} diff --git a/web/src/pages/network/PortTrafficWall.tsx b/web/src/pages/network/PortTrafficWall.tsx index 2f3b061..c92f79d 100644 --- a/web/src/pages/network/PortTrafficWall.tsx +++ b/web/src/pages/network/PortTrafficWall.tsx @@ -2,13 +2,14 @@ import { useEffect, useMemo, useRef } from "react"; import uPlot from "uplot"; import "uplot/dist/uPlot.min.css"; import type { PortTrafficSamplePoint, PortTrafficTarget } from "../../types"; +import { formatSystemTime, parseApiTime } from "../../utils/time"; 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)}`; + return `${Math.round(n)}`; } function formatBwLabel(bps: number): string { @@ -16,14 +17,42 @@ function formatBwLabel(bps: number): string { return `${formatBps(bps)}bit/s`; } +function toSeries(points: PortTrafficSamplePoint[]): [number[], number[], number[]] { + const xs: number[] = []; + const ins: number[] = []; + const outs: number[] = []; + for (const p of points) { + const d = parseApiTime(p.ts); + if (!d) continue; + xs.push(Math.floor(d.getTime() / 1000)); + ins.push(Number(p.in_bps) || 0); + outs.push(Number(p.out_bps) || 0); + } + // Single sample: stretch a short segment so the line is visible. + if (xs.length === 1) { + xs.push(xs[0] + 60); + ins.push(ins[0]); + outs.push(outs[0]); + } + return [xs, ins, outs]; +} + +function chartSize(el: HTMLElement): { width: number; height: number } { + const width = Math.max(320, el.clientWidth || el.parentElement?.clientWidth || 800); + const height = Math.max(260, Math.min(440, Math.floor(width * 0.34))); + return { width, height }; +} + type Props = { target: PortTrafficTarget | null; points: PortTrafficSamplePoint[]; rangeLabel: string; + loading?: boolean; + hint?: string; }; -export function PortTrafficWall({ target, points, rangeLabel }: Props) { - const rootRef = useRef(null); +export function PortTrafficWall({ target, points, rangeLabel, loading, hint }: Props) { + const mountRef = useRef(null); const plotRef = useRef(null); const latest = points.length ? points[points.length - 1] : null; @@ -38,30 +67,30 @@ export function PortTrafficWall({ target, points, rangeLabel }: Props) { }; }, [latest, target]); + // Create plot once; keep resize listener for lifetime of mount. useEffect(() => { - const el = rootRef.current; + const el = mountRef.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 { width, height } = chartSize(el); 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))), + width, + height, series: [ {}, { - label: "In", + label: "In bit/s", stroke: "#2dd4bf", width: 2, - points: { show: false }, + fill: "rgba(45, 212, 191, 0.12)", + points: { show: true, size: 5, fill: "#2dd4bf" }, }, { - label: "Out", + label: "Out bit/s", stroke: "#f59e0b", width: 2, - points: { show: false }, + fill: "rgba(245, 158, 11, 0.10)", + points: { show: true, size: 5, fill: "#f59e0b" }, }, ], axes: [ @@ -75,34 +104,34 @@ export function PortTrafficWall({ target, points, rangeLabel }: Props) { 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, + size: 64, + label: "bit/s", + labelFont: "12px sans-serif", + labelSize: 14, }, ], scales: { x: { time: true }, + y: { + auto: true, + range: (_u, min, max) => { + if (!Number.isFinite(min) || !Number.isFinite(max)) return [0, 1]; + if (min === max) { + const pad = Math.max(10, Math.abs(max) * 0.25 || 10); + return [Math.max(0, min - pad), max + pad]; + } + return [Math.max(0, min * 0.95), max * 1.05]; + }, + }, }, legend: { show: true }, - cursor: { - drag: { x: true, y: false }, - }, + 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); - + plotRef.current = new uPlot(opts, [[], [], []], 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))), - }); + if (!plotRef.current || !mountRef.current) return; + plotRef.current.setSize(chartSize(mountRef.current)); }; window.addEventListener("resize", onResize); return () => { @@ -110,8 +139,19 @@ export function PortTrafficWall({ target, points, rangeLabel }: Props) { plotRef.current?.destroy(); plotRef.current = null; }; + }, []); + + // Push live samples into the existing chart. + useEffect(() => { + const plot = plotRef.current; + if (!plot) return; + const data = toSeries(points); + plot.setData(data); + if (mountRef.current) plot.setSize(chartSize(mountRef.current)); }, [points]); + const empty = points.length === 0; + return (
@@ -119,7 +159,10 @@ export function PortTrafficWall({ target, points, rangeLabel }: Props) { {target ? target.ne_name || target.ne_ip || "—" : "—"} {target?.ifname || "Select interface"}
-
{rangeLabel}
+
+ {rangeLabel} + {latest?.ts ? ` · last ${formatSystemTime(latest.ts)}` : ""} +
@@ -149,8 +192,13 @@ export function PortTrafficWall({ target, points, rangeLabel }: Props) {
-
- {!points.length ?
No samples in range
: null} +
+
+ {empty ? ( +
+ {loading ? "Loading samples…" : hint || "Waiting for samples…"} +
+ ) : null}
); diff --git a/web/src/services/api.ts b/web/src/services/api.ts index 323b8ef..803539a 100644 --- a/web/src/services/api.ts +++ b/web/src/services/api.ts @@ -816,6 +816,12 @@ export const pausePortTrafficTask = (taskId: string) => export const stopPortTrafficTask = (taskId: string) => apiPost(`/v1/port-traffic/tasks/${encodeURIComponent(taskId)}/stop`, {}); +export const collectPortTrafficNow = (taskId: string) => + apiPost( + `/v1/port-traffic/tasks/${encodeURIComponent(taskId)}/collect-now`, + {}, + ); + export const deletePortTrafficTask = (taskId: string) => apiDelete<{ ok: boolean; id: string }>(`/v1/port-traffic/tasks/${encodeURIComponent(taskId)}`); diff --git a/web/src/utils/time.ts b/web/src/utils/time.ts index 0464334..4734fdb 100644 --- a/web/src/utils/time.ts +++ b/web/src/utils/time.ts @@ -21,14 +21,24 @@ function normalizeTimeInput(value: string | number | Date, options?: FormatTimeO return `${text.replace(" ", "T")}Z`; } +/** Parse API timestamps; naive ISO strings are treated as UTC. */ +export function parseApiTime( + value: string | number | Date | null | undefined, + options?: FormatTimeOptions, +): Date | null { + if (value === null || value === undefined) return null; + const normalized = normalizeTimeInput(value, options); + const d = normalized instanceof Date ? normalized : new Date(normalized); + return Number.isNaN(d.getTime()) ? null : d; +} + export function formatSystemTime( value: string | number | Date | null | undefined, options?: FormatTimeOptions, ): string { if (value === null || value === undefined) return ""; - const normalized = normalizeTimeInput(value, options); - const d = normalized instanceof Date ? normalized : new Date(normalized); - if (Number.isNaN(d.getTime())) return String(value); + const d = parseApiTime(value, options); + if (!d) return String(value); return d.toLocaleString(undefined, { hour12: false, timeZone: systemTimeZone,