mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 07:36:38 +08:00
597 lines
16 KiB
Go
597 lines
16 KiB
Go
package handler
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"log"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"go-backend/internal/monitoring"
|
|
"go-backend/internal/store/model"
|
|
)
|
|
|
|
const (
|
|
tunnelQualityProbeTimeout = 8 * time.Second
|
|
tunnelQualityPingTimeoutMs = 5000
|
|
tunnelQualityPruneInterval = 10 * time.Minute
|
|
tunnelQualityReportInterval = 30 * time.Second // DB save interval
|
|
)
|
|
|
|
type TunnelQualityHop struct {
|
|
FromNodeID int64 `json:"fromNodeId"`
|
|
FromNodeName string `json:"fromNodeName"`
|
|
ToNodeID int64 `json:"toNodeId"`
|
|
ToNodeName string `json:"toNodeName"`
|
|
Latency float64 `json:"latency"`
|
|
Loss float64 `json:"loss"`
|
|
TargetIP string `json:"targetIp,omitempty"`
|
|
TargetPort int `json:"targetPort,omitempty"`
|
|
}
|
|
|
|
// 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"`
|
|
ChainDetails string `json:"chainDetails,omitempty"`
|
|
ProbeTargetHost string `json:"probeTargetHost,omitempty"`
|
|
ProbeTargetPort int `json:"probeTargetPort,omitempty"`
|
|
|
|
// internal fields for db reporting
|
|
lastDBWrite int64 `json:"-"`
|
|
}
|
|
|
|
// 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
|
|
wake chan struct{}
|
|
lastPrune int64
|
|
probing int32 // atomic flag: 1 = probeAll running, 0 = idle
|
|
probeNode bestExitProbeFunc
|
|
}
|
|
|
|
// newTunnelQualityProber creates a new prober (not yet running).
|
|
func newTunnelQualityProber(h *Handler) *tunnelQualityProber {
|
|
return &tunnelQualityProber{
|
|
handler: h,
|
|
wake: make(chan struct{}, 1),
|
|
}
|
|
}
|
|
|
|
// 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() {
|
|
if p == nil || p.cancel == nil {
|
|
return
|
|
}
|
|
|
|
p.cancel()
|
|
}
|
|
|
|
func (p *tunnelQualityProber) NotifyConfigChanged() {
|
|
if p == nil || p.wake == nil {
|
|
return
|
|
}
|
|
select {
|
|
case p.wake <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// 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()
|
|
|
|
for {
|
|
timer := time.NewTimer(p.probeInterval())
|
|
select {
|
|
case <-p.ctx.Done():
|
|
stopAndDrainTunnelQualityTimer(timer)
|
|
return
|
|
case <-p.wake:
|
|
stopAndDrainTunnelQualityTimer(timer)
|
|
continue
|
|
case <-timer.C:
|
|
p.probeAll()
|
|
p.maybePrune()
|
|
}
|
|
}
|
|
}
|
|
|
|
func stopAndDrainTunnelQualityTimer(timer *time.Timer) {
|
|
if timer == nil || timer.Stop() {
|
|
return
|
|
}
|
|
select {
|
|
case <-timer.C:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (p *tunnelQualityProber) probeInterval() time.Duration {
|
|
if p == nil || p.handler == nil || p.handler.repo == nil {
|
|
return time.Duration(monitoring.DefaultTunnelQualityProbeIntervalSec) * time.Second
|
|
}
|
|
cfg, err := p.handler.repo.GetConfigsByNames([]string{monitoring.ConfigTunnelQualityProbeIntervalSec})
|
|
if err != nil {
|
|
return time.Duration(monitoring.DefaultTunnelQualityProbeIntervalSec) * time.Second
|
|
}
|
|
seconds := monitoring.TunnelQualityProbeIntervalSecondsFromConfigMap(cfg)
|
|
return time.Duration(seconds) * time.Second
|
|
}
|
|
|
|
func (p *tunnelQualityProber) isEnabled() bool {
|
|
if p == nil || p.handler == nil {
|
|
return true
|
|
}
|
|
|
|
return p.handler.isTunnelQualityMonitoringEnabled()
|
|
}
|
|
|
|
func (p *tunnelQualityProber) retentionDays() int {
|
|
if p == nil || p.handler == nil || p.handler.repo == nil {
|
|
return monitoring.DefaultMonitorRetentionDays
|
|
}
|
|
cfg, err := p.handler.repo.GetConfigsByNames([]string{monitoring.ConfigMonitorRetentionDays})
|
|
if err != nil {
|
|
return monitoring.DefaultMonitorRetentionDays
|
|
}
|
|
return monitoring.MonitoringRetentionDaysFromConfigMap(cfg)
|
|
}
|
|
|
|
// 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(time.Duration(p.retentionDays())*24*time.Hour/time.Millisecond)
|
|
if err := h.repo.PruneTunnelQualityResults(cutoff); err != nil {
|
|
log.Printf("tunnel_quality_prober: prune err=%v", err)
|
|
}
|
|
}
|
|
|
|
func (p *tunnelQualityProber) probeAll() {
|
|
if !p.isEnabled() {
|
|
return
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
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 = 20
|
|
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
|
|
}
|
|
probeTarget := effectiveTunnelProbeTargetValues(tunnel.ProbeTargetHost, tunnel.ProbeTargetPort)
|
|
snap.ProbeTargetHost = probeTarget.Host
|
|
snap.ProbeTargetPort = probeTarget.Port
|
|
|
|
chainRows, err := h.listChainNodesForTunnel(tunnelID)
|
|
if err != nil || len(chainRows) == 0 {
|
|
snap.ErrorMessage = "隧道配置不完整"
|
|
p.storeResult(snap)
|
|
return
|
|
}
|
|
|
|
ipPreference := h.repo.GetTunnelIPPreference(tunnelID)
|
|
inNodes, midNodesGrouped, outNodes := splitChainNodeGroups(chainRows)
|
|
|
|
options := diagnosisExecOptions{
|
|
commandTimeout: tunnelQualityProbeTimeout,
|
|
pingTimeoutMS: tunnelQualityPingTimeoutMs,
|
|
pingCount: 1,
|
|
timeoutMessage: "探测超时",
|
|
}
|
|
p.probeBestExitOwners(tunnelID, inNodes, midNodesGrouped, outNodes, ipPreference, options, probeTarget)
|
|
|
|
entry, _, entryOnline := p.firstOnlineChainNode(inNodes)
|
|
exit, _, exitOnline := p.firstOnlineChainNode(outNodes)
|
|
|
|
switch tunnel.Type {
|
|
case 1:
|
|
// Port forwarding: entry → public probe target only.
|
|
if entryOnline {
|
|
lat, loss, err := p.pingNode(entry.NodeID, probeTarget.Host, probeTarget.Port, options)
|
|
if err == nil {
|
|
snap.ExitToBingLatency = lat
|
|
snap.ExitToBingLoss = loss
|
|
snap.Success = true
|
|
} else {
|
|
snap.ErrorMessage = err.Error()
|
|
}
|
|
} else {
|
|
snap.ErrorMessage = "入口节点均不在线"
|
|
}
|
|
case 2:
|
|
// Tunnel forwarding: entry → exit + exit → Bing
|
|
probeOK := true
|
|
|
|
if !entryOnline {
|
|
probeOK = false
|
|
snap.ErrorMessage = "入口节点均不在线"
|
|
snap.EntryToExitLatency = -1
|
|
snap.EntryToExitLoss = 100
|
|
} else if !exitOnline {
|
|
probeOK = false
|
|
snap.ErrorMessage = "出口节点均不在线"
|
|
snap.EntryToExitLatency = -1
|
|
snap.EntryToExitLoss = 100
|
|
} else {
|
|
var hops []TunnelQualityHop
|
|
var totalLat float64
|
|
remainingSuccessProb := 1.0
|
|
|
|
nodesInPath := make([]chainNodeRecord, 0, 2+len(midNodesGrouped))
|
|
nodesInPath = append(nodesInPath, entry)
|
|
for _, midGroup := range midNodesGrouped {
|
|
mid, _, online := p.firstOnlineChainNode(midGroup)
|
|
if !online {
|
|
probeOK = false
|
|
snap.ErrorMessage = "中间节点组均不在线"
|
|
break
|
|
}
|
|
nodesInPath = append(nodesInPath, mid)
|
|
}
|
|
if probeOK {
|
|
nodesInPath = append(nodesInPath, exit)
|
|
}
|
|
|
|
for i := 0; i < len(nodesInPath)-1; i++ {
|
|
source := nodesInPath[i]
|
|
target := nodesInPath[i+1]
|
|
|
|
hop := TunnelQualityHop{
|
|
FromNodeID: source.NodeID,
|
|
FromNodeName: source.NodeName,
|
|
ToNodeID: target.NodeID,
|
|
ToNodeName: target.NodeName,
|
|
}
|
|
|
|
targetNode, nodeErr := h.getNodeRecord(target.NodeID)
|
|
if nodeErr != nil || !isTunnelProbeNodeOnline(targetNode) {
|
|
snap.ErrorMessage = "节点 " + target.NodeName + " 不可用"
|
|
probeOK = false
|
|
hop.Latency = -1
|
|
hop.Loss = 100
|
|
hops = append(hops, hop)
|
|
break
|
|
}
|
|
|
|
fromNode, _ := h.getNodeRecord(source.NodeID)
|
|
targetIP, targetPort, resolveErr := resolveChainProbeTarget(fromNode, targetNode, target.Port, ipPreference, target.ConnectIP)
|
|
if resolveErr != nil {
|
|
snap.ErrorMessage = "解析节点 " + target.NodeName + " 失败: " + resolveErr.Error()
|
|
probeOK = false
|
|
hop.Latency = -1
|
|
hop.Loss = 100
|
|
hops = append(hops, hop)
|
|
break
|
|
}
|
|
|
|
hop.TargetIP = targetIP
|
|
hop.TargetPort = targetPort
|
|
|
|
lat, loss, err := p.pingNode(source.NodeID, targetIP, targetPort, options)
|
|
if err == nil {
|
|
hop.Latency = lat
|
|
hop.Loss = loss
|
|
totalLat += lat
|
|
remainingSuccessProb *= (1.0 - loss/100.0)
|
|
hops = append(hops, hop)
|
|
} else {
|
|
probeOK = false
|
|
hop.Latency = -1
|
|
hop.Loss = 100
|
|
hops = append(hops, hop)
|
|
if snap.ErrorMessage == "" {
|
|
snap.ErrorMessage = err.Error()
|
|
}
|
|
break
|
|
}
|
|
}
|
|
|
|
if probeOK {
|
|
snap.EntryToExitLatency = totalLat
|
|
snap.EntryToExitLoss = (1.0 - remainingSuccessProb) * 100.0
|
|
} else {
|
|
snap.EntryToExitLatency = -1
|
|
snap.EntryToExitLoss = 100
|
|
}
|
|
|
|
if len(hops) > 0 {
|
|
if b, err := json.Marshal(hops); err == nil {
|
|
snap.ChainDetails = string(b)
|
|
}
|
|
}
|
|
}
|
|
|
|
// Exit → Bing
|
|
if exitOnline {
|
|
lat, loss, err := p.pingNode(exit.NodeID, probeTarget.Host, probeTarget.Port, 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 → public probe target.
|
|
if entryOnline {
|
|
lat, loss, err := p.pingNode(entry.NodeID, probeTarget.Host, probeTarget.Port, options)
|
|
if err == nil {
|
|
snap.ExitToBingLatency = lat
|
|
snap.ExitToBingLoss = loss
|
|
snap.Success = true
|
|
} else {
|
|
snap.ErrorMessage = err.Error()
|
|
}
|
|
} else {
|
|
snap.ErrorMessage = "入口节点均不在线"
|
|
}
|
|
}
|
|
|
|
p.storeResult(snap)
|
|
}
|
|
|
|
func isTunnelProbeNodeOnline(node *nodeRecord) bool {
|
|
return node != nil && (node.IsRemote == 1 || node.Status == 1)
|
|
}
|
|
|
|
func (p *tunnelQualityProber) firstOnlineChainNode(nodes []chainNodeRecord) (chainNodeRecord, *nodeRecord, bool) {
|
|
if p == nil || p.handler == nil {
|
|
return chainNodeRecord{}, nil, false
|
|
}
|
|
for _, candidate := range nodes {
|
|
node, err := p.handler.getNodeRecord(candidate.NodeID)
|
|
if err == nil && isTunnelProbeNodeOnline(node) {
|
|
return candidate, node, true
|
|
}
|
|
}
|
|
return chainNodeRecord{}, nil, false
|
|
}
|
|
|
|
func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chainNodeRecord, chainHops [][]chainNodeRecord, outNodes []chainNodeRecord, ipPreference string, options diagnosisExecOptions, probeTarget tunnelProbeTarget) {
|
|
if p == nil || p.handler == nil || p.handler.bestExit == nil || len(outNodes) <= 1 {
|
|
return
|
|
}
|
|
if !isBestTunnelStrategy(outNodes[0].Strategy) {
|
|
return
|
|
}
|
|
owners := bestExitChainOwners(inNodes, chainHops)
|
|
if len(owners) == 0 {
|
|
return
|
|
}
|
|
nodeMap := make(map[int64]*nodeRecord, len(owners)+len(outNodes))
|
|
for _, owner := range owners {
|
|
if node, err := p.handler.getNodeRecord(owner.NodeID); err == nil && node != nil {
|
|
nodeMap[owner.NodeID] = node
|
|
}
|
|
}
|
|
for _, exit := range outNodes {
|
|
if node, err := p.handler.getNodeRecord(exit.NodeID); err == nil && node != nil {
|
|
nodeMap[exit.NodeID] = node
|
|
}
|
|
}
|
|
// This best-exit decision cache is per decision round; the display-oriented
|
|
// tunnel quality snapshot may still collect its own first-exit public probe.
|
|
roundPinger := newBestExitRoundPinger(p.pingNode)
|
|
for _, owner := range owners {
|
|
if nodeMap[owner.NodeID] == nil {
|
|
continue
|
|
}
|
|
key := bestExitOwnerKey{TunnelID: tunnelID, OwnerNodeID: owner.NodeID}
|
|
p.handler.bestExit.ensureApplied(key, outNodes[0].NodeID, time.Now())
|
|
scores := evaluateBestExitOwner(owner, outNodes, nodeMap, ipPreference, options, probeTarget, roundPinger)
|
|
decision := p.handler.bestExit.observeScores(key, scores, time.Now())
|
|
if decision.Switch {
|
|
now := time.Now()
|
|
if err := p.handler.applyBestExitChainOrder(tunnelID, owner.NodeID, outNodes, decision.Scores, ipPreference); err != nil {
|
|
log.Printf("best_exit: switch apply failed tunnel=%d owner=%d exit=%d err=%v", tunnelID, owner.NodeID, decision.ExitNodeID, err)
|
|
p.handler.bestExit.recordApplyFailure(key, decision.ExitNodeID, now)
|
|
continue
|
|
}
|
|
p.handler.bestExit.setApplied(key, decision.ExitNodeID, time.Now())
|
|
}
|
|
}
|
|
}
|
|
|
|
func (p *tunnelQualityProber) pingNode(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
|
if p != nil && p.probeNode != nil {
|
|
return p.probeNode(nodeID, ip, port, options)
|
|
}
|
|
return p.tcpPingNode(nodeID, ip, port, options)
|
|
}
|
|
|
|
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
|
|
}
|
|
if !isTunnelProbeNodeOnline(node) {
|
|
return 0, 100, errors.New("节点不在线")
|
|
}
|
|
|
|
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)
|
|
// 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 {
|
|
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,
|
|
ChainDetails: snap.ChainDetails,
|
|
}
|
|
if err := h.repo.InsertTunnelQuality(q); err != nil {
|
|
log.Printf("tunnel_quality_prober: insert db err=%v tunnel_id=%d", err, snap.TunnelID)
|
|
}
|
|
}
|