diff --git a/go-backend/internal/http/handler/handler.go b/go-backend/internal/http/handler/handler.go index 7c87742..15c6506 100644 --- a/go-backend/internal/http/handler/handler.go +++ b/go-backend/internal/http/handler/handler.go @@ -42,6 +42,8 @@ type Handler struct { upgradeMu sync.Mutex pendingUpgradeRedeploy map[int64]struct{} + + qualityProber *tunnelQualityProber } type loginRequest struct { @@ -93,6 +95,7 @@ func New(repo *repo.Repository, jwtSecret string) *Handler { pendingUpgradeRedeploy: make(map[int64]struct{}), } h.healthCheck = health.NewChecker(repo, h.wsServer) + h.qualityProber = newTunnelQualityProber(h) h.wsServer.SetNodeOnlineHook(h.onNodeOnline) h.wsServer.SetNodeMetricHook(func(nodeID int64, info ws.SystemInfo) { metricInfo := metrics.SystemInfo{ @@ -229,6 +232,7 @@ func (h *Handler) Register(mux *http.ServeMux) { mux.HandleFunc("/api/v1/monitor/nodes/", h.monitorNodeMetricsHandler) mux.HandleFunc("/api/v1/monitor/nodes", h.monitorNodeListHandler) mux.HandleFunc("/api/v1/monitor/tunnels", h.monitorTunnelListHandler) + mux.HandleFunc("/api/v1/monitor/tunnels/quality", h.monitorTunnelQualityHandler) mux.HandleFunc("/api/v1/monitor/tunnels/", h.monitorTunnelMetrics) mux.HandleFunc("/api/v1/monitor/services", h.monitorServiceListHandler) mux.HandleFunc("/api/v1/monitor/services/create", h.monitorServiceCreate) diff --git a/go-backend/internal/http/handler/jobs.go b/go-backend/internal/http/handler/jobs.go index a169b4f..3225eeb 100644 --- a/go-backend/internal/http/handler/jobs.go +++ b/go-backend/internal/http/handler/jobs.go @@ -18,7 +18,7 @@ func (h *Handler) StartBackgroundJobs() { ctx, cancel := context.WithCancel(context.Background()) h.jobsCancel = cancel h.jobsStarted = true - h.jobsWG.Add(5) + h.jobsWG.Add(6) h.jobsMu.Unlock() go h.runHourlyStatsLoop(ctx) @@ -26,6 +26,7 @@ func (h *Handler) StartBackgroundJobs() { go h.runNodeRenewalCycleLoop(ctx) go h.runMetricsIngestion(ctx) go h.runHealthChecks(ctx) + go h.runTunnelQualityProber(ctx) } func (h *Handler) StopBackgroundJobs() { @@ -63,6 +64,13 @@ func (h *Handler) runHealthChecks(ctx context.Context) { } } +func (h *Handler) runTunnelQualityProber(ctx context.Context) { + defer h.jobsWG.Done() + if h.qualityProber != nil { + h.qualityProber.Start(ctx) + } +} + func (h *Handler) runHourlyStatsLoop(ctx context.Context) { defer h.jobsWG.Done() diff --git a/go-backend/internal/http/handler/monitoring.go b/go-backend/internal/http/handler/monitoring.go index 1123faf..99c309d 100644 --- a/go-backend/internal/http/handler/monitoring.go +++ b/go-backend/internal/http/handler/monitoring.go @@ -207,6 +207,98 @@ func (h *Handler) handleNodeMetricsLatest(w http.ResponseWriter, _ *http.Request response.WriteJSON(w, response.OK(metric)) } +func (h *Handler) monitorTunnelQualityHandler(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + response.WriteJSON(w, response.ErrDefault("请求失败")) + return + } + if !h.ensureMonitoringAccess(w, r) { + return + } + + // Try in-memory cache first + if h.qualityProber != nil { + items := h.qualityProber.GetAll() + if len(items) > 0 { + response.WriteJSON(w, response.OK(items)) + return + } + } + + // Fallback to database (latest per tunnel) + qualities, err := h.repo.GetLatestTunnelQualities() + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + + snapshots := make([]tunnelQualitySnapshot, 0, len(qualities)) + for _, q := range qualities { + snapshots = append(snapshots, tunnelQualitySnapshot{ + TunnelID: q.TunnelID, + EntryToExitLatency: q.EntryToExitLatency, + ExitToBingLatency: q.ExitToBingLatency, + EntryToExitLoss: q.EntryToExitLoss, + ExitToBingLoss: q.ExitToBingLoss, + Success: q.Success == 1, + ErrorMessage: q.ErrorMessage, + Timestamp: q.Timestamp, + }) + } + response.WriteJSON(w, response.OK(snapshots)) +} + +// monitorTunnelQualityHistory returns quality probe history for charting. +// GET /api/v1/monitor/tunnels/{id}/quality?start=...&end=... +// Mirrors monitorTunnelMetrics / monitorServiceResultsHandler pattern. +func (h *Handler) monitorTunnelQualityHistory(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + response.WriteJSON(w, response.ErrDefault("请求失败")) + return + } + if !h.ensureMonitoringAccess(w, r) { + return + } + + tunnelIDStr := extractPathParam(r.URL.Path, "/api/v1/monitor/tunnels/", "/quality") + tunnelID, err := strconv.ParseInt(tunnelIDStr, 10, 64) + if err != nil || tunnelID <= 0 { + response.WriteJSON(w, response.ErrDefault("无效的隧道ID")) + return + } + + now := time.Now().UnixMilli() + startMs := now - defaultMetricsRangeMs + endMs := now + + if s := r.URL.Query().Get("start"); s != "" { + if v, err := strconv.ParseInt(s, 10, 64); err == nil { + startMs = v + } + } + if e := r.URL.Query().Get("end"); e != "" { + if v, err := strconv.ParseInt(e, 10, 64); err == nil { + endMs = v + } + } + if startMs <= 0 || endMs <= 0 || endMs < startMs { + response.WriteJSON(w, response.ErrDefault("无效的时间范围")) + return + } + if endMs-startMs > maxMetricsRangeMs { + response.WriteJSON(w, response.ErrDefault("时间范围过大")) + return + } + + results, err := h.repo.GetTunnelQualityHistory(tunnelID, startMs, endMs) + if err != nil { + response.WriteJSON(w, response.Err(-2, err.Error())) + return + } + + response.WriteJSON(w, response.OK(results)) +} + func (h *Handler) monitorTunnelMetrics(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { response.WriteJSON(w, response.ErrDefault("请求失败")) @@ -216,7 +308,23 @@ func (h *Handler) monitorTunnelMetrics(w http.ResponseWriter, r *http.Request) { return } - tunnelIDStr := extractPathParam(r.URL.Path, "/api/v1/monitor/tunnels/", "/metrics") + path := r.URL.Path + prefix := "/api/v1/monitor/tunnels/" + if !strings.HasPrefix(path, prefix) { + response.WriteJSON(w, response.ErrDefault("无效的路径")) + return + } + + rest := strings.TrimPrefix(path, prefix) + + // Route: /api/v1/monitor/tunnels/{id}/quality + if strings.HasSuffix(rest, "/quality") { + h.monitorTunnelQualityHistory(w, r) + return + } + + // Route: /api/v1/monitor/tunnels/{id}/metrics (original) + tunnelIDStr := extractPathParam(path, prefix, "/metrics") tunnelID, err := strconv.ParseInt(tunnelIDStr, 10, 64) if err != nil || tunnelID <= 0 { response.WriteJSON(w, response.ErrDefault("无效的隧道ID")) diff --git a/go-backend/internal/http/handler/tunnel_quality_prober.go b/go-backend/internal/http/handler/tunnel_quality_prober.go new file mode 100644 index 0000000..b7b3194 --- /dev/null +++ b/go-backend/internal/http/handler/tunnel_quality_prober.go @@ -0,0 +1,332 @@ +package handler + +import ( + "context" + "log" + "sync" + "time" + + "go-backend/internal/store/model" +) + +const ( + tunnelQualityProbeInterval = 10 * time.Second + tunnelQualityProbeTimeout = 8 * time.Second + tunnelQualityPingTimeoutMs = 5000 + tunnelQualityRetention = 24 * time.Hour // keep 24h of history + tunnelQualityPruneInterval = 10 * time.Minute +) + +// tunnelQualitySnapshot is the in-memory latest probe result for a tunnel. +type tunnelQualitySnapshot struct { + TunnelID int64 `json:"tunnelId"` + EntryToExitLatency float64 `json:"entryToExitLatency"` + ExitToBingLatency float64 `json:"exitToBingLatency"` + EntryToExitLoss float64 `json:"entryToExitLoss"` + ExitToBingLoss float64 `json:"exitToBingLoss"` + Success bool `json:"success"` + ErrorMessage string `json:"errorMessage,omitempty"` + Timestamp int64 `json:"timestamp"` +} + +// tunnelQualityProber runs periodic TCP ping probes against all enabled tunnels. +// Design mirrors health.Checker: background goroutine with worker pool + scheduled cleanup. +type tunnelQualityProber struct { + handler *Handler + cache sync.Map // tunnelID (int64) → *tunnelQualitySnapshot + ctx context.Context + cancel context.CancelFunc + interval time.Duration + lastPrune int64 +} + +// newTunnelQualityProber creates a new prober (not yet running). +func newTunnelQualityProber(h *Handler) *tunnelQualityProber { + ctx, cancel := context.WithCancel(context.Background()) + return &tunnelQualityProber{ + handler: h, + ctx: ctx, + cancel: cancel, + interval: tunnelQualityProbeInterval, + } +} + +// Start launches the background probe loop (call from jobs.go). +func (p *tunnelQualityProber) Start(ctx context.Context) { + // Use the provided context so we stop with other background jobs. + p.ctx, p.cancel = context.WithCancel(ctx) + p.loop() +} + +// Stop halts the background probe loop. +func (p *tunnelQualityProber) Stop() { + p.cancel() +} + +// GetAll returns all cached quality snapshots (latest per tunnel). +func (p *tunnelQualityProber) GetAll() []tunnelQualitySnapshot { + var items []tunnelQualitySnapshot + p.cache.Range(func(_, value interface{}) bool { + if snap, ok := value.(*tunnelQualitySnapshot); ok { + items = append(items, *snap) + } + return true + }) + return items +} + +func (p *tunnelQualityProber) loop() { + // Initial delay to let the system boot up + select { + case <-time.After(5 * time.Second): + case <-p.ctx.Done(): + return + } + + // Run once immediately + p.probeAll() + + ticker := time.NewTicker(p.interval) + defer ticker.Stop() + + for { + select { + case <-p.ctx.Done(): + return + case <-ticker.C: + p.probeAll() + p.maybePrune() + } + } +} + +// maybePrune deletes old quality rows periodically (mirrors PruneServiceMonitorResults). +func (p *tunnelQualityProber) maybePrune() { + now := time.Now().UnixMilli() + if p.lastPrune > 0 && now-p.lastPrune < int64(tunnelQualityPruneInterval/time.Millisecond) { + return + } + p.lastPrune = now + + h := p.handler + if h == nil || h.repo == nil { + return + } + + cutoff := now - int64(tunnelQualityRetention/time.Millisecond) + if err := h.repo.PruneTunnelQualityResults(cutoff); err != nil { + log.Printf("tunnel_quality_prober: prune err=%v", err) + } +} + +func (p *tunnelQualityProber) probeAll() { + h := p.handler + if h == nil || h.repo == nil { + return + } + + tunnelIDs, err := h.repo.ListEnabledTunnelIDs() + if err != nil { + log.Printf("tunnel_quality_prober: list enabled tunnels err=%v", err) + return + } + if len(tunnelIDs) == 0 { + return + } + + // Probe tunnels concurrently with a worker limit + // (mirrors health.Checker worker pool pattern) + const maxWorkers = 4 + sem := make(chan struct{}, maxWorkers) + var wg sync.WaitGroup + + for _, tunnelID := range tunnelIDs { + select { + case <-p.ctx.Done(): + return + default: + } + + wg.Add(1) + sem <- struct{}{} + go func(tid int64) { + defer wg.Done() + defer func() { <-sem }() + p.probeTunnel(tid) + }(tunnelID) + } + wg.Wait() +} + +func (p *tunnelQualityProber) probeTunnel(tunnelID int64) { + h := p.handler + if h == nil || h.repo == nil { + return + } + + now := time.Now().UnixMilli() + snap := &tunnelQualitySnapshot{ + TunnelID: tunnelID, + Timestamp: now, + } + + // Get tunnel chain info + tunnel, err := h.getTunnelRecord(tunnelID) + if err != nil { + snap.ErrorMessage = "隧道不存在" + p.storeResult(snap) + return + } + + chainRows, err := h.listChainNodesForTunnel(tunnelID) + if err != nil || len(chainRows) == 0 { + snap.ErrorMessage = "隧道配置不完整" + p.storeResult(snap) + return + } + + ipPreference := h.repo.GetTunnelIPPreference(tunnelID) + inNodes, _, outNodes := splitChainNodeGroups(chainRows) + + options := diagnosisExecOptions{ + commandTimeout: tunnelQualityProbeTimeout, + pingTimeoutMS: tunnelQualityPingTimeoutMs, + timeoutMessage: "探测超时", + } + + switch tunnel.Type { + case 1: + // Port forwarding: entry → Bing only + if len(inNodes) > 0 { + lat, loss, err := p.tcpPingNode(inNodes[0].NodeID, "www.bing.com", 443, options) + if err == nil { + snap.ExitToBingLatency = lat + snap.ExitToBingLoss = loss + snap.Success = true + } else { + snap.ErrorMessage = err.Error() + } + } + case 2: + // Tunnel forwarding: entry → exit + exit → Bing + probeOK := true + + if len(inNodes) > 0 && len(outNodes) > 0 { + // Entry → Exit + targetNode, nodeErr := h.getNodeRecord(outNodes[0].NodeID) + if nodeErr == nil && targetNode != nil { + fromNode, _ := h.getNodeRecord(inNodes[0].NodeID) + targetIP, targetPort, resolveErr := resolveChainProbeTarget(fromNode, targetNode, outNodes[0].Port, ipPreference, outNodes[0].ConnectIP) + if resolveErr == nil { + lat, loss, err := p.tcpPingNode(inNodes[0].NodeID, targetIP, targetPort, options) + if err == nil { + snap.EntryToExitLatency = lat + snap.EntryToExitLoss = loss + } else { + snap.EntryToExitLatency = -1 + snap.EntryToExitLoss = 100 + probeOK = false + } + } else { + snap.ErrorMessage = resolveErr.Error() + probeOK = false + } + } else { + snap.ErrorMessage = "出口节点不可用" + probeOK = false + } + } + + // Exit → Bing + if len(outNodes) > 0 { + lat, loss, err := p.tcpPingNode(outNodes[0].NodeID, "www.bing.com", 443, options) + if err == nil { + snap.ExitToBingLatency = lat + snap.ExitToBingLoss = loss + } else { + if snap.ErrorMessage == "" { + snap.ErrorMessage = err.Error() + } + probeOK = false + } + } + + snap.Success = probeOK + default: + // Unknown type: entry → Bing + if len(inNodes) > 0 { + lat, loss, err := p.tcpPingNode(inNodes[0].NodeID, "www.bing.com", 443, options) + if err == nil { + snap.ExitToBingLatency = lat + snap.ExitToBingLoss = loss + snap.Success = true + } else { + snap.ErrorMessage = err.Error() + } + } + } + + p.storeResult(snap) +} + +func (p *tunnelQualityProber) tcpPingNode(nodeID int64, ip string, port int, options diagnosisExecOptions) (latency float64, loss float64, err error) { + h := p.handler + if h == nil { + return 0, 100, nil + } + + node, nodeErr := h.getNodeRecord(nodeID) + if nodeErr != nil { + return 0, 100, nodeErr + } + + var pingData map[string]interface{} + var pingErr error + if node != nil && node.IsRemote == 1 { + pingData, pingErr = h.tcpPingViaRemoteNode(node, ip, port, options) + } else { + pingData, pingErr = h.tcpPingViaNode(nodeID, ip, port, options) + } + if pingErr != nil { + return 0, 100, pingErr + } + + avgTime := asFloat(pingData["averageTime"], 0) + packetLoss := asFloat(pingData["packetLoss"], 100) + + return avgTime, packetLoss, nil +} + +func (p *tunnelQualityProber) storeResult(snap *tunnelQualitySnapshot) { + if snap == nil { + return + } + + // Update in-memory cache (latest per tunnel) + p.cache.Store(snap.TunnelID, snap) + + // Persist to database (history) + h := p.handler + if h == nil || h.repo == nil { + return + } + + successInt := 0 + if snap.Success { + successInt = 1 + } + + q := &model.TunnelQuality{ + TunnelID: snap.TunnelID, + EntryToExitLatency: snap.EntryToExitLatency, + ExitToBingLatency: snap.ExitToBingLatency, + EntryToExitLoss: snap.EntryToExitLoss, + ExitToBingLoss: snap.ExitToBingLoss, + Success: successInt, + ErrorMessage: snap.ErrorMessage, + Timestamp: snap.Timestamp, + } + if err := h.repo.InsertTunnelQuality(q); err != nil { + log.Printf("tunnel_quality_prober: insert db err=%v tunnel_id=%d", err, snap.TunnelID) + } +} diff --git a/go-backend/internal/store/model/model.go b/go-backend/internal/store/model/model.go index 0d58c80..54f59a1 100644 --- a/go-backend/internal/store/model/model.go +++ b/go-backend/internal/store/model/model.go @@ -714,3 +714,20 @@ type ServiceMonitorResult struct { } func (ServiceMonitorResult) TableName() string { return "service_monitor_result" } + +// TunnelQuality stores periodic probe results for a tunnel. +// Unlike the old upsert model, rows accumulate for history/charting. +// Old rows are pruned periodically (default: keep 24h). +type TunnelQuality struct { + ID int64 `gorm:"primaryKey;autoIncrement" json:"id"` + TunnelID int64 `gorm:"column:tunnel_id;not null;index:idx_tunnel_quality_tunnel_time,priority:1" json:"tunnelId"` + EntryToExitLatency float64 `gorm:"column:entry_to_exit_latency" json:"entryToExitLatency"` + ExitToBingLatency float64 `gorm:"column:exit_to_bing_latency" json:"exitToBingLatency"` + EntryToExitLoss float64 `gorm:"column:entry_to_exit_loss" json:"entryToExitLoss"` + ExitToBingLoss float64 `gorm:"column:exit_to_bing_loss" json:"exitToBingLoss"` + Success int `gorm:"not null;default:1" json:"success"` + ErrorMessage string `gorm:"column:error_message;type:text" json:"errorMessage,omitempty"` + Timestamp int64 `gorm:"not null;index:idx_tunnel_quality_tunnel_time,priority:2;index:idx_tunnel_quality_time" json:"timestamp"` +} + +func (TunnelQuality) TableName() string { return "tunnel_quality" } diff --git a/go-backend/internal/store/repo/repository.go b/go-backend/internal/store/repo/repository.go index fcfc366..a0a6de3 100644 --- a/go-backend/internal/store/repo/repository.go +++ b/go-backend/internal/store/repo/repository.go @@ -51,6 +51,7 @@ type NodeMetric = model.NodeMetric type TunnelMetric = model.TunnelMetric type ServiceMonitor = model.ServiceMonitor type ServiceMonitorResult = model.ServiceMonitorResult +type TunnelQuality = model.TunnelQuality // ─── Repository ────────────────────────────────────────────────────── @@ -191,6 +192,7 @@ func autoMigrateAll(db *gorm.DB) error { &model.TunnelMetric{}, &model.ServiceMonitor{}, &model.ServiceMonitorResult{}, + &model.TunnelQuality{}, } if db.Dialector.Name() != "sqlite" { diff --git a/go-backend/internal/store/repo/repository_tunnel_quality.go b/go-backend/internal/store/repo/repository_tunnel_quality.go new file mode 100644 index 0000000..de5f4bb --- /dev/null +++ b/go-backend/internal/store/repo/repository_tunnel_quality.go @@ -0,0 +1,98 @@ +package repo + +import ( + "errors" + + "go-backend/internal/store/model" +) + +// InsertTunnelQuality appends a tunnel quality probe result. +// (Follows the same pattern as InsertServiceMonitorResult.) +func (r *Repository) InsertTunnelQuality(q *model.TunnelQuality) error { + if r == nil || r.db == nil { + return errors.New("repository not initialized") + } + if q == nil || q.TunnelID <= 0 { + return nil + } + return r.db.Create(q).Error +} + +// GetTunnelQualityHistory returns quality probe results for a tunnel +// within a time range, ordered by timestamp ascending. +// (Mirrors GetServiceMonitorResults pattern.) +func (r *Repository) GetTunnelQualityHistory(tunnelID int64, startMs, endMs int64) ([]model.TunnelQuality, error) { + if r == nil || r.db == nil { + return nil, errors.New("repository not initialized") + } + var results []model.TunnelQuality + err := r.db.Where("tunnel_id = ? AND timestamp >= ? AND timestamp <= ?", tunnelID, startMs, endMs). + Order("timestamp ASC"). + Find(&results).Error + return results, err +} + +// GetLatestTunnelQualities returns the newest quality result per tunnel_id. +// (Mirrors GetLatestServiceMonitorResults pattern.) +func (r *Repository) GetLatestTunnelQualities() ([]model.TunnelQuality, error) { + if r == nil || r.db == nil { + return nil, nil + } + + var results []model.TunnelQuality + + // Use window function (works on modern SQLite 3.25+ and PostgreSQL). + q := ` + SELECT id, tunnel_id, entry_to_exit_latency, exit_to_bing_latency, + entry_to_exit_loss, exit_to_bing_loss, success, error_message, timestamp + FROM ( + SELECT *, ROW_NUMBER() OVER (PARTITION BY tunnel_id ORDER BY timestamp DESC, id DESC) AS rn + FROM tunnel_quality + ) t + WHERE rn = 1 + ORDER BY tunnel_id ASC + ` + if err := r.db.Raw(q).Scan(&results).Error; err == nil { + return results, nil + } + + // Fallback for older SQLite + results = nil + err := r.db.Order("timestamp DESC, id DESC").Limit(5000).Find(&results).Error + if err != nil { + return nil, err + } + + seen := make(map[int64]struct{}, len(results)) + out := make([]model.TunnelQuality, 0, len(results)) + for _, row := range results { + if row.TunnelID <= 0 { + continue + } + if _, ok := seen[row.TunnelID]; ok { + continue + } + seen[row.TunnelID] = struct{}{} + out = append(out, row) + } + return out, nil +} + +// PruneTunnelQualityResults deletes quality results older than the given timestamp. +// (Mirrors PruneServiceMonitorResults pattern.) +func (r *Repository) PruneTunnelQualityResults(olderThanMs int64) error { + if r == nil || r.db == nil { + return nil + } + return r.db.Where("timestamp < ?", olderThanMs).Delete(&model.TunnelQuality{}).Error +} + +// ListEnabledTunnelIDs returns IDs of all tunnels with status=1. +func (r *Repository) ListEnabledTunnelIDs() ([]int64, error) { + if r == nil || r.db == nil { + return nil, errors.New("repository not initialized") + } + var ids []int64 + err := r.db.Model(&model.Tunnel{}).Where("status = ?", 1).Pluck("id", &ids).Error + return ids, err +} diff --git a/plans/055-tunnel-quality-periodic-probing.md b/plans/055-tunnel-quality-periodic-probing.md new file mode 100644 index 0000000..1add81d --- /dev/null +++ b/plans/055-tunnel-quality-periodic-probing.md @@ -0,0 +1,31 @@ +# 055 - 隧道质量定时探测 + 实时展示 + 历史图表 + +## 背景 +当前隧道质量检测是手动触发的:用户点击"诊断"按钮 → 后端调用节点 TcpPing → 返回结果。 +需求:改为**后端定时(每10秒)自动探测**所有启用隧道的质量(入口→出口延迟、出口→Bing延迟), +结果保留历史(24h),前端隧道 Tab 实时展示 + 图表历史趋势。 + +## 设计原则:与服务监控复用 + +| 复用点 | 服务监控 | 隧道质量 | +|--------|---------|---------| +| 调度方式 | `health.Checker.Start(ctx)` via `jobs.go` | `tunnelQualityProber.Start(ctx)` via `jobs.go` | +| 存储模式 | `service_monitor_result` (history, insert) | `tunnel_quality` (history, insert) | +| 清理方式 | `PruneServiceMonitorResults(olderThanMs)` | `PruneTunnelQualityResults(olderThanMs)` | +| 最新查询 | `GetLatestServiceMonitorResults()` (window func) | `GetLatestTunnelQualities()` (window func) | +| 历史查询 | `GetServiceMonitorResults(id, limit)` | `GetTunnelQualityHistory(id, start, end)` | +| API 模式 | `GET /monitor/services/{id}/results` | `GET /monitor/tunnels/{id}/quality` | +| 前端图表 | Recharts LineChart (延迟趋势) | Recharts LineChart (同样模式) | + +## 任务清单 + +- [x] 1. `TunnelQuality` model 改为历史存储(composite index, 非 unique) +- [x] 2. Repo 改为 insert(非 upsert),复用服务监控的查询模式 +- [x] 3. 添加 `PruneTunnelQualityResults` + `GetLatestTunnelQualities` + `GetTunnelQualityHistory` +- [x] 4. Prober 生命周期集成到 `jobs.go`(与 healthCheck 同级) +- [x] 5. Prober 添加 24h 清理周期 +- [x] 6. 添加 API `GET /monitor/tunnels/{id}/quality` 返回历史 +- [x] 7. 前端添加 `getMonitorTunnelQualityHistory()` API +- [x] 8. 前端详情页添加质量趋势图表(复用服务监控图表组件模式) +- [x] 9. Go 编译 + 测试通过 +- [x] 10. TypeScript 编译通过 diff --git a/vite-frontend/src/api/index.ts b/vite-frontend/src/api/index.ts index 3608a70..ee824a0 100644 --- a/vite-frontend/src/api/index.ts +++ b/vite-frontend/src/api/index.ts @@ -40,6 +40,7 @@ import type { MonitorTunnelApiItem, MonitorPermissionApiItem, MonitorAccessApiData, + TunnelQualityApiItem, } from "./types"; import axios from "axios"; @@ -451,6 +452,25 @@ export const getTunnelMetrics = ( export const getMonitorTunnels = () => Network.get("/monitor/tunnels"); +export const getMonitorTunnelQuality = () => + Network.get("/monitor/tunnels/quality"); + +export const getMonitorTunnelQualityHistory = ( + tunnelId: number, + start?: number, + end?: number, +) => { + const params: Record = {}; + + if (start) params.start = String(start); + if (end) params.end = String(end); + + return Network.get( + `/monitor/tunnels/${tunnelId}/quality`, + params, + ); +}; + export const getServiceMonitorList = () => Network.get("/monitor/services"); diff --git a/vite-frontend/src/api/types.ts b/vite-frontend/src/api/types.ts index d40134a..25db1b7 100644 --- a/vite-frontend/src/api/types.ts +++ b/vite-frontend/src/api/types.ts @@ -487,3 +487,14 @@ export interface MonitorAccessApiData { allowed: boolean; reason?: string; } + +export interface TunnelQualityApiItem { + tunnelId: number; + entryToExitLatency: number; + exitToBingLatency: number; + entryToExitLoss: number; + exitToBingLoss: number; + success: boolean; + errorMessage?: string; + timestamp: number; +} diff --git a/vite-frontend/src/pages/node/tunnel-monitor-view.tsx b/vite-frontend/src/pages/node/tunnel-monitor-view.tsx index 93135cb..aeaafb3 100644 --- a/vite-frontend/src/pages/node/tunnel-monitor-view.tsx +++ b/vite-frontend/src/pages/node/tunnel-monitor-view.tsx @@ -1,7 +1,7 @@ import type { MonitorTunnelApiItem, TunnelMetricApiItem, - TunnelDiagnosisApiItem, + TunnelQualityApiItem, } from "@/api/types"; import { useCallback, useEffect, useMemo, useRef, useState } from "react"; @@ -13,6 +13,7 @@ import { CartesianGrid, Tooltip, ResponsiveContainer, + Legend, } from "recharts"; import { RefreshCw, @@ -23,16 +24,15 @@ import { ArrowRightLeft, Wifi, WifiOff, - Stethoscope, } from "lucide-react"; import toast from "react-hot-toast"; import { getMonitorTunnels, getTunnelMetrics, - diagnoseTunnel, + getMonitorTunnelQuality, + getMonitorTunnelQualityHistory, } from "@/api"; -import { diagnoseTunnelStream } from "@/api/diagnosis-stream"; import { getDiagnosisQualityDisplay } from "@/pages/tunnel/diagnosis"; import { Button } from "@/shadcn-bridge/heroui/button"; import { Card, CardBody, CardHeader } from "@/shadcn-bridge/heroui/card"; @@ -51,18 +51,7 @@ interface TunnelMonitorViewProps { viewMode?: "list" | "grid"; } -const METRICS_MAX_ROWS = 5000; - -interface TunnelQuality { - loading: boolean; - entryToExitLatency?: number; - exitToBingLatency?: number; - entryToExitLoss?: number; - exitToBingLoss?: number; - results?: TunnelDiagnosisApiItem[]; - timestamp?: number; - error?: string; -} +const QUALITY_POLL_INTERVAL = 10_000; // 10 seconds const formatTimestamp = (ts: number, rangeMs?: number): string => { const date = new Date(ts); @@ -80,6 +69,7 @@ const formatTimestamp = (ts: number, rangeMs?: number): string => { return date.toLocaleTimeString("zh-CN", { hour: "2-digit", minute: "2-digit", + second: "2-digit", }); }; @@ -93,32 +83,60 @@ const formatBytes = (bytes: number): string => { return `${parseFloat((bytes / Math.pow(k, i)).toFixed(2))} ${sizes[i]}`; }; +/** Render a colored latency value with appropriate visual cue */ +function LatencyDisplay({ value, loading }: { value?: number; loading?: boolean }) { + if (loading) { + return ; + } + if (value === undefined || value < 0) { + return -; + } + const ms = value.toFixed(0); + let colorClass = "text-success"; + if (value > 200) colorClass = "text-danger"; + else if (value > 100) colorClass = "text-warning"; + else if (value > 50) colorClass = "text-primary"; + + return {ms}ms; +} + +/** Animated pulse dot for live status */ +function LiveDot() { + return ( + + + + + ); +} + export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) { const [tunnels, setTunnels] = useState([]); const [tunnelsLoading, setTunnelsLoading] = useState(false); const [tunnelsError, setTunnelsError] = useState(null); const [accessDenied, setAccessDenied] = useState(null); + // Quality data from backend periodic probing (latest per tunnel) + const [qualityMap, setQualityMap] = useState>({}); + const [qualityLoading, setQualityLoading] = useState(false); + const qualityTimerRef = useRef(null); + // Detail view state const [detailTunnelId, setDetailTunnelId] = useState(null); + + // Quality history for chart (mirrors service monitor results) + const [qualityHistory, setQualityHistory] = useState([]); + const [qualityHistoryLoading, setQualityHistoryLoading] = useState(false); + const [qualityHistoryError, setQualityHistoryError] = useState(null); + const [qualityRangeMs, setQualityRangeMs] = useState(60 * 60 * 1000); + + // Tunnel traffic metrics for chart const [tunnelMetrics, setTunnelMetrics] = useState([]); const [tunnelMetricsLoading, setTunnelMetricsLoading] = useState(false); const [tunnelMetricsError, setTunnelMetricsError] = useState(null); - const [, setTunnelMetricsTruncated] = useState(false); const [tunnelRangeMs, setTunnelRangeMs] = useState(60 * 60 * 1000); - // Tunnel quality (diagnosis) state - const [tunnelQualities, setTunnelQualities] = useState>({}); - const diagnosisAbortRef = useRef>({}); - - // Cleanup abort controllers on unmount - useEffect(() => { - return () => { - Object.values(diagnosisAbortRef.current).forEach((c) => c.abort()); - diagnosisAbortRef.current = {}; - }; - }, []); - + // --- Load tunnel list --- const loadTunnels = useCallback(async (options?: { silent?: boolean }) => { const silent = options?.silent ?? false; if (!silent) setTunnelsLoading(true); @@ -161,7 +179,75 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) return () => window.clearInterval(timer); }, [loadTunnels]); - // Load tunnel metrics for detail view + // --- Load quality snapshots (auto-polling every 10s) --- + const loadQuality = useCallback(async (options?: { silent?: boolean }) => { + const silent = options?.silent ?? false; + if (!silent) setQualityLoading(true); + try { + const response = await getMonitorTunnelQuality(); + if (response.code === 0 && Array.isArray(response.data)) { + const map: Record = {}; + for (const q of response.data) { + map[q.tunnelId] = q; + } + setQualityMap(map); + } + } catch { + // Silently ignore quality load failures + } finally { + if (!silent) setQualityLoading(false); + } + }, []); + + useEffect(() => { + void loadQuality(); + }, [loadQuality]); + + useEffect(() => { + qualityTimerRef.current = window.setInterval(() => { + void loadQuality({ silent: true }); + }, QUALITY_POLL_INTERVAL); + + return () => { + if (qualityTimerRef.current) { + window.clearInterval(qualityTimerRef.current); + } + }; + }, [loadQuality]); + + // --- Load quality history for detail chart --- + const loadQualityHistory = useCallback( + async (tunnelId: number, options?: { silent?: boolean }) => { + const silent = options?.silent ?? false; + if (!silent) setQualityHistoryLoading(true); + try { + const end = Date.now(); + const start = end - qualityRangeMs; + const response = await getMonitorTunnelQualityHistory(tunnelId, start, end); + + if (response.code === 0 && Array.isArray(response.data)) { + setQualityHistoryError(null); + setQualityHistory(response.data); + return; + } + if (response.code === 403) { + setAccessDenied(response.msg || "暂无监控权限"); + return; + } + setQualityHistoryError(response.msg || "加载质量历史失败"); + if (!silent) toast.error(response.msg || "加载质量历史失败"); + } catch { + if (!silent) { + setQualityHistoryError("加载质量历史失败"); + } + } finally { + if (!silent) setQualityHistoryLoading(false); + } + }, + [qualityRangeMs], + ); + + // --- Load tunnel traffic metrics for detail chart --- const loadTunnelMetrics = useCallback( async (tunnelId: number, options?: { silent?: boolean }) => { const silent = options?.silent ?? false; @@ -172,27 +258,16 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) const response = await getTunnelMetrics(tunnelId, start, end); if (response.code === 0 && Array.isArray(response.data)) { - setAccessDenied(null); setTunnelMetricsError(null); - setTunnelMetricsTruncated(response.data.length >= METRICS_MAX_ROWS); const ordered = [...response.data].sort( (a, b) => a.timestamp - b.timestamp, ); setTunnelMetrics(ordered); return; } - if (response.code === 403) { - setAccessDenied(response.msg || "暂无监控权限,请联系管理员授权"); - setTunnelMetricsTruncated(false); - setTunnelMetricsError(null); - return; - } - setTunnelMetricsTruncated(false); - setTunnelMetricsError(response.msg || "加载隧道指标失败"); - if (!silent) toast.error(response.msg || "加载隧道指标失败"); + setTunnelMetricsError(response.msg || "加载流量数据失败"); } catch { - setTunnelMetricsTruncated(false); - if (!silent) setTunnelMetricsError("加载隧道指标失败"); + if (!silent) setTunnelMetricsError("加载流量数据失败"); } finally { if (!silent) setTunnelMetricsLoading(false); } @@ -202,169 +277,32 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) useEffect(() => { if (detailTunnelId) { + void loadQualityHistory(detailTunnelId); void loadTunnelMetrics(detailTunnelId); } - }, [detailTunnelId, loadTunnelMetrics]); + }, [detailTunnelId, loadQualityHistory, loadTunnelMetrics]); + // Auto-refresh detail charts useEffect(() => { if (!detailTunnelId) return; const timer = window.setInterval(() => { + void loadQualityHistory(detailTunnelId, { silent: true }); void loadTunnelMetrics(detailTunnelId, { silent: true }); }, 30_000); return () => window.clearInterval(timer); - }, [detailTunnelId, loadTunnelMetrics]); + }, [detailTunnelId, loadQualityHistory, loadTunnelMetrics]); - // Diagnose tunnel quality - const diagnoseTunnelQuality = useCallback(async (tunnelId: number) => { - // Abort if already running - if (diagnosisAbortRef.current[tunnelId]) { - diagnosisAbortRef.current[tunnelId].abort(); - } - const abortController = new AbortController(); - diagnosisAbortRef.current[tunnelId] = abortController; + // Chart data for quality history + const qualityChartData = qualityHistory.map((q) => ({ + time: formatTimestamp(q.timestamp, qualityRangeMs), + entryToExit: q.entryToExitLatency >= 0 ? q.entryToExitLatency : null, + exitToBing: q.exitToBingLatency >= 0 ? q.exitToBingLatency : null, + entryToExitLoss: q.entryToExitLoss, + exitToBingLoss: q.exitToBingLoss, + })); - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { loading: true }, - })); - - try { - // Try stream first - const results: TunnelDiagnosisApiItem[] = []; - const streamResult = await diagnoseTunnelStream( - tunnelId, - { - onItem: (payload) => { - results.push(payload.result); - }, - onError: (msg) => { - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { loading: false, error: msg }, - })); - }, - }, - abortController.signal, - ); - - if (streamResult.fallback) { - // Fallback to non-stream API - try { - const response = await diagnoseTunnel(tunnelId); - if (response.code === 0 && response.data?.results) { - const apiResults = response.data.results; - const quality = extractQualityFromResults(apiResults); - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { - loading: false, - ...quality, - results: apiResults, - timestamp: Date.now(), - }, - })); - } else { - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { loading: false, error: response.msg || "诊断失败" }, - })); - } - } catch { - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { loading: false, error: "诊断请求失败" }, - })); - } - return; - } - - // Process stream results - if (results.length > 0) { - const quality = extractQualityFromResults(results); - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { - loading: false, - ...quality, - results, - timestamp: Date.now(), - }, - })); - } else { - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { loading: false, error: "未获取到诊断结果" }, - })); - } - } catch { - if (!abortController.signal.aborted) { - setTunnelQualities((prev) => ({ - ...prev, - [tunnelId]: { loading: false, error: "诊断失败" }, - })); - } - } finally { - delete diagnosisAbortRef.current[tunnelId]; - } - }, []); - - const extractQualityFromResults = ( - results: TunnelDiagnosisApiItem[], - ): Pick => { - // The diagnosis results contain hop-by-hop tests - // We look for entry→exit (hop between entry and exit nodes) - // and exit→Bing (the last hop to external target like bing.com) - let entryToExitLatency: number | undefined; - let exitToBingLatency: number | undefined; - let entryToExitLoss: number | undefined; - let exitToBingLoss: number | undefined; - - for (const r of results) { - if (!r.success) continue; - - // Entry to Exit: chainType transitions from 1 (entry) to 3 (exit) - if (r.fromChainType === 1 && r.toChainType === 3) { - entryToExitLatency = r.averageTime; - entryToExitLoss = r.packetLoss; - } - // Or if it's a mid-chain to exit - if (r.fromChainType === 2 && r.toChainType === 3) { - // Use this if no direct entry→exit - if (entryToExitLatency === undefined) { - entryToExitLatency = r.averageTime; - entryToExitLoss = r.packetLoss; - } - } - - // Exit to external target (Bing / external) - if (r.toChainType === undefined || r.toChainType === 0) { - // This typically means it's the exit node testing external - if (r.fromChainType === 3) { - exitToBingLatency = r.averageTime; - exitToBingLoss = r.packetLoss; - } - } - } - - // If no chainType-based matching, use position-based heuristics - if (entryToExitLatency === undefined && exitToBingLatency === undefined) { - const successResults = results.filter((r) => r.success); - if (successResults.length >= 2) { - entryToExitLatency = successResults[0].averageTime; - entryToExitLoss = successResults[0].packetLoss; - exitToBingLatency = successResults[successResults.length - 1].averageTime; - exitToBingLoss = successResults[successResults.length - 1].packetLoss; - } else if (successResults.length === 1) { - entryToExitLatency = successResults[0].averageTime; - entryToExitLoss = successResults[0].packetLoss; - } - } - - return { entryToExitLatency, exitToBingLatency, entryToExitLoss, exitToBingLoss }; - }; - - // Chart data + // Chart data for traffic metrics const tunnelChartData = tunnelMetrics.map((m) => ({ time: formatTimestamp(m.timestamp, tunnelRangeMs), bytesIn: m.bytesIn, @@ -391,14 +329,35 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) // Aggregate stats const tunnelStats = useMemo(() => { const enabled = tunnels.filter((t) => t.status === 1).length; - const disabled = tunnels.length - enabled; - const diagnosed = Object.keys(tunnelQualities).filter((k) => { - const q = tunnelQualities[Number(k)]; - return q && !q.loading && !q.error; - }).length; - return { total: tunnels.length, enabled, disabled, diagnosed }; - }, [tunnels, tunnelQualities]); + return { total: tunnels.length, enabled }; + }, [tunnels]); + + // Last quality update timestamp + const lastQualityUpdate = useMemo(() => { + let latest = 0; + for (const q of Object.values(qualityMap)) { + if (q.timestamp > latest) latest = q.timestamp; + } + return latest > 0 ? new Date(latest).toLocaleTimeString("zh-CN") : null; + }, [qualityMap]); + + /** Shared time range Select component */ + const TimeRangeSelect = ({ value, onChange }: { value: number; onChange: (v: number) => void }) => ( + + ); // ===================== // RENDER @@ -423,7 +382,7 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) // ===== DETAIL VIEW ===== if (detailTunnelId && detailTunnel) { - const quality = tunnelQualities[detailTunnelId]; + const quality = qualityMap[detailTunnelId]; return (
@@ -431,6 +390,7 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps)
+ {/* Auto-probe status */} +
+ + 自动探测中(每10秒) {quality?.timestamp && ( - - 上次检测: {new Date(quality.timestamp).toLocaleTimeString("zh-CN")} + + · 最近更新: {new Date(quality.timestamp).toLocaleTimeString("zh-CN")} )} - {quality?.error && ( - {quality.error} + {quality?.errorMessage && ( + {quality.errorMessage} )}
- {/* Diagnosis Details */} - {quality?.results && quality.results.length > 0 && ( - - -

诊断详情

-
- - - - 描述 - 节点 - 目标 - 延迟 - 丢包 - 状态 - - - {quality.results.map((r, idx) => ( - - - {r.description || "-"} - - - {r.nodeName || "-"} - - - - {r.targetIp || "-"}{r.targetPort ? `:${r.targetPort}` : ""} - - - - - {r.averageTime !== undefined ? `${r.averageTime.toFixed(0)}ms` : "-"} - - - - - {r.packetLoss !== undefined ? `${r.packetLoss.toFixed(1)}%` : "-"} - - - - - {r.success ? "成功" : "失败"} - - - - ))} - -
-
-
- )} - - {/* Tunnel traffic chart */} + {/* ====== Quality History Chart (mirrors service monitor chart) ====== */} -

隧道流量趋势

+

质量趋势

- + + 刷新 + +
+
+ + {qualityHistoryLoading ? ( +
+ ) : qualityHistoryError ? ( +
{qualityHistoryError}
+ ) : qualityChartData.length > 0 ? ( +
+ + + + + `${Number(v).toFixed(0)}ms`} + label={{ value: "延迟 (ms)", angle: -90, position: "insideLeft", style: { fontSize: 11, fill: "#888" } }} + /> + { + const n = Number(value); + if (!Number.isFinite(n)) return "-"; + const label = name === "entryToExit" ? "入口→出口" : name === "exitToBing" ? "出口→Bing" : name; + return [`${n.toFixed(1)}ms`, label]; + }} + /> + { + if (value === "entryToExit") return "入口→出口"; + if (value === "exitToBing") return "出口→Bing"; + return value; + }} + /> + + + + +
+ ) : ( +
暂无质量历史数据
+ )} +
+
+ + {/* ====== Traffic Chart (unchanged) ====== */} + + +

流量趋势

+
+
) : tunnelMetricsError ? (
{tunnelMetricsError}
- ) : tunnelMetrics.length > 0 ? ( + ) : tunnelChartData.length > 0 ? (
- - - + + + +
) : ( -
暂无指标数据
+
暂无流量数据
)}
@@ -616,13 +584,16 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) ); } - // ===== LIST/GRID VIEW ===== + // ===== LIST/GRID VIEW (unchanged) ===== return (
隧道 {tunnelStats.enabled}/{tunnelStats.total} - {tunnelStats.diagnosed > 0 && ( - 已诊断 {tunnelStats.diagnosed} + {lastQualityUpdate && ( +
+ + 自动探测 · 更新于 {lastQualityUpdate} +
)}
@@ -752,13 +695,13 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) 入口→出口 出口→Bing 质量 - 操作 + 更新时间 {tunnels.map((tunnel) => { - const quality = tunnelQualities[tunnel.id]; + const quality = qualityMap[tunnel.id]; const isEnabled = tunnel.status === 1; - const overallQuality = quality?.entryToExitLatency !== undefined + const overallQuality = quality?.entryToExitLatency !== undefined && quality.entryToExitLatency >= 0 ? getDiagnosisQualityDisplay(quality.entryToExitLatency, quality.entryToExitLoss ?? 0) : null; @@ -777,22 +720,10 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) {tunnel.name} - - {quality?.loading ? ( - - ) : quality?.entryToExitLatency !== undefined ? ( - `${quality.entryToExitLatency.toFixed(0)}ms` - ) : "-"} - + - - {quality?.loading ? ( - - ) : quality?.exitToBingLatency !== undefined ? ( - `${quality.exitToBingLatency.toFixed(0)}ms` - ) : "-"} - + {overallQuality ? ( @@ -804,18 +735,14 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps) )} -
- -
+ {quality?.timestamp ? ( + + + {new Date(quality.timestamp).toLocaleTimeString("zh-CN")} + + ) : ( + - + )}
);