feat: periodic tunnel quality probing and monitoring

This commit is contained in:
sagitchu
2026-03-20 12:34:09 +08:00
parent ff7c91d277
commit 3c57a5ac84
11 changed files with 928 additions and 370 deletions
@@ -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)
+9 -1
View File
@@ -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()
+109 -1
View File
@@ -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"))
@@ -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)
}
}