feat: implement tunnel quality polling and service monitor tuning (#362)

Implement 1s test, 30s report pattern across all monitoring subsystems.
This commit is contained in:
sagit
2026-03-21 17:30:07 +08:00
committed by GitHub
8 changed files with 156 additions and 30 deletions
+53 -11
View File
@@ -8,6 +8,7 @@ import (
"net" "net"
"strings" "strings"
"sync" "sync"
"sync/atomic"
"time" "time"
"go-backend/internal/monitoring" "go-backend/internal/monitoring"
@@ -20,26 +21,47 @@ type nodeCommander interface {
SendCommand(nodeID int64, cmdType string, data interface{}, timeout time.Duration) (ws.CommandResult, error) SendCommand(nodeID int64, cmdType string, data interface{}, timeout time.Duration) (ws.CommandResult, error)
} }
const serviceMonitorReportInterval = 30 * time.Second // DB write interval per monitor
type Checker struct { type Checker struct {
repo *repo.Repository repo *repo.Repository
commander nodeCommander commander nodeCommander
lastRun map[int64]int64 lastRun map[int64]int64
inFlight map[int64]struct{} inFlight map[int64]struct{}
mu sync.RWMutex // In-memory latest result per monitor (for real-time API reads)
cancel context.CancelFunc latestResults map[int64]*model.ServiceMonitorResult
wg sync.WaitGroup lastDBWrite map[int64]int64 // last DB write timestamp per monitorID
mu sync.RWMutex
cancel context.CancelFunc
wg sync.WaitGroup
checking int32 // atomic flag: 1 = runChecks running, 0 = idle
} }
func NewChecker(repo *repo.Repository, commander nodeCommander) *Checker { func NewChecker(repo *repo.Repository, commander nodeCommander) *Checker {
return &Checker{ return &Checker{
repo: repo, repo: repo,
commander: commander, commander: commander,
lastRun: make(map[int64]int64), lastRun: make(map[int64]int64),
inFlight: make(map[int64]struct{}), inFlight: make(map[int64]struct{}),
latestResults: make(map[int64]*model.ServiceMonitorResult),
lastDBWrite: make(map[int64]int64),
} }
} }
// GetLatestCached returns the in-memory latest results (updated every 1s).
// Returns nil if no results are cached.
func (c *Checker) GetLatestCached() []*model.ServiceMonitorResult {
c.mu.RLock()
defer c.mu.RUnlock()
results := make([]*model.ServiceMonitorResult, 0, len(c.latestResults))
for _, r := range c.latestResults {
results = append(results, r)
}
return results
}
func (c *Checker) Start(ctx context.Context) { func (c *Checker) Start(ctx context.Context) {
c.mu.Lock() c.mu.Lock()
ctx, cancel := context.WithCancel(ctx) ctx, cancel := context.WithCancel(ctx)
@@ -52,7 +74,7 @@ func (c *Checker) Start(ctx context.Context) {
limits := c.loadServiceMonitorLimits() limits := c.loadServiceMonitorLimits()
scanInterval := time.Duration(limits.CheckerScanIntervalSec) * time.Second scanInterval := time.Duration(limits.CheckerScanIntervalSec) * time.Second
if scanInterval <= 0 { if scanInterval <= 0 {
scanInterval = 30 * time.Second scanInterval = 1 * time.Second
} }
timer := time.NewTimer(scanInterval) timer := time.NewTimer(scanInterval)
@@ -87,6 +109,12 @@ func (c *Checker) RunOnce(m *model.ServiceMonitor) (*model.ServiceMonitorResult,
} }
func (c *Checker) runChecks(ctx context.Context) { func (c *Checker) runChecks(ctx context.Context) {
// Skip if previous round is still running (interval < timeout guard)
if !atomic.CompareAndSwapInt32(&c.checking, 0, 1) {
return
}
defer atomic.StoreInt32(&c.checking, 0)
if c == nil || c.repo == nil { if c == nil || c.repo == nil {
return return
} }
@@ -174,6 +202,8 @@ func (c *Checker) runChecks(ctx context.Context) {
} }
close(jobs) close(jobs)
reportIntervalMs := int64(serviceMonitorReportInterval / time.Millisecond)
for i := 0; i < workerLimit; i++ { for i := 0; i < workerLimit; i++ {
c.wg.Add(1) c.wg.Add(1)
go func() { go func() {
@@ -188,14 +218,26 @@ func (c *Checker) runChecks(ctx context.Context) {
} }
ts := time.Now().UnixMilli() ts := time.Now().UnixMilli()
result := c.executeCheck(&m, ts, limits) result := c.executeCheck(&m, ts, limits)
if err := c.repo.InsertServiceMonitorResult(result); err != nil {
log.Printf("monitoring write failed op=service_monitor_result.insert monitor_id=%d err=%v", result.MonitorID, err)
}
// Always update in-memory cache for real-time reads
c.mu.Lock() c.mu.Lock()
c.latestResults[m.ID] = result
c.lastRun[m.ID] = result.Timestamp c.lastRun[m.ID] = result.Timestamp
delete(c.inFlight, m.ID) delete(c.inFlight, m.ID)
// Only write to DB every 30s per monitor
lastWrite := c.lastDBWrite[m.ID]
writeToDB := ts-lastWrite >= reportIntervalMs
if writeToDB {
c.lastDBWrite[m.ID] = ts
}
c.mu.Unlock() c.mu.Unlock()
if writeToDB {
if err := c.repo.InsertServiceMonitorResult(result); err != nil {
log.Printf("monitoring write failed op=service_monitor_result.insert monitor_id=%d err=%v", result.MonitorID, err)
}
}
} }
} }
}() }()
@@ -744,6 +744,16 @@ func (h *Handler) monitorServiceLatestResultsHandler(w http.ResponseWriter, r *h
return return
} }
// Try in-memory cache first (updated every 1s)
if h.healthCheck != nil {
cached := h.healthCheck.GetLatestCached()
if len(cached) > 0 {
response.WriteJSON(w, response.OK(cached))
return
}
}
// Fallback to database
results, err := h.repo.GetLatestServiceMonitorResults() results, err := h.repo.GetLatestServiceMonitorResults()
if err != nil { if err != nil {
response.WriteJSON(w, response.Err(-2, err.Error())) response.WriteJSON(w, response.Err(-2, err.Error()))
@@ -4,17 +4,19 @@ import (
"context" "context"
"log" "log"
"sync" "sync"
"sync/atomic"
"time" "time"
"go-backend/internal/store/model" "go-backend/internal/store/model"
) )
const ( const (
tunnelQualityProbeInterval = 10 * time.Second tunnelQualityProbeInterval = 1 * time.Second
tunnelQualityProbeTimeout = 8 * time.Second tunnelQualityProbeTimeout = 8 * time.Second
tunnelQualityPingTimeoutMs = 5000 tunnelQualityPingTimeoutMs = 5000
tunnelQualityRetention = 24 * time.Hour // keep 24h of history tunnelQualityRetention = 24 * time.Hour // keep 24h of history
tunnelQualityPruneInterval = 10 * time.Minute tunnelQualityPruneInterval = 10 * time.Minute
tunnelQualityReportInterval = 30 * time.Second // DB save interval
) )
// tunnelQualitySnapshot is the in-memory latest probe result for a tunnel. // tunnelQualitySnapshot is the in-memory latest probe result for a tunnel.
@@ -27,6 +29,9 @@ type tunnelQualitySnapshot struct {
Success bool `json:"success"` Success bool `json:"success"`
ErrorMessage string `json:"errorMessage,omitempty"` ErrorMessage string `json:"errorMessage,omitempty"`
Timestamp int64 `json:"timestamp"` Timestamp int64 `json:"timestamp"`
// internal fields for db reporting
lastDBWrite int64 `json:"-"`
} }
// tunnelQualityProber runs periodic TCP ping probes against all enabled tunnels. // tunnelQualityProber runs periodic TCP ping probes against all enabled tunnels.
@@ -38,6 +43,7 @@ type tunnelQualityProber struct {
cancel context.CancelFunc cancel context.CancelFunc
interval time.Duration interval time.Duration
lastPrune int64 lastPrune int64
probing int32 // atomic flag: 1 = probeAll running, 0 = idle
} }
// newTunnelQualityProber creates a new prober (not yet running). // newTunnelQualityProber creates a new prober (not yet running).
@@ -120,6 +126,12 @@ func (p *tunnelQualityProber) maybePrune() {
} }
func (p *tunnelQualityProber) probeAll() { func (p *tunnelQualityProber) probeAll() {
// Skip if previous probe round is still running (interval < timeout guard)
if !atomic.CompareAndSwapInt32(&p.probing, 0, 1) {
return
}
defer atomic.StoreInt32(&p.probing, 0)
h := p.handler h := p.handler
if h == nil || h.repo == nil { if h == nil || h.repo == nil {
return return
@@ -136,7 +148,7 @@ func (p *tunnelQualityProber) probeAll() {
// Probe tunnels concurrently with a worker limit // Probe tunnels concurrently with a worker limit
// (mirrors health.Checker worker pool pattern) // (mirrors health.Checker worker pool pattern)
const maxWorkers = 4 const maxWorkers = 20
sem := make(chan struct{}, maxWorkers) sem := make(chan struct{}, maxWorkers)
var wg sync.WaitGroup var wg sync.WaitGroup
@@ -303,8 +315,28 @@ func (p *tunnelQualityProber) storeResult(snap *tunnelQualitySnapshot) {
} }
// Update in-memory cache (latest per tunnel) // Update in-memory cache (latest per tunnel)
// Retain the lastDBWrite timestamp if it exists, so we only DB write every 30s
var lastWrite int64
if existing, ok := p.cache.Load(snap.TunnelID); ok {
if eg, ok := existing.(*tunnelQualitySnapshot); ok {
lastWrite = eg.lastDBWrite
}
}
snap.lastDBWrite = lastWrite
now := time.Now().UnixMilli()
writeToDB := false
if now-snap.lastDBWrite >= int64(tunnelQualityReportInterval/time.Millisecond) {
writeToDB = true
snap.lastDBWrite = now
}
p.cache.Store(snap.TunnelID, snap) p.cache.Store(snap.TunnelID, snap)
if !writeToDB {
return
}
// Persist to database (history) // Persist to database (history)
h := p.handler h := p.handler
if h == nil || h.repo == nil { if h == nil || h.repo == nil {
+4 -4
View File
@@ -29,10 +29,10 @@ const (
func DefaultServiceMonitorLimits() ServiceMonitorLimits { func DefaultServiceMonitorLimits() ServiceMonitorLimits {
return ServiceMonitorLimits{ return ServiceMonitorLimits{
CheckerScanIntervalSec: 30, CheckerScanIntervalSec: 1,
WorkerLimit: 5, WorkerLimit: 20,
MinIntervalSec: 30, MinIntervalSec: 1,
DefaultIntervalSec: 60, DefaultIntervalSec: 1,
MinTimeoutSec: 1, MinTimeoutSec: 1,
DefaultTimeoutSec: 5, DefaultTimeoutSec: 5,
MaxTimeoutSec: 60, MaxTimeoutSec: 60,
+1 -1
View File
@@ -190,7 +190,7 @@ func NewWebSocketReporter(serverURL string, secret string) *WebSocketReporter {
return &WebSocketReporter{ return &WebSocketReporter{
url: serverURL, url: serverURL,
curBackoff: initialBackoff, // 当前退避间隔 curBackoff: initialBackoff, // 当前退避间隔
pingInterval: 5 * time.Second, // 指标上报间隔 pingInterval: 1 * time.Second, // 指标上报间隔(每秒采集)
configInterval: 10 * time.Minute, // 配置上报间隔 configInterval: 10 * time.Minute, // 配置上报间隔
ctx: ctx, ctx: ctx,
cancel: cancel, cancel: cancel,
+42
View File
@@ -0,0 +1,42 @@
# 060 Nezha-style Monitoring (1s test, 30s report)
## Objective
Update all monitoring subsystems to test every 1 second and report (write to DB) every 30 seconds, matching Nezha-style monitoring behavior.
## Changes
### 1. Tunnel Quality Prober (`go-backend/internal/http/handler/tunnel_quality_prober.go`)
- [x] Change `tunnelQualityProbeInterval` from 10s to 1s
- [x] Add `tunnelQualityReportInterval = 30s` for DB write throttling
- [x] Update `storeResult` to cache in-memory every tick, write to DB only every 30s per tunnel
- [x] Add atomic `probing` flag to prevent overlapping `probeAll()` goroutine pile-up
- [x] Increase `maxWorkers` from 4 to 20
### 2. Service Monitor Checker (`go-backend/internal/health/checker.go`)
- [x] Add `serviceMonitorReportInterval = 30s` for DB write throttling
- [x] Add `latestResults` in-memory map and `lastDBWrite` map per monitor
- [x] Add `GetLatestCached()` method for real-time API reads
- [x] Add atomic `checking` flag to prevent overlapping `runChecks()` goroutine pile-up
- [x] Modify worker goroutines to always update in-memory cache, only write to DB every 30s
### 3. Service Monitor Limits (`go-backend/internal/monitoring/limits.go`)
- [x] Change `CheckerScanIntervalSec` default from 30 to 1
- [x] Change `WorkerLimit` default from 5 to 20
- [x] Change `MinIntervalSec` default from 30 to 1
- [x] Change `DefaultIntervalSec` default from 60 to 1
### 4. Monitoring API Handler (`go-backend/internal/http/handler/monitoring.go`)
- [x] Update `monitorServiceLatestResultsHandler` to prefer in-memory cached results from `healthCheck.GetLatestCached()`
### 5. Agent WebSocket Reporter (`go-gost/x/socket/websocket_reporter.go`)
- [x] Change `pingInterval` (metric reporting) from 5s to 1s
### 6. Frontend - Tunnel Monitor (`vite-frontend/src/pages/node/tunnel-monitor-view.tsx`)
- [x] Change `QUALITY_POLL_INTERVAL` from 10s to 1s
- [x] Update detail view text: "自动探测中(每秒测试,30秒上报)"
- [x] Update list view text: "每秒探测 · 更新于 ..."
### 7. Frontend - Service Monitor (`vite-frontend/src/pages/node/monitor-view.tsx`)
- [x] Change `DEFAULT_SERVICE_MONITOR_LIMITS` defaults to match backend (1s intervals)
- [x] Change service monitor + latest results polling from 30s to 1s
- [x] Update info bar text: "每秒测试,30秒上报"
@@ -264,9 +264,9 @@ type MetricType =
const METRICS_MAX_ROWS = 5000; const METRICS_MAX_ROWS = 5000;
const DEFAULT_SERVICE_MONITOR_LIMITS: ServiceMonitorLimitsApiData = { const DEFAULT_SERVICE_MONITOR_LIMITS: ServiceMonitorLimitsApiData = {
checkerScanIntervalSec: 30, checkerScanIntervalSec: 1,
minIntervalSec: 30, minIntervalSec: 1,
defaultIntervalSec: 60, defaultIntervalSec: 1,
minTimeoutSec: 1, minTimeoutSec: 1,
defaultTimeoutSec: 5, defaultTimeoutSec: 5,
maxTimeoutSec: 60, maxTimeoutSec: 60,
@@ -619,7 +619,7 @@ export function MonitorView({ nodeMap, viewMode = "grid" }: MonitorViewProps) {
const timer = window.setInterval(() => { const timer = window.setInterval(() => {
void loadServiceMonitors({ silent: true }); void loadServiceMonitors({ silent: true });
void loadLatestMonitorResults(); void loadLatestMonitorResults();
}, 30_000); }, 1_000);
return () => window.clearInterval(timer); return () => window.clearInterval(timer);
}, [loadLatestMonitorResults, loadServiceMonitors]); }, [loadLatestMonitorResults, loadServiceMonitors]);
@@ -1412,7 +1412,7 @@ export function MonitorView({ nodeMap, viewMode = "grid" }: MonitorViewProps) {
<div className="flex items-center gap-3 min-w-0 flex-wrap"> <div className="flex items-center gap-3 min-w-0 flex-wrap">
<Chip size="sm" color="primary" variant="flat">{resolvedActiveMonitor.type.toUpperCase()}</Chip> <Chip size="sm" color="primary" variant="flat">{resolvedActiveMonitor.type.toUpperCase()}</Chip>
<span className="font-mono text-xs text-default-500">{resolvedActiveMonitor.target}</span> <span className="font-mono text-xs text-default-500">{resolvedActiveMonitor.target}</span>
<span className="text-xs text-default-500">间隔 {resolvedActiveMonitor.intervalSec}s</span> <span className="text-xs text-default-500">每秒测试,30秒上报</span>
{activeLatestResult && Number.isFinite(activeLatestResult.latencyMs) ? ( {activeLatestResult && Number.isFinite(activeLatestResult.latencyMs) ? (
<span className="font-mono text-xs font-semibold text-success">{activeLatestResult.latencyMs.toFixed(0)}ms</span> <span className="font-mono text-xs font-semibold text-success">{activeLatestResult.latencyMs.toFixed(0)}ms</span>
) : null} ) : null}
@@ -50,7 +50,7 @@ interface TunnelMonitorViewProps {
viewMode?: "list" | "grid"; viewMode?: "list" | "grid";
} }
const QUALITY_POLL_INTERVAL = 10_000; // 10 seconds const QUALITY_POLL_INTERVAL = 1_000; // 1 second
const formatTimestamp = (ts: number, rangeMs?: number): string => { const formatTimestamp = (ts: number, rangeMs?: number): string => {
const date = new Date(ts); const date = new Date(ts);
@@ -572,7 +572,7 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps)
{/* Auto-probe status */} {/* Auto-probe status */}
<div className="flex items-center gap-2 text-xs text-default-500"> <div className="flex items-center gap-2 text-xs text-default-500">
<LiveDot /> <LiveDot />
<span>自动探测中(每10秒)</span> <span>自动探测中(每秒测试,30秒上报)</span>
{quality?.timestamp && ( {quality?.timestamp && (
<span className="text-default-400"> <span className="text-default-400">
· 最近更新: {new Date(quality.timestamp).toLocaleTimeString("zh-CN")} · 最近更新: {new Date(quality.timestamp).toLocaleTimeString("zh-CN")}
@@ -712,7 +712,7 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps)
{lastQualityUpdate && ( {lastQualityUpdate && (
<div className="flex items-center gap-1.5 text-xs text-default-500"> <div className="flex items-center gap-1.5 text-xs text-default-500">
<LiveDot /> <LiveDot />
<span>自动探测 · 更新于 {lastQualityUpdate}</span> <span>每秒探测 · 更新于 {lastQualityUpdate}</span>
</div> </div>
)} )}
<div className="ml-auto"> <div className="ml-auto">