mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 07:36:38 +08:00
feat: implement tunnel quality polling and service monitor tuning to 1s/30s intervals
This commit is contained in:
@@ -8,6 +8,7 @@ import (
|
||||
"net"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"go-backend/internal/monitoring"
|
||||
@@ -20,26 +21,47 @@ type nodeCommander interface {
|
||||
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 {
|
||||
repo *repo.Repository
|
||||
commander nodeCommander
|
||||
lastRun map[int64]int64
|
||||
inFlight map[int64]struct{}
|
||||
|
||||
mu sync.RWMutex
|
||||
cancel context.CancelFunc
|
||||
wg sync.WaitGroup
|
||||
// In-memory latest result per monitor (for real-time API reads)
|
||||
latestResults map[int64]*model.ServiceMonitorResult
|
||||
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 {
|
||||
return &Checker{
|
||||
repo: repo,
|
||||
commander: commander,
|
||||
lastRun: make(map[int64]int64),
|
||||
inFlight: make(map[int64]struct{}),
|
||||
repo: repo,
|
||||
commander: commander,
|
||||
lastRun: make(map[int64]int64),
|
||||
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) {
|
||||
c.mu.Lock()
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
@@ -52,7 +74,7 @@ func (c *Checker) Start(ctx context.Context) {
|
||||
limits := c.loadServiceMonitorLimits()
|
||||
scanInterval := time.Duration(limits.CheckerScanIntervalSec) * time.Second
|
||||
if scanInterval <= 0 {
|
||||
scanInterval = 30 * time.Second
|
||||
scanInterval = 1 * time.Second
|
||||
}
|
||||
|
||||
timer := time.NewTimer(scanInterval)
|
||||
@@ -87,6 +109,12 @@ func (c *Checker) RunOnce(m *model.ServiceMonitor) (*model.ServiceMonitorResult,
|
||||
}
|
||||
|
||||
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 {
|
||||
return
|
||||
}
|
||||
@@ -174,6 +202,8 @@ func (c *Checker) runChecks(ctx context.Context) {
|
||||
}
|
||||
close(jobs)
|
||||
|
||||
reportIntervalMs := int64(serviceMonitorReportInterval / time.Millisecond)
|
||||
|
||||
for i := 0; i < workerLimit; i++ {
|
||||
c.wg.Add(1)
|
||||
go func() {
|
||||
@@ -188,14 +218,26 @@ func (c *Checker) runChecks(ctx context.Context) {
|
||||
}
|
||||
ts := time.Now().UnixMilli()
|
||||
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.latestResults[m.ID] = result
|
||||
c.lastRun[m.ID] = result.Timestamp
|
||||
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()
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
// 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()
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
|
||||
@@ -4,17 +4,19 @@ import (
|
||||
"context"
|
||||
"log"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"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
|
||||
tunnelQualityProbeInterval = 1 * time.Second
|
||||
tunnelQualityProbeTimeout = 8 * time.Second
|
||||
tunnelQualityPingTimeoutMs = 5000
|
||||
tunnelQualityRetention = 24 * time.Hour // keep 24h of history
|
||||
tunnelQualityPruneInterval = 10 * time.Minute
|
||||
tunnelQualityReportInterval = 30 * time.Second // DB save interval
|
||||
)
|
||||
|
||||
// tunnelQualitySnapshot is the in-memory latest probe result for a tunnel.
|
||||
@@ -27,6 +29,9 @@ type tunnelQualitySnapshot struct {
|
||||
Success bool `json:"success"`
|
||||
ErrorMessage string `json:"errorMessage,omitempty"`
|
||||
Timestamp int64 `json:"timestamp"`
|
||||
|
||||
// internal fields for db reporting
|
||||
lastDBWrite int64 `json:"-"`
|
||||
}
|
||||
|
||||
// tunnelQualityProber runs periodic TCP ping probes against all enabled tunnels.
|
||||
@@ -38,6 +43,7 @@ type tunnelQualityProber struct {
|
||||
cancel context.CancelFunc
|
||||
interval time.Duration
|
||||
lastPrune int64
|
||||
probing int32 // atomic flag: 1 = probeAll running, 0 = idle
|
||||
}
|
||||
|
||||
// newTunnelQualityProber creates a new prober (not yet running).
|
||||
@@ -120,6 +126,12 @@ func (p *tunnelQualityProber) maybePrune() {
|
||||
}
|
||||
|
||||
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
|
||||
if h == nil || h.repo == nil {
|
||||
return
|
||||
@@ -136,7 +148,7 @@ func (p *tunnelQualityProber) probeAll() {
|
||||
|
||||
// Probe tunnels concurrently with a worker limit
|
||||
// (mirrors health.Checker worker pool pattern)
|
||||
const maxWorkers = 4
|
||||
const maxWorkers = 20
|
||||
sem := make(chan struct{}, maxWorkers)
|
||||
var wg sync.WaitGroup
|
||||
|
||||
@@ -303,8 +315,28 @@ func (p *tunnelQualityProber) storeResult(snap *tunnelQualitySnapshot) {
|
||||
}
|
||||
|
||||
// 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)
|
||||
|
||||
if !writeToDB {
|
||||
return
|
||||
}
|
||||
|
||||
// Persist to database (history)
|
||||
h := p.handler
|
||||
if h == nil || h.repo == nil {
|
||||
|
||||
@@ -29,10 +29,10 @@ const (
|
||||
|
||||
func DefaultServiceMonitorLimits() ServiceMonitorLimits {
|
||||
return ServiceMonitorLimits{
|
||||
CheckerScanIntervalSec: 30,
|
||||
WorkerLimit: 5,
|
||||
MinIntervalSec: 30,
|
||||
DefaultIntervalSec: 60,
|
||||
CheckerScanIntervalSec: 1,
|
||||
WorkerLimit: 20,
|
||||
MinIntervalSec: 1,
|
||||
DefaultIntervalSec: 1,
|
||||
MinTimeoutSec: 1,
|
||||
DefaultTimeoutSec: 5,
|
||||
MaxTimeoutSec: 60,
|
||||
|
||||
@@ -190,7 +190,7 @@ func NewWebSocketReporter(serverURL string, secret string) *WebSocketReporter {
|
||||
return &WebSocketReporter{
|
||||
url: serverURL,
|
||||
curBackoff: initialBackoff, // 当前退避间隔
|
||||
pingInterval: 5 * time.Second, // 指标上报间隔
|
||||
pingInterval: 1 * time.Second, // 指标上报间隔(每秒采集)
|
||||
configInterval: 10 * time.Minute, // 配置上报间隔
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
|
||||
@@ -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 DEFAULT_SERVICE_MONITOR_LIMITS: ServiceMonitorLimitsApiData = {
|
||||
checkerScanIntervalSec: 30,
|
||||
minIntervalSec: 30,
|
||||
defaultIntervalSec: 60,
|
||||
checkerScanIntervalSec: 1,
|
||||
minIntervalSec: 1,
|
||||
defaultIntervalSec: 1,
|
||||
minTimeoutSec: 1,
|
||||
defaultTimeoutSec: 5,
|
||||
maxTimeoutSec: 60,
|
||||
@@ -619,7 +619,7 @@ export function MonitorView({ nodeMap, viewMode = "grid" }: MonitorViewProps) {
|
||||
const timer = window.setInterval(() => {
|
||||
void loadServiceMonitors({ silent: true });
|
||||
void loadLatestMonitorResults();
|
||||
}, 30_000);
|
||||
}, 1_000);
|
||||
|
||||
return () => window.clearInterval(timer);
|
||||
}, [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">
|
||||
<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="text-xs text-default-500">间隔 {resolvedActiveMonitor.intervalSec}s</span>
|
||||
<span className="text-xs text-default-500">每秒测试,30秒上报</span>
|
||||
{activeLatestResult && Number.isFinite(activeLatestResult.latencyMs) ? (
|
||||
<span className="font-mono text-xs font-semibold text-success">{activeLatestResult.latencyMs.toFixed(0)}ms</span>
|
||||
) : null}
|
||||
|
||||
@@ -50,7 +50,7 @@ interface TunnelMonitorViewProps {
|
||||
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 date = new Date(ts);
|
||||
@@ -572,7 +572,7 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps)
|
||||
{/* Auto-probe status */}
|
||||
<div className="flex items-center gap-2 text-xs text-default-500">
|
||||
<LiveDot />
|
||||
<span>自动探测中(每10秒)</span>
|
||||
<span>自动探测中(每秒测试,30秒上报)</span>
|
||||
{quality?.timestamp && (
|
||||
<span className="text-default-400">
|
||||
· 最近更新: {new Date(quality.timestamp).toLocaleTimeString("zh-CN")}
|
||||
@@ -712,7 +712,7 @@ export function TunnelMonitorView({ viewMode = "grid" }: TunnelMonitorViewProps)
|
||||
{lastQualityUpdate && (
|
||||
<div className="flex items-center gap-1.5 text-xs text-default-500">
|
||||
<LiveDot />
|
||||
<span>自动探测 · 更新于 {lastQualityUpdate}</span>
|
||||
<span>每秒探测 · 更新于 {lastQualityUpdate}</span>
|
||||
</div>
|
||||
)}
|
||||
<div className="ml-auto">
|
||||
|
||||
Reference in New Issue
Block a user