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}