mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 07:36:38 +08:00
Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cbec9a63da | |||
| 0b23d6f7d7 | |||
| 9e6f80019d | |||
| a8fd01d4d8 | |||
| cbe2fc492e | |||
| ae370382d3 | |||
| e112d81697 |
+10
-1
@@ -64,6 +64,14 @@ curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/panel_instal
|
||||
curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh -o install.sh && chmod +x install.sh && ./install.sh
|
||||
```
|
||||
|
||||
Alpine Linux 最小化安装若未包含 `curl`,可使用系统自带的 `wget` 下载:
|
||||
|
||||
```bash
|
||||
wget -O install.sh https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh && chmod +x install.sh && ./install.sh
|
||||
```
|
||||
|
||||
脚本会在 Alpine 上自动安装 Bash,并使用 OpenRC 注册、启动和管理 `flux_agent` 服务;其他受支持的 Linux 发行版继续使用 systemd。
|
||||
|
||||
**安装过程中会提示输入:**
|
||||
- **服务器地址**: 面板端的通信地址(通常是 `http://<面板IP>:<后端端口>`,例如 `http://1.2.3.4:6365`)。
|
||||
- **密钥**: 刚才在面板中获取的节点密钥。
|
||||
@@ -77,7 +85,8 @@ curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh -
|
||||
|
||||
### 3. 验证安装
|
||||
安装完成后,服务会自动启动。
|
||||
- 查看状态: `systemctl status flux_agent`
|
||||
- systemd 查看状态: `systemctl status flux_agent`
|
||||
- Alpine/OpenRC 查看状态: `rc-service flux_agent status`
|
||||
- 回到面板 **节点管理** 页面,该节点状态应显示为 **在线**。
|
||||
|
||||
---
|
||||
|
||||
@@ -59,6 +59,7 @@ type diagnosisWorkItem struct {
|
||||
type diagnosisExecOptions struct {
|
||||
commandTimeout time.Duration
|
||||
pingTimeoutMS int
|
||||
pingCount int
|
||||
timeoutMessage string
|
||||
}
|
||||
|
||||
@@ -1596,10 +1597,14 @@ func (h *Handler) tcpPingViaNode(nodeID int64, ip string, port int, options diag
|
||||
if options.pingTimeoutMS <= 0 {
|
||||
options.pingTimeoutMS = int(diagnosisCommandTimeout / time.Millisecond)
|
||||
}
|
||||
pingCount := options.pingCount
|
||||
if pingCount <= 0 {
|
||||
pingCount = 4
|
||||
}
|
||||
res, err := h.sendNodeCommandWithTimeout(nodeID, "TcpPing", map[string]interface{}{
|
||||
"ip": ip,
|
||||
"port": port,
|
||||
"count": 4,
|
||||
"count": pingCount,
|
||||
"timeout": options.pingTimeoutMS,
|
||||
}, options.commandTimeout, false, false)
|
||||
if err != nil {
|
||||
@@ -1626,12 +1631,16 @@ func (h *Handler) tcpPingViaRemoteNode(node *nodeRecord, ip string, port int, op
|
||||
if options.pingTimeoutMS <= 0 {
|
||||
options.pingTimeoutMS = int(diagnosisCommandTimeout / time.Millisecond)
|
||||
}
|
||||
pingCount := options.pingCount
|
||||
if pingCount <= 0 {
|
||||
pingCount = 4
|
||||
}
|
||||
|
||||
fc := client.NewFederationClientWithTimeout(options.commandTimeout)
|
||||
return fc.Diagnose(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeDiagnoseRequest{
|
||||
IP: strings.TrimSpace(ip),
|
||||
Port: port,
|
||||
Count: 4,
|
||||
Count: pingCount,
|
||||
Timeout: options.pingTimeoutMS,
|
||||
Protocol: "tcp",
|
||||
})
|
||||
|
||||
@@ -1015,6 +1015,7 @@ func (h *Handler) updateConfigs(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
h.notifyTunnelQualityConfigChanged(key)
|
||||
}
|
||||
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
@@ -1062,6 +1063,7 @@ func (h *Handler) updateSingleConfig(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
h.notifyTunnelQualityConfigChanged(name)
|
||||
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
@@ -1110,11 +1112,23 @@ func normalizeAndValidateConfigValue(key, value string) (string, error) {
|
||||
}
|
||||
case monitoring.ConfigMonitorRetentionDays:
|
||||
return monitoring.NormalizeMonitoringRetentionDays(value)
|
||||
case monitoring.ConfigTunnelQualityProbeIntervalSec:
|
||||
return monitoring.NormalizeTunnelQualityProbeIntervalSeconds(value)
|
||||
default:
|
||||
return value, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) notifyTunnelQualityConfigChanged(key string) {
|
||||
if h == nil || h.qualityProber == nil {
|
||||
return
|
||||
}
|
||||
switch strings.TrimSpace(key) {
|
||||
case monitorTunnelQualityEnabledConfigKey, monitoring.ConfigTunnelQualityProbeIntervalSec:
|
||||
h.qualityProber.NotifyConfigChanged()
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) isTunnelQualityMonitoringEnabled() bool {
|
||||
if h == nil || h.repo == nil {
|
||||
return true
|
||||
|
||||
@@ -3,6 +3,7 @@ package handler
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
"net/netip"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -75,6 +76,12 @@ func IsValidNodeAddress(addr string) error {
|
||||
if strings.ContainsAny(addr, "/?") {
|
||||
return fmt.Errorf("address must not contain path or query parameters")
|
||||
}
|
||||
// A bare IPv6 literal contains multiple colons, so net.SplitHostPort treats
|
||||
// it as a malformed host:port pair. Accept IP literals before attempting
|
||||
// host:port parsing; netip also handles scoped IPv6 addresses.
|
||||
if _, err := netip.ParseAddr(addr); err == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
_, _, err := net.SplitHostPort(addr)
|
||||
if err != nil {
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
package handler
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestIssue515IsValidNodeAddressAcceptsBareIPv6(t *testing.T) {
|
||||
for _, addr := range []string{
|
||||
"2001:db8::1",
|
||||
"::1",
|
||||
"fe80::1%eth0",
|
||||
} {
|
||||
t.Run(addr, func(t *testing.T) {
|
||||
if err := IsValidNodeAddress(addr); err != nil {
|
||||
t.Fatalf("expected bare IPv6 address %q to be accepted: %v", addr, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsValidNodeAddressKeepsExistingAddressForms(t *testing.T) {
|
||||
for _, addr := range []string{
|
||||
"203.0.113.10",
|
||||
"node.example.com",
|
||||
"node.example.com:6365",
|
||||
"[2001:db8::1]:6365",
|
||||
} {
|
||||
t.Run(addr, func(t *testing.T) {
|
||||
if err := IsValidNodeAddress(addr); err != nil {
|
||||
t.Fatalf("expected node address %q to be accepted: %v", addr, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsValidNodeAddressRejectsURLComponents(t *testing.T) {
|
||||
for _, addr := range []string{
|
||||
"https://node.example.com",
|
||||
"node.example.com/path",
|
||||
"node.example.com?transport=tcp",
|
||||
} {
|
||||
t.Run(addr, func(t *testing.T) {
|
||||
if err := IsValidNodeAddress(addr); err == nil {
|
||||
t.Fatalf("expected node address %q to be rejected", addr)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -152,7 +152,11 @@ func evaluateBestExitOwner(owner chainNodeRecord, exits []chainNodeRecord, nodes
|
||||
ownerNode := nodes[owner.NodeID]
|
||||
for _, exit := range exits {
|
||||
exitNode := nodes[exit.NodeID]
|
||||
if exitNode == nil {
|
||||
if !isTunnelProbeNodeOnline(ownerNode) {
|
||||
scores = append(scores, failedBestExitCandidate(owner.NodeID, exit, "owner node offline"))
|
||||
continue
|
||||
}
|
||||
if !isTunnelProbeNodeOnline(exitNode) {
|
||||
scores = append(scores, failedBestExitCandidate(owner.NodeID, exit, "exit node unavailable"))
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -358,9 +358,9 @@ func TestEvaluateBestExitOwnerScoresAllCandidates(t *testing.T) {
|
||||
{NodeID: 31, NodeName: "exit-b", Port: 30031},
|
||||
}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
31: {ID: 31, ServerIP: "10.0.0.31", ServerIPv4: "10.0.0.31", TCPListenAddr: "[::]"},
|
||||
10: {ID: 10, Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Status: 1, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
31: {ID: 31, Status: 1, ServerIP: "10.0.0.31", ServerIPv4: "10.0.0.31", TCPListenAddr: "[::]"},
|
||||
}
|
||||
pinger := func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
switch {
|
||||
@@ -387,12 +387,30 @@ func TestEvaluateBestExitOwnerScoresAllCandidates(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluateBestExitOwnerSkipsOfflineCandidate(t *testing.T) {
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30030}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10"},
|
||||
30: {ID: 30, Status: 0, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30"},
|
||||
}
|
||||
ping := func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
t.Fatalf("offline best-exit candidate should not be probed: node=%d target=%s:%d", nodeID, ip, port)
|
||||
return 0, 100, nil
|
||||
}
|
||||
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, defaultTunnelProbeTarget(), ping)
|
||||
if len(scores) != 1 || scores[0].Success {
|
||||
t.Fatalf("expected one failed offline candidate, got %+v", scores)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluateBestExitOwnerUsesConfiguredPublicProbeTarget(t *testing.T) {
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry-a"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30001}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, Name: "entry-a", ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10"},
|
||||
30: {ID: 30, Name: "exit-a", ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30"},
|
||||
10: {ID: 10, Name: "entry-a", Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10"},
|
||||
30: {ID: 30, Name: "exit-a", Status: 1, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30"},
|
||||
}
|
||||
target := tunnelProbeTarget{Host: "speed.example.com", Port: 8443}
|
||||
var calls []string
|
||||
@@ -419,8 +437,8 @@ func TestEvaluateBestExitOwnerMarksCandidateFailedWhenOwnerToExitFails(t *testin
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30030}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
10: {ID: 10, Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Status: 1, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
}
|
||||
pinger := func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
return 0, 100, errBestExitProbeForTest
|
||||
@@ -436,8 +454,8 @@ func TestEvaluateBestExitOwnerMarksCandidateFailedWhenTargetResolutionFails(t *t
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-v6", Port: 30030}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, Name: "entry", ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Name: "exit-v6", ServerIP: "2001:db8::30", ServerIPv6: "2001:db8::30", TCPListenAddr: "[::]"},
|
||||
10: {ID: 10, Name: "entry", Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Name: "exit-v6", Status: 1, ServerIP: "2001:db8::30", ServerIPv6: "2001:db8::30", TCPListenAddr: "[::]"},
|
||||
}
|
||||
pinger := func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
t.Fatalf("ping should not be called when target resolution fails: node=%d ip=%s port=%d", nodeID, ip, port)
|
||||
|
||||
@@ -3,6 +3,8 @@ package handler
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -13,7 +15,6 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
tunnelQualityProbeInterval = 1 * time.Second
|
||||
tunnelQualityProbeTimeout = 8 * time.Second
|
||||
tunnelQualityPingTimeoutMs = 5000
|
||||
tunnelQualityPruneInterval = 10 * time.Minute
|
||||
@@ -31,6 +32,26 @@ type TunnelQualityHop struct {
|
||||
TargetPort int `json:"targetPort,omitempty"`
|
||||
}
|
||||
|
||||
type TunnelQualityCandidateHop struct {
|
||||
TunnelQualityHop
|
||||
FromRole string `json:"fromRole"`
|
||||
ToRole string `json:"toRole"`
|
||||
HopIndex int `json:"hopIndex"`
|
||||
Selected bool `json:"selected"`
|
||||
ErrorMessage string `json:"errorMessage,omitempty"`
|
||||
}
|
||||
|
||||
type tunnelQualityChainDetails struct {
|
||||
PrimaryPath []TunnelQualityHop `json:"primaryPath,omitempty"`
|
||||
CandidateHops []TunnelQualityCandidateHop `json:"candidateHops,omitempty"`
|
||||
}
|
||||
|
||||
type tunnelQualityCandidateGroup struct {
|
||||
role string
|
||||
roleIndex int
|
||||
nodes []chainNodeRecord
|
||||
}
|
||||
|
||||
// tunnelQualitySnapshot is the in-memory latest probe result for a tunnel.
|
||||
type tunnelQualitySnapshot struct {
|
||||
TunnelID int64 `json:"tunnelId"`
|
||||
@@ -56,7 +77,7 @@ type tunnelQualityProber struct {
|
||||
cache sync.Map // tunnelID (int64) → *tunnelQualitySnapshot
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
interval time.Duration
|
||||
wake chan struct{}
|
||||
lastPrune int64
|
||||
probing int32 // atomic flag: 1 = probeAll running, 0 = idle
|
||||
probeNode bestExitProbeFunc
|
||||
@@ -65,8 +86,8 @@ type tunnelQualityProber struct {
|
||||
// newTunnelQualityProber creates a new prober (not yet running).
|
||||
func newTunnelQualityProber(h *Handler) *tunnelQualityProber {
|
||||
return &tunnelQualityProber{
|
||||
handler: h,
|
||||
interval: tunnelQualityProbeInterval,
|
||||
handler: h,
|
||||
wake: make(chan struct{}, 1),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +107,16 @@ func (p *tunnelQualityProber) Stop() {
|
||||
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
|
||||
@@ -109,20 +140,44 @@ func (p *tunnelQualityProber) loop() {
|
||||
// Run once immediately
|
||||
p.probeAll()
|
||||
|
||||
ticker := time.NewTicker(p.interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
timer := time.NewTimer(p.probeInterval())
|
||||
select {
|
||||
case <-p.ctx.Done():
|
||||
stopAndDrainTunnelQualityTimer(timer)
|
||||
return
|
||||
case <-ticker.C:
|
||||
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
|
||||
@@ -246,15 +301,28 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
options := diagnosisExecOptions{
|
||||
commandTimeout: tunnelQualityProbeTimeout,
|
||||
pingTimeoutMS: tunnelQualityPingTimeoutMs,
|
||||
pingCount: 1,
|
||||
timeoutMessage: "探测超时",
|
||||
}
|
||||
p.probeBestExitOwners(tunnelID, inNodes, midNodesGrouped, outNodes, ipPreference, options, probeTarget)
|
||||
roundPinger := newBestExitRoundPinger(p.pingNode)
|
||||
p.probeBestExitOwners(tunnelID, inNodes, midNodesGrouped, outNodes, ipPreference, options, probeTarget, roundPinger)
|
||||
|
||||
entry, _, entryOnline := p.firstOnlineChainNode(inNodes)
|
||||
exit, _, exitOnline := p.firstOnlineChainNode(outNodes)
|
||||
selectedNodeIDs := make(map[string]int64, 2+len(midNodesGrouped))
|
||||
if entryOnline {
|
||||
selectedNodeIDs[tunnelQualityGroupKey("entry", 0)] = entry.NodeID
|
||||
}
|
||||
if exitOnline {
|
||||
selectedNodeIDs[tunnelQualityGroupKey("exit", 0)] = exit.NodeID
|
||||
}
|
||||
var primaryHops []TunnelQualityHop
|
||||
|
||||
switch tunnel.Type {
|
||||
case 1:
|
||||
// Port forwarding: entry → public probe target only.
|
||||
if len(inNodes) > 0 {
|
||||
lat, loss, err := p.pingNode(inNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if entryOnline {
|
||||
lat, loss, err := roundPinger(entry.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -262,24 +330,42 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
} else {
|
||||
snap.ErrorMessage = err.Error()
|
||||
}
|
||||
} else {
|
||||
snap.ErrorMessage = "入口节点均不在线"
|
||||
}
|
||||
case 2:
|
||||
// Tunnel forwarding: entry → exit + exit → Bing
|
||||
probeOK := true
|
||||
|
||||
if len(inNodes) > 0 && len(outNodes) > 0 {
|
||||
var hops []TunnelQualityHop
|
||||
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 totalLat float64
|
||||
remainingSuccessProb := 1.0
|
||||
|
||||
nodesInPath := make([]chainNodeRecord, 0, 2+len(midNodesGrouped))
|
||||
nodesInPath = append(nodesInPath, inNodes[0])
|
||||
for _, midGroup := range midNodesGrouped {
|
||||
if len(midGroup) > 0 {
|
||||
nodesInPath = append(nodesInPath, midGroup[0])
|
||||
nodesInPath = append(nodesInPath, entry)
|
||||
for midIndex, midGroup := range midNodesGrouped {
|
||||
mid, _, online := p.firstOnlineChainNode(midGroup)
|
||||
if !online {
|
||||
probeOK = false
|
||||
snap.ErrorMessage = "中间节点组均不在线"
|
||||
break
|
||||
}
|
||||
nodesInPath = append(nodesInPath, mid)
|
||||
selectedNodeIDs[tunnelQualityGroupKey("middle", midIndex)] = mid.NodeID
|
||||
}
|
||||
if probeOK {
|
||||
nodesInPath = append(nodesInPath, exit)
|
||||
}
|
||||
nodesInPath = append(nodesInPath, outNodes[0])
|
||||
|
||||
for i := 0; i < len(nodesInPath)-1; i++ {
|
||||
source := nodesInPath[i]
|
||||
@@ -293,12 +379,12 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
}
|
||||
|
||||
targetNode, nodeErr := h.getNodeRecord(target.NodeID)
|
||||
if nodeErr != nil || targetNode == nil {
|
||||
if nodeErr != nil || !isTunnelProbeNodeOnline(targetNode) {
|
||||
snap.ErrorMessage = "节点 " + target.NodeName + " 不可用"
|
||||
probeOK = false
|
||||
hop.Latency = -1
|
||||
hop.Loss = 100
|
||||
hops = append(hops, hop)
|
||||
primaryHops = append(primaryHops, hop)
|
||||
break
|
||||
}
|
||||
|
||||
@@ -309,25 +395,25 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
probeOK = false
|
||||
hop.Latency = -1
|
||||
hop.Loss = 100
|
||||
hops = append(hops, hop)
|
||||
primaryHops = append(primaryHops, hop)
|
||||
break
|
||||
}
|
||||
|
||||
hop.TargetIP = targetIP
|
||||
hop.TargetPort = targetPort
|
||||
|
||||
lat, loss, err := p.pingNode(source.NodeID, targetIP, targetPort, options)
|
||||
lat, loss, err := roundPinger(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)
|
||||
primaryHops = append(primaryHops, hop)
|
||||
} else {
|
||||
probeOK = false
|
||||
hop.Latency = -1
|
||||
hop.Loss = 100
|
||||
hops = append(hops, hop)
|
||||
primaryHops = append(primaryHops, hop)
|
||||
if snap.ErrorMessage == "" {
|
||||
snap.ErrorMessage = err.Error()
|
||||
}
|
||||
@@ -342,17 +428,11 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
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 len(outNodes) > 0 {
|
||||
lat, loss, err := p.pingNode(outNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if exitOnline {
|
||||
lat, loss, err := roundPinger(exit.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -367,8 +447,8 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
snap.Success = probeOK
|
||||
default:
|
||||
// Unknown type: entry → public probe target.
|
||||
if len(inNodes) > 0 {
|
||||
lat, loss, err := p.pingNode(inNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if entryOnline {
|
||||
lat, loss, err := roundPinger(entry.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -376,13 +456,215 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
} else {
|
||||
snap.ErrorMessage = err.Error()
|
||||
}
|
||||
} else {
|
||||
snap.ErrorMessage = "入口节点均不在线"
|
||||
}
|
||||
}
|
||||
|
||||
candidateHops := p.probeTunnelCandidateHops(
|
||||
tunnel.Type,
|
||||
inNodes,
|
||||
midNodesGrouped,
|
||||
outNodes,
|
||||
selectedNodeIDs,
|
||||
ipPreference,
|
||||
options,
|
||||
probeTarget,
|
||||
roundPinger,
|
||||
)
|
||||
if len(primaryHops) > 0 || len(candidateHops) > 0 {
|
||||
details := tunnelQualityChainDetails{
|
||||
PrimaryPath: primaryHops,
|
||||
CandidateHops: candidateHops,
|
||||
}
|
||||
if b, err := json.Marshal(details); err == nil {
|
||||
snap.ChainDetails = string(b)
|
||||
}
|
||||
}
|
||||
|
||||
p.storeResult(snap)
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chainNodeRecord, chainHops [][]chainNodeRecord, outNodes []chainNodeRecord, ipPreference string, options diagnosisExecOptions, probeTarget tunnelProbeTarget) {
|
||||
func tunnelQualityGroupKey(role string, index int) string {
|
||||
return fmt.Sprintf("%s:%d", role, index)
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probeTunnelCandidateHops(
|
||||
tunnelType int,
|
||||
inNodes []chainNodeRecord,
|
||||
chainHops [][]chainNodeRecord,
|
||||
outNodes []chainNodeRecord,
|
||||
selectedNodeIDs map[string]int64,
|
||||
ipPreference string,
|
||||
options diagnosisExecOptions,
|
||||
probeTarget tunnelProbeTarget,
|
||||
ping bestExitProbeFunc,
|
||||
) []TunnelQualityCandidateHop {
|
||||
if p == nil || p.handler == nil || ping == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
if tunnelType != 2 {
|
||||
return p.probePublicTargetCandidates("entry", 0, inNodes, selectedNodeIDs, options, probeTarget, ping)
|
||||
}
|
||||
|
||||
groups := make([]tunnelQualityCandidateGroup, 0, 2+len(chainHops))
|
||||
groups = append(groups, tunnelQualityCandidateGroup{role: "entry", roleIndex: 0, nodes: inNodes})
|
||||
for i, hop := range chainHops {
|
||||
groups = append(groups, tunnelQualityCandidateGroup{role: "middle", roleIndex: i, nodes: hop})
|
||||
}
|
||||
groups = append(groups, tunnelQualityCandidateGroup{role: "exit", roleIndex: 0, nodes: outNodes})
|
||||
|
||||
var items []TunnelQualityCandidateHop
|
||||
for i := 0; i < len(groups)-1; i++ {
|
||||
items = append(items, p.probeCandidateGroupLinks(
|
||||
groups[i],
|
||||
groups[i+1],
|
||||
i,
|
||||
selectedNodeIDs,
|
||||
ipPreference,
|
||||
options,
|
||||
ping,
|
||||
)...)
|
||||
}
|
||||
items = append(items, p.probePublicTargetCandidates(
|
||||
"exit",
|
||||
0,
|
||||
outNodes,
|
||||
selectedNodeIDs,
|
||||
options,
|
||||
probeTarget,
|
||||
ping,
|
||||
)...)
|
||||
return items
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probeCandidateGroupLinks(
|
||||
fromGroup tunnelQualityCandidateGroup,
|
||||
toGroup tunnelQualityCandidateGroup,
|
||||
hopIndex int,
|
||||
selectedNodeIDs map[string]int64,
|
||||
ipPreference string,
|
||||
options diagnosisExecOptions,
|
||||
ping bestExitProbeFunc,
|
||||
) []TunnelQualityCandidateHop {
|
||||
items := make([]TunnelQualityCandidateHop, 0, len(fromGroup.nodes)*len(toGroup.nodes))
|
||||
for _, source := range fromGroup.nodes {
|
||||
for _, target := range toGroup.nodes {
|
||||
item := TunnelQualityCandidateHop{
|
||||
TunnelQualityHop: TunnelQualityHop{
|
||||
FromNodeID: source.NodeID,
|
||||
FromNodeName: source.NodeName,
|
||||
ToNodeID: target.NodeID,
|
||||
ToNodeName: target.NodeName,
|
||||
Latency: -1,
|
||||
Loss: 100,
|
||||
},
|
||||
FromRole: fromGroup.role,
|
||||
ToRole: toGroup.role,
|
||||
HopIndex: hopIndex,
|
||||
Selected: selectedNodeIDs[tunnelQualityGroupKey(fromGroup.role, fromGroup.roleIndex)] == source.NodeID &&
|
||||
selectedNodeIDs[tunnelQualityGroupKey(toGroup.role, toGroup.roleIndex)] == target.NodeID,
|
||||
}
|
||||
|
||||
sourceNode, sourceErr := p.handler.getNodeRecord(source.NodeID)
|
||||
if sourceErr != nil || !isTunnelProbeNodeOnline(sourceNode) {
|
||||
item.ErrorMessage = "来源节点不在线"
|
||||
items = append(items, item)
|
||||
continue
|
||||
}
|
||||
targetNode, targetErr := p.handler.getNodeRecord(target.NodeID)
|
||||
if targetErr != nil || !isTunnelProbeNodeOnline(targetNode) {
|
||||
item.ErrorMessage = "目标节点不在线"
|
||||
items = append(items, item)
|
||||
continue
|
||||
}
|
||||
|
||||
targetIP, targetPort, resolveErr := resolveChainProbeTarget(sourceNode, targetNode, target.Port, ipPreference, target.ConnectIP)
|
||||
if resolveErr != nil {
|
||||
item.ErrorMessage = resolveErr.Error()
|
||||
items = append(items, item)
|
||||
continue
|
||||
}
|
||||
item.TargetIP = targetIP
|
||||
item.TargetPort = targetPort
|
||||
latency, loss, probeErr := ping(source.NodeID, targetIP, targetPort, options)
|
||||
if probeErr != nil {
|
||||
item.ErrorMessage = probeErr.Error()
|
||||
items = append(items, item)
|
||||
continue
|
||||
}
|
||||
item.Latency = latency
|
||||
item.Loss = loss
|
||||
items = append(items, item)
|
||||
}
|
||||
}
|
||||
return items
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probePublicTargetCandidates(
|
||||
fromRole string,
|
||||
fromIndex int,
|
||||
nodes []chainNodeRecord,
|
||||
selectedNodeIDs map[string]int64,
|
||||
options diagnosisExecOptions,
|
||||
probeTarget tunnelProbeTarget,
|
||||
ping bestExitProbeFunc,
|
||||
) []TunnelQualityCandidateHop {
|
||||
items := make([]TunnelQualityCandidateHop, 0, len(nodes))
|
||||
for _, source := range nodes {
|
||||
item := TunnelQualityCandidateHop{
|
||||
TunnelQualityHop: TunnelQualityHop{
|
||||
FromNodeID: source.NodeID,
|
||||
FromNodeName: source.NodeName,
|
||||
ToNodeName: formatTunnelProbeTarget(probeTarget),
|
||||
Latency: -1,
|
||||
Loss: 100,
|
||||
TargetIP: probeTarget.Host,
|
||||
TargetPort: probeTarget.Port,
|
||||
},
|
||||
FromRole: fromRole,
|
||||
ToRole: "target",
|
||||
HopIndex: fromIndex,
|
||||
Selected: selectedNodeIDs[tunnelQualityGroupKey(fromRole, fromIndex)] == source.NodeID,
|
||||
}
|
||||
sourceNode, sourceErr := p.handler.getNodeRecord(source.NodeID)
|
||||
if sourceErr != nil || !isTunnelProbeNodeOnline(sourceNode) {
|
||||
item.ErrorMessage = "来源节点不在线"
|
||||
items = append(items, item)
|
||||
continue
|
||||
}
|
||||
latency, loss, probeErr := ping(source.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if probeErr != nil {
|
||||
item.ErrorMessage = probeErr.Error()
|
||||
items = append(items, item)
|
||||
continue
|
||||
}
|
||||
item.Latency = latency
|
||||
item.Loss = loss
|
||||
items = append(items, item)
|
||||
}
|
||||
return items
|
||||
}
|
||||
|
||||
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, roundPinger bestExitProbeFunc) {
|
||||
if p == nil || p.handler == nil || p.handler.bestExit == nil || len(outNodes) <= 1 {
|
||||
return
|
||||
}
|
||||
@@ -404,9 +686,6 @@ func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chai
|
||||
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
|
||||
@@ -444,6 +723,9 @@ func (p *tunnelQualityProber) tcpPingNode(nodeID int64, ip string, port int, opt
|
||||
if nodeErr != nil {
|
||||
return 0, 100, nodeErr
|
||||
}
|
||||
if !isTunnelProbeNodeOnline(node) {
|
||||
return 0, 100, errors.New("节点不在线")
|
||||
}
|
||||
|
||||
var pingData map[string]interface{}
|
||||
var pingErr error
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"slices"
|
||||
"testing"
|
||||
@@ -26,6 +27,9 @@ func TestTunnelQualityProberUsesConfiguredProbeTarget(t *testing.T) {
|
||||
p := newTunnelQualityProber(h)
|
||||
var calls []string
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
if options.pingCount != 1 {
|
||||
t.Fatalf("expected real-time quality probe count 1, got %d", options.pingCount)
|
||||
}
|
||||
calls = append(calls, fmt.Sprintf("%d|%s|%d", nodeID, ip, port))
|
||||
return 10, 0, nil
|
||||
}
|
||||
@@ -46,6 +50,114 @@ func TestTunnelQualityProberUsesConfiguredProbeTarget(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberSkipsAllOfflineExits(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedQualityForwardTunnel(t, h, 81, []int{0, 0, 0})
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
probeCalls := 0
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
probeCalls++
|
||||
return 0, 100, fmt.Errorf("unexpected probe node=%d target=%s:%d", nodeID, ip, port)
|
||||
}
|
||||
p.probeTunnel(81)
|
||||
|
||||
if probeCalls != 0 {
|
||||
t.Fatalf("expected no TCP probes when all exits are offline, got %d", probeCalls)
|
||||
}
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 {
|
||||
t.Fatalf("expected one quality snapshot, got %+v", snaps)
|
||||
}
|
||||
if snaps[0].Success || snaps[0].ErrorMessage != "出口节点均不在线" {
|
||||
t.Fatalf("expected offline exit snapshot, got %+v", snaps[0])
|
||||
}
|
||||
if snaps[0].EntryToExitLoss != 100 {
|
||||
t.Fatalf("expected 100%% entry-to-exit loss, got %+v", snaps[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberUsesOnlineBackupExit(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedQualityForwardTunnel(t, h, 82, []int{0, 1})
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
var calls []string
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
if options.pingCount != 1 {
|
||||
t.Fatalf("expected real-time quality probe count 1, got %d", options.pingCount)
|
||||
}
|
||||
calls = append(calls, fmt.Sprintf("%d|%s|%d", nodeID, ip, port))
|
||||
return 10, 0, nil
|
||||
}
|
||||
p.probeTunnel(82)
|
||||
|
||||
if slices.Contains(calls, "10|10.0.0.30|30030") {
|
||||
t.Fatalf("did not expect probe to offline primary exit, calls=%+v", calls)
|
||||
}
|
||||
if !slices.Contains(calls, "10|10.0.0.31|30031") {
|
||||
t.Fatalf("expected entry probe to online backup exit, calls=%+v", calls)
|
||||
}
|
||||
if !slices.Contains(calls, "31|www.bing.com|443") {
|
||||
t.Fatalf("expected public probe from online backup exit, calls=%+v", calls)
|
||||
}
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 || !snaps[0].Success {
|
||||
t.Fatalf("expected successful backup exit snapshot, got %+v", snaps)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberReportsAllExitCandidateLatencies(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedQualityForwardTunnel(t, h, 83, []int{1, 1})
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
switch fmt.Sprintf("%d|%s|%d", nodeID, ip, port) {
|
||||
case "10|10.0.0.30|30030":
|
||||
return 20, 0, nil
|
||||
case "10|10.0.0.31|30031":
|
||||
return 35, 0, nil
|
||||
case "30|www.bing.com|443":
|
||||
return 50, 0, nil
|
||||
case "31|www.bing.com|443":
|
||||
return 65, 0, nil
|
||||
default:
|
||||
return 0, 100, fmt.Errorf("unexpected probe node=%d target=%s:%d", nodeID, ip, port)
|
||||
}
|
||||
}
|
||||
p.probeTunnel(83)
|
||||
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 {
|
||||
t.Fatalf("expected one quality snapshot, got %+v", snaps)
|
||||
}
|
||||
if snaps[0].EntryToExitLatency != 20 || snaps[0].ExitToBingLatency != 50 {
|
||||
t.Fatalf("expected primary path metrics to remain unchanged, got %+v", snaps[0])
|
||||
}
|
||||
|
||||
var details tunnelQualityChainDetails
|
||||
if err := json.Unmarshal([]byte(snaps[0].ChainDetails), &details); err != nil {
|
||||
t.Fatalf("decode chain details: %v", err)
|
||||
}
|
||||
assertCandidateHop := func(fromID, toID int64, latency float64, selected bool) {
|
||||
t.Helper()
|
||||
for _, hop := range details.CandidateHops {
|
||||
if hop.FromNodeID == fromID && hop.ToNodeID == toID {
|
||||
if hop.Latency != latency || hop.Selected != selected || hop.ErrorMessage != "" {
|
||||
t.Fatalf("unexpected candidate hop: %+v", hop)
|
||||
}
|
||||
return
|
||||
}
|
||||
}
|
||||
t.Fatalf("candidate hop %d -> %d not found in %+v", fromID, toID, details.CandidateHops)
|
||||
}
|
||||
assertCandidateHop(10, 30, 20, true)
|
||||
assertCandidateHop(10, 31, 35, false)
|
||||
assertCandidateHop(30, 0, 50, true)
|
||||
assertCandidateHop(31, 0, 65, false)
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberStoresProbeTargetWhenChainIncomplete(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 78, "quality-target-incomplete", "speed.example.com", 8443)
|
||||
@@ -67,3 +179,69 @@ func TestTunnelQualityProberStoresProbeTargetWhenChainIncomplete(t *testing.T) {
|
||||
t.Fatalf("unexpected snapshot target metadata: %+v", snaps[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberUsesConfiguredInterval(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
if err := h.repo.UpsertConfig("monitor_tunnel_quality_interval_sec", "15", time.Now().UnixMilli()); err != nil {
|
||||
t.Fatalf("upsert interval config: %v", err)
|
||||
}
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
if got := p.probeInterval(); got != 15*time.Second {
|
||||
t.Fatalf("probe interval = %s, want 15s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberConfigNotificationIsCoalesced(t *testing.T) {
|
||||
p := newTunnelQualityProber(nil)
|
||||
p.NotifyConfigChanged()
|
||||
p.NotifyConfigChanged()
|
||||
|
||||
if got := len(p.wake); got != 1 {
|
||||
t.Fatalf("wake notifications = %d, want 1", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelQualityProbeIntervalConfigValue(t *testing.T) {
|
||||
got, err := normalizeAndValidateConfigValue("monitor_tunnel_quality_interval_sec", " 15 ")
|
||||
if err != nil || got != "15" {
|
||||
t.Fatalf("normalize interval = %q, %v", got, err)
|
||||
}
|
||||
if _, err := normalizeAndValidateConfigValue("monitor_tunnel_quality_interval_sec", "0"); err == nil {
|
||||
t.Fatalf("expected invalid interval to be rejected")
|
||||
}
|
||||
}
|
||||
|
||||
func seedQualityForwardTunnel(t *testing.T, h *Handler, tunnelID int64, exitStatuses []int) {
|
||||
t.Helper()
|
||||
now := time.Now().UnixMilli()
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, inx, ip_preference, probe_target_host, probe_target_port)
|
||||
VALUES(?, ?, 1, 2, 'tls', 1, ?, ?, 1, ?, '', '', 0)
|
||||
`, tunnelID, fmt.Sprintf("quality-forward-%d", tunnelID), now, now, tunnelID).Error; err != nil {
|
||||
t.Fatalf("insert forwarding tunnel: %v", err)
|
||||
}
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol)
|
||||
VALUES(?, '1', 10, 30001, 'fifo', 1, 'tls')
|
||||
`, tunnelID).Error; err != nil {
|
||||
t.Fatalf("insert entry chain: %v", err)
|
||||
}
|
||||
for i, status := range exitStatuses {
|
||||
nodeID := int64(30 + i)
|
||||
port := 30030 + i
|
||||
ip := fmt.Sprintf("10.0.0.%d", nodeID)
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO node(id, name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx)
|
||||
VALUES(?, ?, ?, ?, ?, '', '30000-30100', '', 'v1', 1, 1, 1, ?, ?, ?, '[::]', '[::]', 0)
|
||||
`, nodeID, fmt.Sprintf("exit-%d", i+1), fmt.Sprintf("exit-secret-%d", i+1), ip, ip, now, now, status).Error; err != nil {
|
||||
t.Fatalf("insert exit node %d: %v", nodeID, err)
|
||||
}
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol)
|
||||
VALUES(?, '3', ?, ?, 'fifo', ?, 'tls')
|
||||
`, tunnelID, nodeID, port, i+1).Error; err != nil {
|
||||
t.Fatalf("insert exit chain %d: %v", nodeID, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
package monitoring
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
const (
|
||||
ConfigTunnelQualityProbeIntervalSec = "monitor_tunnel_quality_interval_sec"
|
||||
DefaultTunnelQualityProbeIntervalSec = 1
|
||||
MinTunnelQualityProbeIntervalSec = 1
|
||||
MaxTunnelQualityProbeIntervalSec = 3600
|
||||
)
|
||||
|
||||
func TunnelQualityProbeIntervalSecondsFromConfigMap(cfg map[string]string) int {
|
||||
if cfg == nil {
|
||||
return DefaultTunnelQualityProbeIntervalSec
|
||||
}
|
||||
seconds, err := parseTunnelQualityProbeIntervalSeconds(cfg[ConfigTunnelQualityProbeIntervalSec])
|
||||
if err != nil {
|
||||
return DefaultTunnelQualityProbeIntervalSec
|
||||
}
|
||||
return seconds
|
||||
}
|
||||
|
||||
func NormalizeTunnelQualityProbeIntervalSeconds(value string) (string, error) {
|
||||
seconds, err := parseTunnelQualityProbeIntervalSeconds(value)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return strconv.Itoa(seconds), nil
|
||||
}
|
||||
|
||||
func parseTunnelQualityProbeIntervalSeconds(value string) (int, error) {
|
||||
trimmed := strings.TrimSpace(value)
|
||||
if trimmed == "" {
|
||||
return 0, fmt.Errorf("隧道质量探测间隔不能为空")
|
||||
}
|
||||
seconds, err := strconv.Atoi(trimmed)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("隧道质量探测间隔必须是整数")
|
||||
}
|
||||
if seconds < MinTunnelQualityProbeIntervalSec || seconds > MaxTunnelQualityProbeIntervalSec {
|
||||
return 0, fmt.Errorf(
|
||||
"隧道质量探测间隔必须在 %d 到 %d 秒之间",
|
||||
MinTunnelQualityProbeIntervalSec,
|
||||
MaxTunnelQualityProbeIntervalSec,
|
||||
)
|
||||
}
|
||||
return seconds, nil
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package monitoring
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestTunnelQualityProbeIntervalSecondsFromConfigMap(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
cfg map[string]string
|
||||
want int
|
||||
}{
|
||||
{name: "missing config", cfg: nil, want: DefaultTunnelQualityProbeIntervalSec},
|
||||
{name: "configured", cfg: map[string]string{ConfigTunnelQualityProbeIntervalSec: "15"}, want: 15},
|
||||
{name: "invalid", cfg: map[string]string{ConfigTunnelQualityProbeIntervalSec: "0"}, want: DefaultTunnelQualityProbeIntervalSec},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if got := TunnelQualityProbeIntervalSecondsFromConfigMap(tt.cfg); got != tt.want {
|
||||
t.Fatalf("interval = %d, want %d", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelQualityProbeIntervalSeconds(t *testing.T) {
|
||||
for _, value := range []string{"1", "15", "3600"} {
|
||||
if got, err := NormalizeTunnelQualityProbeIntervalSeconds(value); err != nil || got != value {
|
||||
t.Fatalf("normalize %q = %q, %v", value, got, err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, value := range []string{"", "0", "3601", "1.5", "abc"} {
|
||||
if got, err := NormalizeTunnelQualityProbeIntervalSeconds(value); err == nil {
|
||||
t.Fatalf("normalize %q unexpectedly succeeded with %q", value, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -132,3 +132,31 @@ func TestUpsertTunnelMetricBucketsIsSafeUnderConcurrency(t *testing.T) {
|
||||
t.Fatalf("expected bytesOut %d, got %d", wantOut, rows[0].BytesOut)
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetLatestTunnelQualitiesIncludesChainDetails(t *testing.T) {
|
||||
r, err := Open(":memory:")
|
||||
if err != nil {
|
||||
t.Fatalf("open repo: %v", err)
|
||||
}
|
||||
defer r.Close()
|
||||
|
||||
if err := r.InsertTunnelQuality(&model.TunnelQuality{
|
||||
TunnelID: 7,
|
||||
Timestamp: time.Now().UnixMilli(),
|
||||
Success: 1,
|
||||
ChainDetails: `{"primaryPath":[],"candidateHops":[{"fromNodeId":10,"toNodeId":31}]}`,
|
||||
}); err != nil {
|
||||
t.Fatalf("insert tunnel quality: %v", err)
|
||||
}
|
||||
|
||||
items, err := r.GetLatestTunnelQualities()
|
||||
if err != nil {
|
||||
t.Fatalf("get latest tunnel qualities: %v", err)
|
||||
}
|
||||
if len(items) != 1 {
|
||||
t.Fatalf("expected one latest tunnel quality, got %+v", items)
|
||||
}
|
||||
if items[0].ChainDetails == "" {
|
||||
t.Fatalf("expected chain details in latest quality row, got %+v", items[0])
|
||||
}
|
||||
}
|
||||
|
||||
@@ -44,7 +44,8 @@ func (r *Repository) GetLatestTunnelQualities() ([]model.TunnelQuality, error) {
|
||||
// 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
|
||||
entry_to_exit_loss, exit_to_bing_loss, success, error_message, timestamp,
|
||||
chain_details
|
||||
FROM (
|
||||
SELECT *, ROW_NUMBER() OVER (PARTITION BY tunnel_id ORDER BY timestamp DESC, id DESC) AS rn
|
||||
FROM tunnel_quality
|
||||
|
||||
@@ -2,6 +2,8 @@ package chain
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
|
||||
"github.com/go-gost/core/chain"
|
||||
"github.com/go-gost/core/hop"
|
||||
@@ -38,11 +40,12 @@ type chainNamer interface {
|
||||
}
|
||||
|
||||
type Chain struct {
|
||||
name string
|
||||
hops []hop.Hop
|
||||
marker selector.Marker
|
||||
metadata metadata.Metadata
|
||||
logger logger.Logger
|
||||
name string
|
||||
hops []hop.Hop
|
||||
ownedHops []hop.Hop
|
||||
marker selector.Marker
|
||||
metadata metadata.Metadata
|
||||
logger logger.Logger
|
||||
}
|
||||
|
||||
func NewChain(name string, opts ...ChainOption) *Chain {
|
||||
@@ -61,8 +64,15 @@ func NewChain(name string, opts ...ChainOption) *Chain {
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Chain) AddHop(hop hop.Hop) {
|
||||
func (c *Chain) AddHop(hop hop.Hop, owned ...bool) {
|
||||
c.hops = append(c.hops, hop)
|
||||
isOwned := true
|
||||
if len(owned) > 0 {
|
||||
isOwned = owned[0]
|
||||
}
|
||||
if isOwned {
|
||||
c.ownedHops = append(c.ownedHops, hop)
|
||||
}
|
||||
}
|
||||
|
||||
// Metadata implements metadata.Metadatable interface.
|
||||
@@ -112,6 +122,36 @@ func (c *Chain) Route(ctx context.Context, network, address string, opts ...chai
|
||||
return rt
|
||||
}
|
||||
|
||||
// Retire gracefully drains resources owned by a chain that has been replaced.
|
||||
func (c *Chain) Retire() {
|
||||
if c == nil {
|
||||
return
|
||||
}
|
||||
for _, h := range c.ownedHops {
|
||||
if retirer, ok := h.(interface{ Retire() }); ok {
|
||||
retirer.Retire()
|
||||
continue
|
||||
}
|
||||
if closer, ok := h.(io.Closer); ok {
|
||||
_ = closer.Close()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Close immediately releases all resources owned by the chain.
|
||||
func (c *Chain) Close() error {
|
||||
if c == nil {
|
||||
return nil
|
||||
}
|
||||
var errs []error
|
||||
for _, h := range c.ownedHops {
|
||||
if closer, ok := h.(io.Closer); ok {
|
||||
errs = append(errs, closer.Close())
|
||||
}
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
type chainGroup struct {
|
||||
chains []chain.Chainer
|
||||
selector selector.Selector[chain.Chainer]
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
package chain
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
corechain "github.com/go-gost/core/chain"
|
||||
corehop "github.com/go-gost/core/hop"
|
||||
)
|
||||
|
||||
type lifecycleTestHop struct {
|
||||
selected int
|
||||
retired int
|
||||
closed int
|
||||
}
|
||||
|
||||
func (h *lifecycleTestHop) Select(context.Context, ...corehop.SelectOption) *corechain.Node {
|
||||
h.selected++
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *lifecycleTestHop) Retire() {
|
||||
h.retired++
|
||||
}
|
||||
|
||||
func (h *lifecycleTestHop) Close() error {
|
||||
h.closed++
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestChainRoutesThroughSharedHopWithoutOwningLifecycle(t *testing.T) {
|
||||
hop := &lifecycleTestHop{}
|
||||
chain := NewChain("shared-hop")
|
||||
chain.AddHop(hop, false)
|
||||
|
||||
if route := chain.Route(context.Background(), "tcp", "example.com:443"); route == nil {
|
||||
t.Fatal("route is nil")
|
||||
}
|
||||
if hop.selected != 1 {
|
||||
t.Fatalf("shared hop selected %d times, want 1", hop.selected)
|
||||
}
|
||||
|
||||
chain.Retire()
|
||||
if err := chain.Close(); err != nil {
|
||||
t.Fatalf("close chain: %v", err)
|
||||
}
|
||||
if hop.retired != 0 || hop.closed != 0 {
|
||||
t.Fatalf("shared hop lifecycle changed: retired=%d closed=%d", hop.retired, hop.closed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChainRetiresAndClosesOwnedHop(t *testing.T) {
|
||||
hop := &lifecycleTestHop{}
|
||||
chain := NewChain("owned-hop")
|
||||
chain.AddHop(hop)
|
||||
|
||||
chain.Retire()
|
||||
if err := chain.Close(); err != nil {
|
||||
t.Fatalf("close chain: %v", err)
|
||||
}
|
||||
if hop.retired != 1 || hop.closed != 1 {
|
||||
t.Fatalf("owned hop lifecycle: retired=%d closed=%d, want 1/1", hop.retired, hop.closed)
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,8 @@ package chain
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"net"
|
||||
|
||||
"github.com/go-gost/core/chain"
|
||||
@@ -102,5 +104,34 @@ func (tr *Transport) Options() *chain.TransportOptions {
|
||||
func (tr *Transport) Copy() chain.Transporter {
|
||||
tr2 := &Transport{}
|
||||
*tr2 = *tr
|
||||
return tr
|
||||
return tr2
|
||||
}
|
||||
|
||||
// Retire prevents long-lived dialer sessions owned by an obsolete chain from
|
||||
// accepting new streams while allowing existing streams to drain.
|
||||
func (tr *Transport) Retire() {
|
||||
if tr == nil {
|
||||
return
|
||||
}
|
||||
if retirer, ok := tr.dialer.(interface{ Retire() }); ok {
|
||||
retirer.Retire()
|
||||
}
|
||||
if retirer, ok := tr.connector.(interface{ Retire() }); ok {
|
||||
retirer.Retire()
|
||||
}
|
||||
}
|
||||
|
||||
// Close immediately releases transport-owned dialer and connector resources.
|
||||
func (tr *Transport) Close() error {
|
||||
if tr == nil {
|
||||
return nil
|
||||
}
|
||||
var errs []error
|
||||
if closer, ok := tr.dialer.(io.Closer); ok {
|
||||
errs = append(errs, closer.Close())
|
||||
}
|
||||
if closer, ok := tr.connector.(io.Closer); ok {
|
||||
errs = append(errs, closer.Close())
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
package chain
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"testing"
|
||||
|
||||
corechain "github.com/go-gost/core/chain"
|
||||
)
|
||||
|
||||
type copyTestRoute struct{}
|
||||
|
||||
func (copyTestRoute) Dial(context.Context, string, string, ...corechain.DialOption) (net.Conn, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (copyTestRoute) Bind(context.Context, string, string, ...corechain.BindOption) (net.Listener, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (copyTestRoute) Nodes() []*corechain.Node {
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestTransportCopyReturnsIndependentTransport(t *testing.T) {
|
||||
originalRoute := copyTestRoute{}
|
||||
replacementRoute := ©TestRoute{}
|
||||
original := NewTransport(nil, nil, corechain.RouteTransportOption(originalRoute))
|
||||
|
||||
copied, ok := original.Copy().(*Transport)
|
||||
if !ok {
|
||||
t.Fatalf("copy type = %T, want *Transport", original.Copy())
|
||||
}
|
||||
if copied == original {
|
||||
t.Fatal("Copy returned the original transport")
|
||||
}
|
||||
|
||||
copied.Options().Route = replacementRoute
|
||||
if original.Options().Route != originalRoute {
|
||||
t.Fatal("mutating copied transport changed original route")
|
||||
}
|
||||
if copied.Options().Route != replacementRoute {
|
||||
t.Fatal("copied transport did not retain its independent route")
|
||||
}
|
||||
}
|
||||
@@ -35,16 +35,18 @@ func ParseChain(cfg *config.ChainConfig, log logger.Logger) (chain.Chainer, erro
|
||||
for _, ch := range cfg.Hops {
|
||||
var hop hop.Hop
|
||||
var err error
|
||||
owned := false
|
||||
|
||||
if ch.Nodes != nil || ch.Plugin != nil {
|
||||
if hop, err = hop_parser.ParseHop(ch, log); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
owned = true
|
||||
} else {
|
||||
hop = registry.HopRegistry().Get(ch.Name)
|
||||
}
|
||||
if hop != nil {
|
||||
c.AddHop(hop)
|
||||
c.AddHop(hop, owned)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -19,20 +19,23 @@ func (session *muxSession) Accept() (net.Conn, error) {
|
||||
}
|
||||
|
||||
func (session *muxSession) Close() error {
|
||||
if session.session == nil {
|
||||
if session == nil || session.session == nil {
|
||||
return nil
|
||||
}
|
||||
return session.session.Close()
|
||||
}
|
||||
|
||||
func (session *muxSession) IsClosed() bool {
|
||||
if session.session == nil {
|
||||
if session == nil || session.session == nil {
|
||||
return true
|
||||
}
|
||||
return session.session.IsClosed()
|
||||
}
|
||||
|
||||
func (session *muxSession) NumStreams() int {
|
||||
if session == nil || session.session == nil {
|
||||
return 0
|
||||
}
|
||||
return session.session.NumStreams()
|
||||
}
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"github.com/go-gost/core/logger"
|
||||
md "github.com/go-gost/core/metadata"
|
||||
kcp_util "github.com/go-gost/x/internal/util/kcp"
|
||||
"github.com/go-gost/x/internal/util/sessionretire"
|
||||
mdutil "github.com/go-gost/x/metadata/util"
|
||||
"github.com/go-gost/x/registry"
|
||||
"github.com/xtaci/kcp-go/v5"
|
||||
@@ -25,6 +26,7 @@ func init() {
|
||||
type kcpDialer struct {
|
||||
sessions map[string]*muxSession
|
||||
sessionMutex sync.Mutex
|
||||
retired bool
|
||||
logger logger.Logger
|
||||
md metadata
|
||||
options dialer.Options
|
||||
@@ -64,6 +66,9 @@ func (d *kcpDialer) Dial(ctx context.Context, addr string, opts ...dialer.DialOp
|
||||
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
session, ok := d.sessions[addr]
|
||||
if session != nil && session.IsClosed() {
|
||||
@@ -171,3 +176,36 @@ func (d *kcpDialer) initSession(ctx context.Context, addr net.Addr, conn net.Pac
|
||||
func (d *kcpDialer) Multiplex() bool {
|
||||
return true
|
||||
}
|
||||
|
||||
// Retire drains existing streams and closes their backing sessions once idle.
|
||||
func (d *kcpDialer) Retire() {
|
||||
for _, session := range d.detachSessions() {
|
||||
sessionretire.Gracefully(session)
|
||||
}
|
||||
}
|
||||
|
||||
// Close immediately releases all cached multiplex sessions.
|
||||
func (d *kcpDialer) Close() error {
|
||||
var errs []error
|
||||
for _, session := range d.detachSessions() {
|
||||
errs = append(errs, session.Close())
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
func (d *kcpDialer) detachSessions() []*muxSession {
|
||||
if d == nil {
|
||||
return nil
|
||||
}
|
||||
d.sessionMutex.Lock()
|
||||
d.retired = true
|
||||
sessions := make([]*muxSession, 0, len(d.sessions))
|
||||
for _, session := range d.sessions {
|
||||
if session != nil {
|
||||
sessions = append(sessions, session)
|
||||
}
|
||||
}
|
||||
d.sessions = make(map[string]*muxSession)
|
||||
d.sessionMutex.Unlock()
|
||||
return sessions
|
||||
}
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
package kcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRetiredDialerRejectsNewConnections(t *testing.T) {
|
||||
dialer := NewDialer().(*kcpDialer)
|
||||
dialer.Retire()
|
||||
|
||||
if _, err := dialer.Dial(context.Background(), "127.0.0.1:1"); !errors.Is(err, net.ErrClosed) {
|
||||
t.Fatalf("Dial error = %v, want net.ErrClosed", err)
|
||||
}
|
||||
}
|
||||
@@ -20,13 +20,24 @@ func (session *muxSession) Accept() (net.Conn, error) {
|
||||
}
|
||||
|
||||
func (session *muxSession) Close() error {
|
||||
if session == nil {
|
||||
return nil
|
||||
}
|
||||
if session.session == nil {
|
||||
if session.conn != nil {
|
||||
conn := session.conn
|
||||
session.conn = nil
|
||||
return conn.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
return session.session.Close()
|
||||
}
|
||||
|
||||
func (session *muxSession) IsClosed() bool {
|
||||
if session == nil {
|
||||
return true
|
||||
}
|
||||
if session.session == nil {
|
||||
return true
|
||||
}
|
||||
@@ -34,5 +45,8 @@ func (session *muxSession) IsClosed() bool {
|
||||
}
|
||||
|
||||
func (session *muxSession) NumStreams() int {
|
||||
if session == nil || session.session == nil {
|
||||
return 0
|
||||
}
|
||||
return session.session.NumStreams()
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"github.com/go-gost/core/logger"
|
||||
md "github.com/go-gost/core/metadata"
|
||||
"github.com/go-gost/x/internal/util/mux"
|
||||
"github.com/go-gost/x/internal/util/sessionretire"
|
||||
"github.com/go-gost/x/registry"
|
||||
)
|
||||
|
||||
@@ -21,6 +22,7 @@ func init() {
|
||||
type mtcpDialer struct {
|
||||
sessions map[string]*muxSession
|
||||
sessionMutex sync.Mutex
|
||||
retired bool
|
||||
logger logger.Logger
|
||||
md metadata
|
||||
options dialer.Options
|
||||
@@ -55,6 +57,9 @@ func (d *mtcpDialer) Multiplex() bool {
|
||||
func (d *mtcpDialer) Dial(ctx context.Context, addr string, opts ...dialer.DialOption) (conn net.Conn, err error) {
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
session, ok := d.sessions[addr]
|
||||
if session != nil && session.IsClosed() {
|
||||
@@ -88,6 +93,10 @@ func (d *mtcpDialer) Handshake(ctx context.Context, conn net.Conn, options ...di
|
||||
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
conn.Close()
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
if d.md.handshakeTimeout > 0 {
|
||||
conn.SetDeadline(time.Now().Add(d.md.handshakeTimeout))
|
||||
@@ -129,3 +138,34 @@ func (d *mtcpDialer) initSession(ctx context.Context, conn net.Conn) (*muxSessio
|
||||
}
|
||||
return &muxSession{conn: conn, session: session}, nil
|
||||
}
|
||||
|
||||
func (d *mtcpDialer) Retire() {
|
||||
for _, session := range d.detachSessions() {
|
||||
sessionretire.Gracefully(session)
|
||||
}
|
||||
}
|
||||
|
||||
func (d *mtcpDialer) Close() error {
|
||||
var errs []error
|
||||
for _, session := range d.detachSessions() {
|
||||
errs = append(errs, session.Close())
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
func (d *mtcpDialer) detachSessions() []*muxSession {
|
||||
if d == nil {
|
||||
return nil
|
||||
}
|
||||
d.sessionMutex.Lock()
|
||||
d.retired = true
|
||||
sessions := make([]*muxSession, 0, len(d.sessions))
|
||||
for _, session := range d.sessions {
|
||||
if session != nil {
|
||||
sessions = append(sessions, session)
|
||||
}
|
||||
}
|
||||
d.sessions = make(map[string]*muxSession)
|
||||
d.sessionMutex.Unlock()
|
||||
return sessions
|
||||
}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
package mtcp
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRetiredDialerRejectsNewConnections(t *testing.T) {
|
||||
dialer := NewDialer().(*mtcpDialer)
|
||||
dialer.Retire()
|
||||
|
||||
if _, err := dialer.Dial(context.Background(), "127.0.0.1:1"); !errors.Is(err, net.ErrClosed) {
|
||||
t.Fatalf("Dial error = %v, want net.ErrClosed", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSessionCloseReleasesPreHandshakeConnection(t *testing.T) {
|
||||
conn, peer := net.Pipe()
|
||||
defer peer.Close()
|
||||
session := &muxSession{conn: conn}
|
||||
|
||||
if err := session.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
if !session.IsClosed() {
|
||||
t.Fatal("pre-handshake session still reports open after Close")
|
||||
}
|
||||
}
|
||||
@@ -20,13 +20,24 @@ func (session *muxSession) Accept() (net.Conn, error) {
|
||||
}
|
||||
|
||||
func (session *muxSession) Close() error {
|
||||
if session == nil {
|
||||
return nil
|
||||
}
|
||||
if session.session == nil {
|
||||
if session.conn != nil {
|
||||
conn := session.conn
|
||||
session.conn = nil
|
||||
return conn.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
return session.session.Close()
|
||||
}
|
||||
|
||||
func (session *muxSession) IsClosed() bool {
|
||||
if session == nil {
|
||||
return true
|
||||
}
|
||||
if session.session == nil {
|
||||
return true
|
||||
}
|
||||
@@ -34,5 +45,8 @@ func (session *muxSession) IsClosed() bool {
|
||||
}
|
||||
|
||||
func (session *muxSession) NumStreams() int {
|
||||
if session == nil || session.session == nil {
|
||||
return 0
|
||||
}
|
||||
return session.session.NumStreams()
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"github.com/go-gost/core/logger"
|
||||
md "github.com/go-gost/core/metadata"
|
||||
"github.com/go-gost/x/internal/util/mux"
|
||||
"github.com/go-gost/x/internal/util/sessionretire"
|
||||
"github.com/go-gost/x/registry"
|
||||
)
|
||||
|
||||
@@ -22,6 +23,7 @@ func init() {
|
||||
type mtlsDialer struct {
|
||||
sessions map[string]*muxSession
|
||||
sessionMutex sync.Mutex
|
||||
retired bool
|
||||
logger logger.Logger
|
||||
md metadata
|
||||
options dialer.Options
|
||||
@@ -56,6 +58,9 @@ func (d *mtlsDialer) Multiplex() bool {
|
||||
func (d *mtlsDialer) Dial(ctx context.Context, addr string, opts ...dialer.DialOption) (conn net.Conn, err error) {
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
session, ok := d.sessions[addr]
|
||||
if session != nil && session.IsClosed() {
|
||||
@@ -89,6 +94,10 @@ func (d *mtlsDialer) Handshake(ctx context.Context, conn net.Conn, options ...di
|
||||
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
conn.Close()
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
if d.md.handshakeTimeout > 0 {
|
||||
conn.SetDeadline(time.Now().Add(d.md.handshakeTimeout))
|
||||
@@ -136,3 +145,34 @@ func (d *mtlsDialer) initSession(ctx context.Context, conn net.Conn) (*muxSessio
|
||||
}
|
||||
return &muxSession{conn: conn, session: session}, nil
|
||||
}
|
||||
|
||||
func (d *mtlsDialer) Retire() {
|
||||
for _, session := range d.detachSessions() {
|
||||
sessionretire.Gracefully(session)
|
||||
}
|
||||
}
|
||||
|
||||
func (d *mtlsDialer) Close() error {
|
||||
var errs []error
|
||||
for _, session := range d.detachSessions() {
|
||||
errs = append(errs, session.Close())
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
func (d *mtlsDialer) detachSessions() []*muxSession {
|
||||
if d == nil {
|
||||
return nil
|
||||
}
|
||||
d.sessionMutex.Lock()
|
||||
d.retired = true
|
||||
sessions := make([]*muxSession, 0, len(d.sessions))
|
||||
for _, session := range d.sessions {
|
||||
if session != nil {
|
||||
sessions = append(sessions, session)
|
||||
}
|
||||
}
|
||||
d.sessions = make(map[string]*muxSession)
|
||||
d.sessionMutex.Unlock()
|
||||
return sessions
|
||||
}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
package mtls
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRetiredDialerRejectsNewConnections(t *testing.T) {
|
||||
dialer := NewDialer().(*mtlsDialer)
|
||||
dialer.Retire()
|
||||
|
||||
if _, err := dialer.Dial(context.Background(), "127.0.0.1:1"); !errors.Is(err, net.ErrClosed) {
|
||||
t.Fatalf("Dial error = %v, want net.ErrClosed", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSessionCloseReleasesPreHandshakeConnection(t *testing.T) {
|
||||
conn, peer := net.Pipe()
|
||||
defer peer.Close()
|
||||
session := &muxSession{conn: conn}
|
||||
|
||||
if err := session.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
if !session.IsClosed() {
|
||||
t.Fatal("pre-handshake session still reports open after Close")
|
||||
}
|
||||
}
|
||||
@@ -20,13 +20,24 @@ func (session *muxSession) Accept() (net.Conn, error) {
|
||||
}
|
||||
|
||||
func (session *muxSession) Close() error {
|
||||
if session == nil {
|
||||
return nil
|
||||
}
|
||||
if session.session == nil {
|
||||
if session.conn != nil {
|
||||
conn := session.conn
|
||||
session.conn = nil
|
||||
return conn.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
return session.session.Close()
|
||||
}
|
||||
|
||||
func (session *muxSession) IsClosed() bool {
|
||||
if session == nil {
|
||||
return true
|
||||
}
|
||||
if session.session == nil {
|
||||
return true
|
||||
}
|
||||
@@ -34,5 +45,8 @@ func (session *muxSession) IsClosed() bool {
|
||||
}
|
||||
|
||||
func (session *muxSession) NumStreams() int {
|
||||
if session == nil || session.session == nil {
|
||||
return 0
|
||||
}
|
||||
return session.session.NumStreams()
|
||||
}
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"github.com/go-gost/core/logger"
|
||||
md "github.com/go-gost/core/metadata"
|
||||
"github.com/go-gost/x/internal/util/mux"
|
||||
"github.com/go-gost/x/internal/util/sessionretire"
|
||||
ws_util "github.com/go-gost/x/internal/util/ws"
|
||||
"github.com/go-gost/x/registry"
|
||||
"github.com/gorilla/websocket"
|
||||
@@ -25,6 +26,7 @@ func init() {
|
||||
type mwsDialer struct {
|
||||
sessions map[string]*muxSession
|
||||
sessionMutex sync.Mutex
|
||||
retired bool
|
||||
tlsEnabled bool
|
||||
md metadata
|
||||
options dialer.Options
|
||||
@@ -70,6 +72,9 @@ func (d *mwsDialer) Multiplex() bool {
|
||||
func (d *mwsDialer) Dial(ctx context.Context, addr string, opts ...dialer.DialOption) (conn net.Conn, err error) {
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
session, ok := d.sessions[addr]
|
||||
if session != nil && session.IsClosed() {
|
||||
@@ -108,6 +113,10 @@ func (d *mwsDialer) Handshake(ctx context.Context, conn net.Conn, options ...dia
|
||||
|
||||
d.sessionMutex.Lock()
|
||||
defer d.sessionMutex.Unlock()
|
||||
if d.retired {
|
||||
conn.Close()
|
||||
return nil, net.ErrClosed
|
||||
}
|
||||
|
||||
session, ok := d.sessions[opts.Addr]
|
||||
if session != nil && session.conn != conn {
|
||||
@@ -208,3 +217,34 @@ func (d *mwsDialer) keepAlive(conn ws_util.WebsocketConn) {
|
||||
conn.SetWriteDeadline(time.Time{})
|
||||
}
|
||||
}
|
||||
|
||||
func (d *mwsDialer) Retire() {
|
||||
for _, session := range d.detachSessions() {
|
||||
sessionretire.Gracefully(session)
|
||||
}
|
||||
}
|
||||
|
||||
func (d *mwsDialer) Close() error {
|
||||
var errs []error
|
||||
for _, session := range d.detachSessions() {
|
||||
errs = append(errs, session.Close())
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
func (d *mwsDialer) detachSessions() []*muxSession {
|
||||
if d == nil {
|
||||
return nil
|
||||
}
|
||||
d.sessionMutex.Lock()
|
||||
d.retired = true
|
||||
sessions := make([]*muxSession, 0, len(d.sessions))
|
||||
for _, session := range d.sessions {
|
||||
if session != nil {
|
||||
sessions = append(sessions, session)
|
||||
}
|
||||
}
|
||||
d.sessions = make(map[string]*muxSession)
|
||||
d.sessionMutex.Unlock()
|
||||
return sessions
|
||||
}
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
package mws
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRetiredDialerRejectsNewConnections(t *testing.T) {
|
||||
dialer := NewDialer().(*mwsDialer)
|
||||
dialer.Retire()
|
||||
|
||||
if _, err := dialer.Dial(context.Background(), "127.0.0.1:1"); !errors.Is(err, net.ErrClosed) {
|
||||
t.Fatalf("Dial error = %v, want net.ErrClosed", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSessionCloseReleasesPreHandshakeConnection(t *testing.T) {
|
||||
conn, peer := net.Pipe()
|
||||
defer peer.Close()
|
||||
session := &muxSession{conn: conn}
|
||||
|
||||
if err := session.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
if !session.IsClosed() {
|
||||
t.Fatal("pre-handshake session still reports open after Close")
|
||||
}
|
||||
}
|
||||
+49
-8
@@ -3,6 +3,7 @@ package hop
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net"
|
||||
"sort"
|
||||
@@ -92,6 +93,7 @@ type chainHop struct {
|
||||
nodes []*chain.Node
|
||||
mu sync.RWMutex
|
||||
cancelFunc context.CancelFunc
|
||||
stopOnce sync.Once
|
||||
options options
|
||||
}
|
||||
|
||||
@@ -383,13 +385,52 @@ func (p *chainHop) parseNode(r io.Reader) ([]*chain.Node, error) {
|
||||
return nodes, nil
|
||||
}
|
||||
|
||||
func (p *chainHop) Close() error {
|
||||
p.cancelFunc()
|
||||
if p.options.fileLoader != nil {
|
||||
p.options.fileLoader.Close()
|
||||
func (p *chainHop) stopReload() {
|
||||
if p == nil {
|
||||
return
|
||||
}
|
||||
if p.options.redisLoader != nil {
|
||||
p.options.redisLoader.Close()
|
||||
}
|
||||
return nil
|
||||
p.stopOnce.Do(func() {
|
||||
p.cancelFunc()
|
||||
if p.options.fileLoader != nil {
|
||||
p.options.fileLoader.Close()
|
||||
}
|
||||
if p.options.redisLoader != nil {
|
||||
p.options.redisLoader.Close()
|
||||
}
|
||||
if p.options.httpLoader != nil {
|
||||
p.options.httpLoader.Close()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func (p *chainHop) Retire() {
|
||||
if p == nil {
|
||||
return
|
||||
}
|
||||
p.stopReload()
|
||||
for _, node := range p.Nodes() {
|
||||
if node == nil || node.Options().Transport == nil {
|
||||
continue
|
||||
}
|
||||
if retirer, ok := node.Options().Transport.(interface{ Retire() }); ok {
|
||||
retirer.Retire()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (p *chainHop) Close() error {
|
||||
if p == nil {
|
||||
return nil
|
||||
}
|
||||
p.stopReload()
|
||||
var errs []error
|
||||
for _, node := range p.Nodes() {
|
||||
if node == nil || node.Options().Transport == nil {
|
||||
continue
|
||||
}
|
||||
if closer, ok := node.Options().Transport.(io.Closer); ok {
|
||||
errs = append(errs, closer.Close())
|
||||
}
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
package sessionretire
|
||||
|
||||
import "time"
|
||||
|
||||
const (
|
||||
defaultIdleGrace = time.Second
|
||||
defaultPollPeriod = 100 * time.Millisecond
|
||||
)
|
||||
|
||||
// Session is the lifecycle surface shared by the multiplexed dialers.
|
||||
type Session interface {
|
||||
Close() error
|
||||
IsClosed() bool
|
||||
NumStreams() int
|
||||
}
|
||||
|
||||
// Gracefully closes a retired session after all existing streams have drained.
|
||||
// A short idle grace covers the Dial/Handshake hand-off used by several dialers.
|
||||
func Gracefully(session Session) {
|
||||
if session == nil {
|
||||
return
|
||||
}
|
||||
go waitUntilIdle(session, defaultIdleGrace, defaultPollPeriod)
|
||||
}
|
||||
|
||||
func waitUntilIdle(session Session, idleGrace, pollPeriod time.Duration) {
|
||||
if session == nil {
|
||||
return
|
||||
}
|
||||
if idleGrace <= 0 {
|
||||
idleGrace = defaultIdleGrace
|
||||
}
|
||||
if pollPeriod <= 0 {
|
||||
pollPeriod = defaultPollPeriod
|
||||
}
|
||||
|
||||
ticker := time.NewTicker(pollPeriod)
|
||||
defer ticker.Stop()
|
||||
|
||||
var idleSince time.Time
|
||||
for {
|
||||
if session.IsClosed() {
|
||||
_ = session.Close()
|
||||
return
|
||||
}
|
||||
if session.NumStreams() == 0 {
|
||||
if idleSince.IsZero() {
|
||||
idleSince = time.Now()
|
||||
} else if time.Since(idleSince) >= idleGrace {
|
||||
_ = session.Close()
|
||||
return
|
||||
}
|
||||
} else {
|
||||
idleSince = time.Time{}
|
||||
}
|
||||
<-ticker.C
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package sessionretire
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
type testSession struct {
|
||||
mu sync.Mutex
|
||||
streams int
|
||||
closed bool
|
||||
}
|
||||
|
||||
func (s *testSession) Close() error {
|
||||
s.mu.Lock()
|
||||
s.closed = true
|
||||
s.mu.Unlock()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *testSession) IsClosed() bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.closed
|
||||
}
|
||||
|
||||
func (s *testSession) NumStreams() int {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.streams
|
||||
}
|
||||
|
||||
func TestWaitUntilIdlePreservesActiveStreams(t *testing.T) {
|
||||
session := &testSession{streams: 1}
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
waitUntilIdle(session, 20*time.Millisecond, time.Millisecond)
|
||||
close(done)
|
||||
}()
|
||||
|
||||
time.Sleep(30 * time.Millisecond)
|
||||
if session.IsClosed() {
|
||||
t.Fatal("active session was closed")
|
||||
}
|
||||
|
||||
session.mu.Lock()
|
||||
session.streams = 0
|
||||
session.mu.Unlock()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(250 * time.Millisecond):
|
||||
t.Fatal("idle session was not closed")
|
||||
}
|
||||
if !session.IsClosed() {
|
||||
t.Fatal("retired session did not close after becoming idle")
|
||||
}
|
||||
}
|
||||
@@ -28,7 +28,13 @@ func (r *chainRegistry) Register(name string, v chain.Chainer) error {
|
||||
}
|
||||
|
||||
func (r *chainRegistry) replace(name string, v chain.Chainer) {
|
||||
r.m.Store(name, v)
|
||||
old, loaded := r.m.Swap(name, v)
|
||||
if !loaded {
|
||||
return
|
||||
}
|
||||
if retirer, ok := old.(interface{ Retire() }); ok {
|
||||
retirer.Retire()
|
||||
}
|
||||
}
|
||||
|
||||
func (r *chainRegistry) Get(name string) chain.Chainer {
|
||||
|
||||
@@ -16,6 +16,15 @@ func (c testChainer) Route(context.Context, string, string, ...chain.RouteOption
|
||||
return c.route
|
||||
}
|
||||
|
||||
type retiringTestChainer struct {
|
||||
testChainer
|
||||
retired bool
|
||||
}
|
||||
|
||||
func (c *retiringTestChainer) Retire() {
|
||||
c.retired = true
|
||||
}
|
||||
|
||||
type testRoute struct {
|
||||
nodes []*chain.Node
|
||||
}
|
||||
@@ -49,3 +58,20 @@ func TestReplaceChainOverwritesExistingRegistration(t *testing.T) {
|
||||
t.Fatalf("expected replacement chain route, got %#v", route)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReplaceChainRetiresPreviousRegistration(t *testing.T) {
|
||||
name := "replace_chain_retire_tdd"
|
||||
ChainRegistry().Unregister(name)
|
||||
defer ChainRegistry().Unregister(name)
|
||||
|
||||
old := &retiringTestChainer{}
|
||||
if err := ChainRegistry().Register(name, old); err != nil {
|
||||
t.Fatalf("register old chain: %v", err)
|
||||
}
|
||||
if err := ReplaceChain(name, testChainer{}); err != nil {
|
||||
t.Fatalf("replace chain: %v", err)
|
||||
}
|
||||
if !old.retired {
|
||||
t.Fatal("previous chain was not retired")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -151,6 +151,7 @@ const (
|
||||
initialBackoff = 2 * time.Second // 重连初始退避
|
||||
maxBackoff = 2 * time.Minute // 重连最大退避
|
||||
defaultMetricReportInterval = 5 * time.Second
|
||||
maxConcurrentTCPPings = 8
|
||||
)
|
||||
|
||||
type WebSocketReporter struct {
|
||||
@@ -172,6 +173,7 @@ type WebSocketReporter struct {
|
||||
connecting bool // 正在连接状态
|
||||
connMutex sync.Mutex // 连接状态锁
|
||||
aesCrypto *crypto.AESCrypto // AES加密器
|
||||
tcpPingSem chan struct{} // 限制诊断探测并发,避免离线目标耗尽连接
|
||||
}
|
||||
|
||||
var wsDial = func(dialer *websocket.Dialer, rawURL string) (*websocket.Conn, *http.Response, error) {
|
||||
@@ -201,6 +203,29 @@ func NewWebSocketReporter(serverURL string, secret string) *WebSocketReporter {
|
||||
connected: false,
|
||||
connecting: false,
|
||||
aesCrypto: aesCrypto,
|
||||
tcpPingSem: make(chan struct{}, maxConcurrentTCPPings),
|
||||
}
|
||||
}
|
||||
|
||||
func (w *WebSocketReporter) tryAcquireTCPPingSlot() bool {
|
||||
if w == nil || w.tcpPingSem == nil {
|
||||
return false
|
||||
}
|
||||
select {
|
||||
case w.tcpPingSem <- struct{}{}:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
func (w *WebSocketReporter) releaseTCPPingSlot() {
|
||||
if w == nil || w.tcpPingSem == nil {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-w.tcpPingSem:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
@@ -840,9 +865,14 @@ func (w *WebSocketReporter) routeCommand(cmd CommandMessage) {
|
||||
|
||||
// TCP Ping 诊断命令(只读,不需要保存配置)
|
||||
case "TcpPing":
|
||||
response.Type = "TcpPingResponse"
|
||||
if !w.tryAcquireTCPPingSlot() {
|
||||
err = fmt.Errorf("TCP探测任务过多,请稍后重试")
|
||||
break
|
||||
}
|
||||
defer w.releaseTCPPingSlot()
|
||||
var tcpPingResult TcpPingResponse
|
||||
tcpPingResult, err = w.handleTcpPing(cmd.Data)
|
||||
response.Type = "TcpPingResponse"
|
||||
response.Data = tcpPingResult
|
||||
// needSaveConfig = false (默认值)
|
||||
|
||||
|
||||
@@ -148,6 +148,25 @@ func TestNewWebSocketReporterUsesReducedMetricInterval(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestWebSocketReporterLimitsConcurrentTCPPings(t *testing.T) {
|
||||
reporter := &WebSocketReporter{tcpPingSem: make(chan struct{}, maxConcurrentTCPPings)}
|
||||
for i := 0; i < maxConcurrentTCPPings; i++ {
|
||||
if !reporter.tryAcquireTCPPingSlot() {
|
||||
t.Fatalf("expected TCP ping slot %d to be available", i)
|
||||
}
|
||||
}
|
||||
if reporter.tryAcquireTCPPingSlot() {
|
||||
t.Fatalf("expected TCP ping concurrency limit at %d", maxConcurrentTCPPings)
|
||||
}
|
||||
for i := 0; i < maxConcurrentTCPPings; i++ {
|
||||
reporter.releaseTCPPingSlot()
|
||||
}
|
||||
if !reporter.tryAcquireTCPPingSlot() {
|
||||
t.Fatalf("expected released TCP ping slot to be reusable")
|
||||
}
|
||||
reporter.releaseTCPPingSlot()
|
||||
}
|
||||
|
||||
func TestFormatWebSocketDialErrorIncludesHTTPStatus(t *testing.T) {
|
||||
err := errors.New("websocket: bad handshake")
|
||||
resp := &http.Response{
|
||||
|
||||
+247
-40
@@ -1,4 +1,31 @@
|
||||
#!/bin/bash
|
||||
#!/bin/sh
|
||||
# shellcheck shell=bash
|
||||
|
||||
# Alpine 默认不带 Bash。先用系统自带的 /bin/sh 安装/切换到 Bash,
|
||||
# 后续主体继续使用 Bash 语法,避免要求用户手动准备运行环境。
|
||||
if [ -z "${BASH_VERSION:-}" ]; then
|
||||
if command -v bash >/dev/null 2>&1; then
|
||||
exec bash "$0" "$@"
|
||||
fi
|
||||
|
||||
if [ -f /etc/alpine-release ] && command -v apk >/dev/null 2>&1; then
|
||||
if [ "$(id -u)" -eq 0 ]; then
|
||||
apk add --no-cache bash
|
||||
elif command -v sudo >/dev/null 2>&1; then
|
||||
sudo apk add --no-cache bash
|
||||
elif command -v doas >/dev/null 2>&1; then
|
||||
doas apk add --no-cache bash
|
||||
else
|
||||
echo "❌ Alpine 安装需要 root 权限,或已配置 sudo/doas。" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
exec bash "$0" "$@"
|
||||
fi
|
||||
|
||||
echo "❌ 此安装脚本需要 Bash。" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# GitHub repo used for release downloads
|
||||
REPO="Sagit-chu/flux-panel"
|
||||
@@ -24,11 +51,14 @@ get_architecture() {
|
||||
|
||||
# 安装目录
|
||||
INSTALL_DIR="/etc/flux_agent"
|
||||
FLUX_AGENT_SYSTEMD_SERVICE_FILE="/etc/systemd/system/flux_agent.service"
|
||||
FLUX_AGENT_OPENRC_SERVICE_FILE="/etc/init.d/flux_agent"
|
||||
LEGACY_GOST_BINARY="/usr/local/bin/gost"
|
||||
LEGACY_GOST_CONFIG_DIR="/etc/gost"
|
||||
LEGACY_GOST_SERVICE_FILE_ETC="/etc/systemd/system/gost.service"
|
||||
LEGACY_GOST_SERVICE_FILE_LIB="/lib/systemd/system/gost.service"
|
||||
LEGACY_GOST_SERVICE_FILE_USR_LIB="/usr/lib/systemd/system/gost.service"
|
||||
SERVICE_MANAGER="${SERVICE_MANAGER:-}"
|
||||
|
||||
# 镜像加速配置(可由面板传入或交互式询问)
|
||||
PROXY_ENABLED="${PROXY_ENABLED:-}"
|
||||
@@ -256,6 +286,199 @@ write_flux_agent_config() {
|
||||
"$(json_escape "$SECRET")" > "$path"
|
||||
}
|
||||
|
||||
ensure_service_manager() {
|
||||
if [[ -n "$SERVICE_MANAGER" ]]; then
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd|openrc)
|
||||
return 0
|
||||
;;
|
||||
*)
|
||||
echo "❌ 不支持的服务管理器: $SERVICE_MANAGER" >&2
|
||||
return 1
|
||||
;;
|
||||
esac
|
||||
fi
|
||||
|
||||
if command -v systemctl >/dev/null 2>&1 && [[ -d /run/systemd/system ]]; then
|
||||
SERVICE_MANAGER="systemd"
|
||||
return 0
|
||||
fi
|
||||
if command -v rc-service >/dev/null 2>&1 && command -v rc-update >/dev/null 2>&1; then
|
||||
SERVICE_MANAGER="openrc"
|
||||
return 0
|
||||
fi
|
||||
|
||||
echo "❌ 未检测到受支持的服务管理器(systemd 或 OpenRC)。" >&2
|
||||
return 1
|
||||
}
|
||||
|
||||
flux_agent_service_exists() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
[[ -f "$FLUX_AGENT_SYSTEMD_SERVICE_FILE" ]] || \
|
||||
systemctl list-units --full -all 2>/dev/null | grep -Fq "flux_agent.service"
|
||||
;;
|
||||
openrc)
|
||||
[[ -f "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]]
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
stop_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl stop flux_agent 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent stop 2>/dev/null || true
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
disable_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl disable flux_agent 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rc-update del flux_agent default 2>/dev/null || true
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
write_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
mkdir -p "$(dirname "$FLUX_AGENT_SYSTEMD_SERVICE_FILE")"
|
||||
cat > "$FLUX_AGENT_SYSTEMD_SERVICE_FILE" <<EOF
|
||||
[Unit]
|
||||
Description=Flux_agent Proxy Service
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
WorkingDirectory=$INSTALL_DIR
|
||||
ExecStart=$INSTALL_DIR/flux_agent
|
||||
Restart=on-failure
|
||||
StandardOutput=null
|
||||
StandardError=null
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
EOF
|
||||
;;
|
||||
openrc)
|
||||
mkdir -p "$(dirname "$FLUX_AGENT_OPENRC_SERVICE_FILE")"
|
||||
cat > "$FLUX_AGENT_OPENRC_SERVICE_FILE" <<EOF
|
||||
#!/sbin/openrc-run
|
||||
|
||||
name="flux_agent"
|
||||
description="Flux_agent Proxy Service"
|
||||
command="$INSTALL_DIR/flux_agent"
|
||||
directory="$INSTALL_DIR"
|
||||
command_background="yes"
|
||||
pidfile="/run/\${RC_SVCNAME}.pid"
|
||||
output_log="/dev/null"
|
||||
error_log="/dev/null"
|
||||
|
||||
depend() {
|
||||
need net
|
||||
}
|
||||
EOF
|
||||
chmod +x "$FLUX_AGENT_OPENRC_SERVICE_FILE"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
enable_and_start_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl daemon-reload
|
||||
systemctl enable flux_agent
|
||||
systemctl start flux_agent
|
||||
;;
|
||||
openrc)
|
||||
rc-update add flux_agent default
|
||||
rc-service flux_agent start
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
start_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl start flux_agent
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent start
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
flux_agent_service_is_active() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl is-active --quiet flux_agent
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent status >/dev/null 2>&1
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
flux_agent_service_status() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl is-active flux_agent 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent status 2>/dev/null || true
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
remove_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
rm -f "$FLUX_AGENT_SYSTEMD_SERVICE_FILE"
|
||||
systemctl daemon-reload 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rm -f "$FLUX_AGENT_OPENRC_SERVICE_FILE"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
flux_agent_service_status_hint() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
echo "systemctl status flux_agent --no-pager"
|
||||
;;
|
||||
openrc)
|
||||
echo "rc-service flux_agent status"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
cleanup_legacy_gost_installation() {
|
||||
local matched_service_files=()
|
||||
local service_file=""
|
||||
@@ -341,19 +564,21 @@ install_flux_agent() {
|
||||
|
||||
get_config_params
|
||||
|
||||
# 检查并安装 tcpkill
|
||||
# 检查并安装 tcpkill
|
||||
check_and_install_tcpkill
|
||||
|
||||
ensure_service_manager || exit 1
|
||||
|
||||
mkdir -p "$INSTALL_DIR"
|
||||
|
||||
local tmp_binary="$INSTALL_DIR/flux_agent.new"
|
||||
|
||||
# 停止并禁用已有服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
if flux_agent_service_exists; then
|
||||
echo "🔍 检测到已存在的flux_agent服务"
|
||||
systemctl stop flux_agent 2>/dev/null && echo "🛑 停止服务"
|
||||
systemctl disable flux_agent 2>/dev/null && echo "🚫 禁用自启"
|
||||
stop_flux_agent_service
|
||||
echo "🛑 停止服务"
|
||||
disable_flux_agent_service
|
||||
echo "🚫 禁用自启"
|
||||
fi
|
||||
|
||||
# 下载 flux_agent
|
||||
@@ -392,38 +617,21 @@ EOF
|
||||
# 加强权限
|
||||
chmod 600 "$INSTALL_DIR"/*.json
|
||||
|
||||
# 创建 systemd 服务
|
||||
SERVICE_FILE="/etc/systemd/system/flux_agent.service"
|
||||
cat > "$SERVICE_FILE" <<EOF
|
||||
[Unit]
|
||||
Description=Flux_agent Proxy Service
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
WorkingDirectory=$INSTALL_DIR
|
||||
ExecStart=$INSTALL_DIR/flux_agent
|
||||
Restart=on-failure
|
||||
StandardOutput=null
|
||||
StandardError=null
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
EOF
|
||||
# 创建 systemd 或 OpenRC 服务
|
||||
write_flux_agent_service
|
||||
|
||||
# 启动服务
|
||||
systemctl daemon-reload
|
||||
systemctl enable flux_agent
|
||||
systemctl start flux_agent
|
||||
enable_and_start_flux_agent_service
|
||||
|
||||
# 检查状态
|
||||
echo "🔄 检查服务状态..."
|
||||
if systemctl is-active --quiet flux_agent; then
|
||||
if flux_agent_service_is_active; then
|
||||
echo "✅ 安装完成,flux_agent服务已启动并设置为开机启动。"
|
||||
echo "📁 配置目录: $INSTALL_DIR"
|
||||
echo "🔧 服务状态: $(systemctl is-active flux_agent)"
|
||||
echo "🔧 服务状态: $(flux_agent_service_status)"
|
||||
else
|
||||
echo "❌ flux_agent服务启动失败,请执行以下命令查看状态:"
|
||||
echo "systemctl status flux_agent --no-pager"
|
||||
flux_agent_service_status_hint
|
||||
fi
|
||||
}
|
||||
|
||||
@@ -443,6 +651,7 @@ update_flux_agent() {
|
||||
|
||||
# 检查并安装 tcpkill
|
||||
check_and_install_tcpkill
|
||||
ensure_service_manager || return 1
|
||||
|
||||
# 先下载新版本
|
||||
echo "⬇️ 下载最新版本..."
|
||||
@@ -455,9 +664,9 @@ update_flux_agent() {
|
||||
cleanup_legacy_gost_installation
|
||||
|
||||
# 停止服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
if flux_agent_service_exists; then
|
||||
echo "🛑 停止 flux_agent 服务..."
|
||||
systemctl stop flux_agent
|
||||
stop_flux_agent_service
|
||||
fi
|
||||
|
||||
# 替换文件
|
||||
@@ -469,7 +678,7 @@ update_flux_agent() {
|
||||
|
||||
# 重启服务
|
||||
echo "🔄 重启服务..."
|
||||
systemctl start flux_agent
|
||||
start_flux_agent_service
|
||||
|
||||
echo "✅ 更新完成,服务已重新启动。"
|
||||
}
|
||||
@@ -477,6 +686,7 @@ update_flux_agent() {
|
||||
# 卸载功能
|
||||
uninstall_flux_agent() {
|
||||
echo "🗑️ 开始卸载 flux_agent..."
|
||||
ensure_service_manager || return 1
|
||||
|
||||
read -p "确认卸载 flux_agent 吗?此操作将删除所有相关文件 (y/N): " confirm
|
||||
if [[ "$confirm" != "y" && "$confirm" != "Y" ]]; then
|
||||
@@ -485,15 +695,15 @@ uninstall_flux_agent() {
|
||||
fi
|
||||
|
||||
# 停止并禁用服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
if flux_agent_service_exists; then
|
||||
echo "🛑 停止并禁用服务..."
|
||||
systemctl stop flux_agent 2>/dev/null
|
||||
systemctl disable flux_agent 2>/dev/null
|
||||
stop_flux_agent_service
|
||||
disable_flux_agent_service
|
||||
fi
|
||||
|
||||
# 删除服务文件
|
||||
if [[ -f "/etc/systemd/system/flux_agent.service" ]]; then
|
||||
rm -f "/etc/systemd/system/flux_agent.service"
|
||||
if [[ -f "$FLUX_AGENT_SYSTEMD_SERVICE_FILE" || -f "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]]; then
|
||||
remove_flux_agent_service
|
||||
echo "🧹 删除服务文件"
|
||||
fi
|
||||
|
||||
@@ -503,9 +713,6 @@ uninstall_flux_agent() {
|
||||
echo "🧹 删除安装目录: $INSTALL_DIR"
|
||||
fi
|
||||
|
||||
# 重载 systemd
|
||||
systemctl daemon-reload
|
||||
|
||||
echo "✅ 卸载完成"
|
||||
}
|
||||
|
||||
|
||||
@@ -90,6 +90,7 @@ test_update_flux_agent_asks_for_proxy_config() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
@@ -165,6 +166,7 @@ test_install_flux_agent_preserves_legacy_gost_when_download_fails() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
@@ -199,6 +201,7 @@ test_update_flux_agent_preserves_legacy_gost_when_download_fails() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
@@ -236,7 +239,9 @@ test_install_flux_agent_writes_json_safe_config() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
FLUX_AGENT_SYSTEMD_SERVICE_FILE="$INSTALL_DIR/flux_agent.service"
|
||||
SERVER_ADDR='panel"addr'
|
||||
SECRET='sec\ret"1'
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
@@ -274,6 +279,117 @@ EOF
|
||||
assert_equals "$expected" "$actual" "install_flux_agent should JSON-escape config values"
|
||||
)
|
||||
|
||||
test_install_script_bootstraps_bash_for_alpine() (
|
||||
set -euo pipefail
|
||||
|
||||
local shebang
|
||||
shebang=$(head -n 1 "$ROOT_DIR/install.sh")
|
||||
assert_equals "#!/bin/sh" "$shebang" "install.sh should start with Alpine's default shell"
|
||||
|
||||
grep -Fq 'apk add --no-cache bash' "$ROOT_DIR/install.sh" || \
|
||||
fail "install.sh should bootstrap Bash through apk on Alpine"
|
||||
)
|
||||
|
||||
test_install_flux_agent_uses_openrc() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
local temp_root
|
||||
temp_root=$(mktemp -d)
|
||||
INSTALL_DIR="$temp_root/flux_agent"
|
||||
FLUX_AGENT_OPENRC_SERVICE_FILE="$temp_root/init.d/flux_agent"
|
||||
SERVICE_MANAGER="openrc"
|
||||
SERVER_ADDR="panel.example.com:443"
|
||||
SECRET="secret"
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
|
||||
local rc_service_calls=""
|
||||
local rc_update_calls=""
|
||||
|
||||
ask_proxy_config() { :; }
|
||||
ensure_download_url_initialized() { :; }
|
||||
get_config_params() { :; }
|
||||
check_and_install_tcpkill() { :; }
|
||||
cleanup_legacy_gost_installation() { :; }
|
||||
curl() {
|
||||
local output=""
|
||||
while [[ $# -gt 0 ]]; do
|
||||
if [[ "$1" == "-o" ]]; then
|
||||
output="$2"
|
||||
shift 2
|
||||
continue
|
||||
fi
|
||||
shift
|
||||
done
|
||||
|
||||
cat > "$output" <<'EOF'
|
||||
#!/bin/sh
|
||||
echo "new version"
|
||||
EOF
|
||||
chmod +x "$output"
|
||||
}
|
||||
rc-service() {
|
||||
rc_service_calls+=$'\n'"$*"
|
||||
if [[ "$2" == "status" ]]; then
|
||||
echo "status: started"
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
rc-update() {
|
||||
rc_update_calls+=$'\n'"$*"
|
||||
return 0
|
||||
}
|
||||
|
||||
install_flux_agent >/dev/null
|
||||
|
||||
[[ -x "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]] || fail "OpenRC service file should be executable"
|
||||
grep -Fq '#!/sbin/openrc-run' "$FLUX_AGENT_OPENRC_SERVICE_FILE" || \
|
||||
fail "OpenRC service should use openrc-run"
|
||||
grep -Fq "command=\"$INSTALL_DIR/flux_agent\"" "$FLUX_AGENT_OPENRC_SERVICE_FILE" || \
|
||||
fail "OpenRC service should launch the installed flux_agent binary"
|
||||
grep -Fq 'command_background="yes"' "$FLUX_AGENT_OPENRC_SERVICE_FILE" || \
|
||||
fail "OpenRC service should run flux_agent in the background"
|
||||
if command -v openrc-run >/dev/null 2>&1; then
|
||||
"$FLUX_AGENT_OPENRC_SERVICE_FILE" describe >/dev/null 2>&1
|
||||
fi
|
||||
[[ "$rc_update_calls" == *"add flux_agent default"* ]] || \
|
||||
fail "OpenRC install should enable flux_agent in the default runlevel"
|
||||
[[ "$rc_service_calls" == *"start"* ]] || fail "OpenRC install should start flux_agent"
|
||||
[[ "$rc_service_calls" == *"status"* ]] || fail "OpenRC install should verify flux_agent status"
|
||||
)
|
||||
|
||||
test_remove_flux_agent_service_uses_openrc() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
local temp_root
|
||||
temp_root=$(mktemp -d)
|
||||
SERVICE_MANAGER="openrc"
|
||||
FLUX_AGENT_OPENRC_SERVICE_FILE="$temp_root/init.d/flux_agent"
|
||||
mkdir -p "$(dirname "$FLUX_AGENT_OPENRC_SERVICE_FILE")"
|
||||
: > "$FLUX_AGENT_OPENRC_SERVICE_FILE"
|
||||
|
||||
local rc_service_calls=""
|
||||
local rc_update_calls=""
|
||||
rc-service() {
|
||||
rc_service_calls+=$'\n'"$*"
|
||||
return 0
|
||||
}
|
||||
rc-update() {
|
||||
rc_update_calls+=$'\n'"$*"
|
||||
return 0
|
||||
}
|
||||
|
||||
stop_flux_agent_service
|
||||
disable_flux_agent_service
|
||||
remove_flux_agent_service
|
||||
|
||||
[[ "$rc_service_calls" == *"stop"* ]] || fail "OpenRC uninstall should stop flux_agent"
|
||||
[[ "$rc_update_calls" == *"del flux_agent default"* ]] || \
|
||||
fail "OpenRC uninstall should remove flux_agent from the default runlevel"
|
||||
[[ ! -e "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]] || fail "OpenRC uninstall should remove its service file"
|
||||
)
|
||||
|
||||
test_cleanup_legacy_gost_installation_removes_service_and_binary() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
@@ -504,6 +620,9 @@ test_update_flux_agent_skips_proxy_prompt_when_not_installed
|
||||
test_install_flux_agent_preserves_legacy_gost_when_download_fails
|
||||
test_update_flux_agent_preserves_legacy_gost_when_download_fails
|
||||
test_install_flux_agent_writes_json_safe_config
|
||||
test_install_script_bootstraps_bash_for_alpine
|
||||
test_install_flux_agent_uses_openrc
|
||||
test_remove_flux_agent_service_uses_openrc
|
||||
test_cleanup_legacy_gost_installation_removes_service_and_binary
|
||||
test_cleanup_legacy_gost_installation_preserves_unrelated_gost
|
||||
test_install_script_accepts_proxy_url_env_without_prompt
|
||||
@@ -514,4 +633,4 @@ test_panel_install_script_uses_default_proxy
|
||||
test_panel_install_script_accepts_proxy_url_env_without_prompt
|
||||
test_panel_install_script_defaults_proxy_on_eof
|
||||
|
||||
echo "install script proxy tests passed"
|
||||
echo "install script tests passed"
|
||||
|
||||
@@ -595,6 +595,20 @@ export interface TunnelQualityHopApiItem {
|
||||
targetPort?: number;
|
||||
}
|
||||
|
||||
export interface TunnelQualityCandidateHopApiItem
|
||||
extends TunnelQualityHopApiItem {
|
||||
fromRole: "entry" | "middle" | "exit";
|
||||
toRole: "middle" | "exit" | "target";
|
||||
hopIndex: number;
|
||||
selected: boolean;
|
||||
errorMessage?: string;
|
||||
}
|
||||
|
||||
export interface TunnelQualityChainDetailsApiItem {
|
||||
primaryPath?: TunnelQualityHopApiItem[];
|
||||
candidateHops?: TunnelQualityCandidateHopApiItem[];
|
||||
}
|
||||
|
||||
export interface TunnelQualityApiItem {
|
||||
tunnelId: number;
|
||||
entryToExitLatency: number;
|
||||
|
||||
@@ -28,12 +28,13 @@ function DialogClose({
|
||||
return <DialogPrimitive.Close data-slot="dialog-close" {...props} />;
|
||||
}
|
||||
|
||||
function DialogOverlay({
|
||||
className,
|
||||
...props
|
||||
}: React.ComponentProps<typeof DialogPrimitive.Overlay>) {
|
||||
const DialogOverlay = React.forwardRef<
|
||||
React.ElementRef<typeof DialogPrimitive.Overlay>,
|
||||
React.ComponentPropsWithoutRef<typeof DialogPrimitive.Overlay>
|
||||
>(({ className, ...props }, ref) => {
|
||||
return (
|
||||
<DialogPrimitive.Overlay
|
||||
ref={ref}
|
||||
className={cn(
|
||||
"fixed inset-0 z-50 bg-black/30 backdrop-blur-md data-[state=open]:animate-in data-[state=closed]:animate-out data-[state=closed]:fade-out-0 data-[state=open]:fade-in-0",
|
||||
className,
|
||||
@@ -42,7 +43,9 @@ function DialogOverlay({
|
||||
{...props}
|
||||
/>
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
DialogOverlay.displayName = DialogPrimitive.Overlay.displayName;
|
||||
|
||||
function DialogContent({
|
||||
className,
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
export const TUNNEL_QUALITY_INTERVAL_CONFIG_KEY =
|
||||
"monitor_tunnel_quality_interval_sec";
|
||||
export const DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC = 1;
|
||||
export const MIN_TUNNEL_QUALITY_INTERVAL_SEC = 1;
|
||||
export const MAX_TUNNEL_QUALITY_INTERVAL_SEC = 3600;
|
||||
|
||||
export const parseTunnelQualityIntervalSeconds = (value: unknown): number => {
|
||||
const seconds = Number(value);
|
||||
|
||||
return Number.isInteger(seconds) &&
|
||||
seconds >= MIN_TUNNEL_QUALITY_INTERVAL_SEC &&
|
||||
seconds <= MAX_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
? seconds
|
||||
: DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC;
|
||||
};
|
||||
|
||||
export const validateTunnelQualityInterval = (value: string): string | null => {
|
||||
const normalized = value.trim();
|
||||
|
||||
if (!normalized) {
|
||||
return "请输入探测间隔";
|
||||
}
|
||||
const seconds = Number(normalized);
|
||||
|
||||
if (!Number.isInteger(seconds)) {
|
||||
return "探测间隔必须是整数";
|
||||
}
|
||||
if (
|
||||
seconds < MIN_TUNNEL_QUALITY_INTERVAL_SEC ||
|
||||
seconds > MAX_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
) {
|
||||
return `探测间隔必须在 ${MIN_TUNNEL_QUALITY_INTERVAL_SEC} 到 ${MAX_TUNNEL_QUALITY_INTERVAL_SEC} 秒之间`;
|
||||
}
|
||||
|
||||
return null;
|
||||
};
|
||||
|
||||
export const tunnelQualityIntervalLabel = (seconds: number): string =>
|
||||
seconds === 1 ? "每秒" : `每 ${seconds} 秒`;
|
||||
@@ -43,6 +43,14 @@ import { BackIcon, SettingsIcon } from "@/components/icons";
|
||||
import { ThemeSettings } from "@/components/theme-settings";
|
||||
import { isAdmin } from "@/utils/auth";
|
||||
import { getCachedConfigs, configCache, updateSiteConfig } from "@/config/site";
|
||||
import {
|
||||
DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
MAX_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
MIN_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
parseTunnelQualityIntervalSeconds,
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
validateTunnelQualityInterval,
|
||||
} from "@/config/tunnel-quality";
|
||||
import {
|
||||
type UpdateReleaseChannel,
|
||||
getUpdateReleaseChannel,
|
||||
@@ -157,6 +165,16 @@ const CONFIG_ITEMS: ConfigItem[] = [
|
||||
"关闭后,前端停止自动刷新,后端停止实时隧道质量探测(全局配置)",
|
||||
type: "switch",
|
||||
},
|
||||
{
|
||||
key: TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
label: "隧道质量探测间隔",
|
||||
placeholder: String(DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC),
|
||||
description:
|
||||
"设置实时隧道质量检测的执行频率,单位为秒;允许 1–3600 秒,默认 1 秒。",
|
||||
type: "input",
|
||||
dependsOn: "monitor_tunnel_quality_enabled",
|
||||
dependsValue: "true",
|
||||
},
|
||||
{
|
||||
key: "monitor_retention_days",
|
||||
label: "监控数据保留天数",
|
||||
@@ -239,6 +257,7 @@ const getInitialConfigs = (): Record<string, string> => {
|
||||
"cloudflare_secret_key",
|
||||
"forward_compact_mode",
|
||||
"monitor_tunnel_quality_enabled",
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
"monitor_retention_days",
|
||||
"ip",
|
||||
"panel_domain",
|
||||
@@ -622,6 +641,19 @@ export default function ConfigPage() {
|
||||
|
||||
// 保存配置
|
||||
const handleSave = async () => {
|
||||
const intervalValue = configs[TUNNEL_QUALITY_INTERVAL_CONFIG_KEY];
|
||||
const intervalChanged =
|
||||
intervalValue !== originalConfigs[TUNNEL_QUALITY_INTERVAL_CONFIG_KEY];
|
||||
const intervalError = intervalChanged
|
||||
? validateTunnelQualityInterval(intervalValue || "")
|
||||
: null;
|
||||
|
||||
if (intervalError) {
|
||||
toast.error(intervalError);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
setSaving(true);
|
||||
try {
|
||||
const changedKeys = Object.keys(configs).filter(
|
||||
@@ -667,12 +699,22 @@ export default function ConfigPage() {
|
||||
}),
|
||||
);
|
||||
|
||||
// 如果隧道质量检测开关变更,通知 tunnel-monitor-view
|
||||
if (changedKeys.includes("monitor_tunnel_quality_enabled")) {
|
||||
// 如果隧道质量检测配置变更,通知 tunnel-monitor-view
|
||||
if (
|
||||
changedKeys.some((key) =>
|
||||
[
|
||||
"monitor_tunnel_quality_enabled",
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
].includes(key),
|
||||
)
|
||||
) {
|
||||
window.dispatchEvent(
|
||||
new CustomEvent("monitorTunnelQualityEnabledChanged", {
|
||||
detail: {
|
||||
enabled: configs["monitor_tunnel_quality_enabled"] === "true",
|
||||
intervalSec: parseTunnelQualityIntervalSeconds(
|
||||
configs[TUNNEL_QUALITY_INTERVAL_CONFIG_KEY],
|
||||
),
|
||||
},
|
||||
}),
|
||||
);
|
||||
@@ -1075,11 +1117,21 @@ export default function ConfigPage() {
|
||||
case "bg_image":
|
||||
return renderBgImageUploader();
|
||||
|
||||
case "input":
|
||||
case "input": {
|
||||
if (isBrandPreviewKey(item.key)) {
|
||||
return renderBrandAssetUploader(item.key, isChanged);
|
||||
}
|
||||
|
||||
const isTunnelQualityInterval =
|
||||
item.key === TUNNEL_QUALITY_INTERVAL_CONFIG_KEY;
|
||||
const intervalValue = isTunnelQualityInterval
|
||||
? (configs[item.key] ?? String(DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC))
|
||||
: (configs[item.key] ?? "");
|
||||
const intervalError =
|
||||
isTunnelQualityInterval && configs[item.key] !== undefined
|
||||
? validateTunnelQualityInterval(intervalValue)
|
||||
: null;
|
||||
|
||||
return (
|
||||
<Input
|
||||
classNames={{
|
||||
@@ -1091,14 +1143,30 @@ export default function ConfigPage() {
|
||||
description={
|
||||
isCommercialDisabled ? "需商业版授权才能修改此项" : undefined
|
||||
}
|
||||
endContent={isTunnelQualityInterval ? "秒" : undefined}
|
||||
errorMessage={intervalError || undefined}
|
||||
isDisabled={isCommercialDisabled}
|
||||
isInvalid={Boolean(intervalError)}
|
||||
max={
|
||||
isTunnelQualityInterval
|
||||
? MAX_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
: undefined
|
||||
}
|
||||
min={
|
||||
isTunnelQualityInterval
|
||||
? MIN_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
: undefined
|
||||
}
|
||||
placeholder={item.placeholder}
|
||||
size="md"
|
||||
value={configs[item.key] || ""}
|
||||
step={isTunnelQualityInterval ? 1 : undefined}
|
||||
type={isTunnelQualityInterval ? "number" : "text"}
|
||||
value={intervalValue}
|
||||
variant="bordered"
|
||||
onChange={(e) => handleConfigChange(item.key, e.target.value)}
|
||||
/>
|
||||
);
|
||||
}
|
||||
|
||||
case "switch":
|
||||
return (
|
||||
|
||||
@@ -2,6 +2,8 @@ import type {
|
||||
MonitorTunnelApiItem,
|
||||
TunnelMetricApiItem,
|
||||
TunnelQualityApiItem,
|
||||
TunnelQualityCandidateHopApiItem,
|
||||
TunnelQualityChainDetailsApiItem,
|
||||
TunnelQualityHopApiItem,
|
||||
} from "@/api/types";
|
||||
|
||||
@@ -53,12 +55,17 @@ import {
|
||||
TableRow,
|
||||
TableCell,
|
||||
} from "@/shadcn-bridge/heroui/table";
|
||||
import {
|
||||
DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
parseTunnelQualityIntervalSeconds,
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
tunnelQualityIntervalLabel,
|
||||
} from "@/config/tunnel-quality";
|
||||
|
||||
interface TunnelMonitorViewProps {
|
||||
viewMode?: "list" | "grid";
|
||||
}
|
||||
|
||||
const QUALITY_POLL_INTERVAL = 1_000; // 1 second
|
||||
const MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY =
|
||||
"monitor_tunnel_quality_enabled";
|
||||
const MONITOR_TUNNEL_QUALITY_ENABLED_EVENT =
|
||||
@@ -463,18 +470,249 @@ const TrafficChartCard = React.memo(function TrafficChartCard({
|
||||
);
|
||||
});
|
||||
|
||||
function ForwardingChainTopology({ hopsStr }: { hopsStr?: string }) {
|
||||
if (!hopsStr) return null;
|
||||
|
||||
let hops: TunnelQualityHopApiItem[] = [];
|
||||
|
||||
const parseTunnelQualityChainDetails = (
|
||||
raw?: string,
|
||||
): TunnelQualityChainDetailsApiItem => {
|
||||
if (!raw) return {};
|
||||
try {
|
||||
hops = JSON.parse(hopsStr);
|
||||
const parsed: unknown = JSON.parse(raw);
|
||||
|
||||
// Backward-compatible with historical rows that stored the primary path
|
||||
// directly as a JSON array.
|
||||
if (Array.isArray(parsed)) {
|
||||
return { primaryPath: parsed as TunnelQualityHopApiItem[] };
|
||||
}
|
||||
if (parsed && typeof parsed === "object") {
|
||||
return parsed as TunnelQualityChainDetailsApiItem;
|
||||
}
|
||||
} catch {
|
||||
return null;
|
||||
return {};
|
||||
}
|
||||
|
||||
if (!Array.isArray(hops) || hops.length === 0) return null;
|
||||
return {};
|
||||
};
|
||||
|
||||
type TunnelTopologyHop = TunnelQualityHopApiItem & {
|
||||
errorMessage?: string;
|
||||
};
|
||||
|
||||
interface TunnelTopologyPath {
|
||||
key: string;
|
||||
hops: TunnelTopologyHop[];
|
||||
alternativeNodeIndex?: number;
|
||||
}
|
||||
|
||||
const tunnelTopologyHopKey = (fromNodeId: number, toNodeId: number) =>
|
||||
`${fromNodeId}:${toNodeId}`;
|
||||
|
||||
const buildTunnelTopologyPaths = (
|
||||
details: TunnelQualityChainDetailsApiItem,
|
||||
): TunnelTopologyPath[] => {
|
||||
const primaryHops = details.primaryPath ?? [];
|
||||
const candidates = details.candidateHops ?? [];
|
||||
|
||||
if (primaryHops.length === 0) {
|
||||
const publicCandidates = candidates.filter(
|
||||
(candidate) => candidate.toRole === "target",
|
||||
);
|
||||
const selected = publicCandidates.find((candidate) => candidate.selected);
|
||||
const paths: TunnelTopologyPath[] = [];
|
||||
|
||||
if (selected) {
|
||||
paths.push({ key: "primary-public", hops: [selected] });
|
||||
}
|
||||
for (const candidate of publicCandidates) {
|
||||
if (candidate.selected) continue;
|
||||
paths.push({
|
||||
key: `alternative-public-${candidate.fromNodeId}`,
|
||||
hops: [candidate],
|
||||
alternativeNodeIndex: 0,
|
||||
});
|
||||
}
|
||||
|
||||
return paths;
|
||||
}
|
||||
|
||||
const primaryNodeIds = [
|
||||
primaryHops[0].fromNodeId,
|
||||
...primaryHops.map((hop) => hop.toNodeId),
|
||||
];
|
||||
const internalCandidates = candidates.filter(
|
||||
(candidate) => candidate.toRole !== "target",
|
||||
);
|
||||
const candidateHopMap = new Map<string, TunnelQualityCandidateHopApiItem>();
|
||||
const alternativeNodes = new Map<
|
||||
string,
|
||||
{ column: number; nodeId: number }
|
||||
>();
|
||||
|
||||
for (const candidate of internalCandidates) {
|
||||
candidateHopMap.set(
|
||||
tunnelTopologyHopKey(candidate.fromNodeId, candidate.toNodeId),
|
||||
candidate,
|
||||
);
|
||||
|
||||
const sourceColumn = candidate.hopIndex;
|
||||
const targetColumn = candidate.hopIndex + 1;
|
||||
|
||||
if (
|
||||
sourceColumn >= 0 &&
|
||||
sourceColumn < primaryNodeIds.length &&
|
||||
candidate.fromNodeId !== primaryNodeIds[sourceColumn]
|
||||
) {
|
||||
alternativeNodes.set(`${sourceColumn}:${candidate.fromNodeId}`, {
|
||||
column: sourceColumn,
|
||||
nodeId: candidate.fromNodeId,
|
||||
});
|
||||
}
|
||||
if (
|
||||
targetColumn >= 0 &&
|
||||
targetColumn < primaryNodeIds.length &&
|
||||
candidate.toNodeId !== primaryNodeIds[targetColumn]
|
||||
) {
|
||||
alternativeNodes.set(`${targetColumn}:${candidate.toNodeId}`, {
|
||||
column: targetColumn,
|
||||
nodeId: candidate.toNodeId,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
const paths: TunnelTopologyPath[] = [{ key: "primary", hops: primaryHops }];
|
||||
|
||||
for (const alternative of alternativeNodes.values()) {
|
||||
const nodeIds = [...primaryNodeIds];
|
||||
|
||||
nodeIds[alternative.column] = alternative.nodeId;
|
||||
const hops: TunnelTopologyHop[] = [];
|
||||
|
||||
for (let index = 0; index < nodeIds.length - 1; index += 1) {
|
||||
const fromNodeId = nodeIds[index];
|
||||
const toNodeId = nodeIds[index + 1];
|
||||
const usesPrimaryEdge =
|
||||
fromNodeId === primaryNodeIds[index] &&
|
||||
toNodeId === primaryNodeIds[index + 1];
|
||||
const hop = usesPrimaryEdge
|
||||
? primaryHops[index]
|
||||
: candidateHopMap.get(tunnelTopologyHopKey(fromNodeId, toNodeId));
|
||||
|
||||
if (!hop) break;
|
||||
hops.push(hop);
|
||||
}
|
||||
|
||||
if (hops.length === primaryHops.length) {
|
||||
paths.push({
|
||||
key: `alternative-${alternative.column}-${alternative.nodeId}`,
|
||||
hops,
|
||||
alternativeNodeIndex: alternative.column,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
return paths;
|
||||
};
|
||||
|
||||
function TunnelTopologyPathRow({ path }: { path: TunnelTopologyPath }) {
|
||||
return (
|
||||
<div className="flex min-w-max items-center py-2">
|
||||
{path.hops.map((hop, index) => {
|
||||
const hasError =
|
||||
Boolean(hop.errorMessage) || hop.latency < 0 || hop.loss > 0;
|
||||
const colorClass =
|
||||
hop.latency < 0 || hop.errorMessage
|
||||
? "text-danger"
|
||||
: hop.loss > 0
|
||||
? "text-warning"
|
||||
: "text-success";
|
||||
const borderColor = hasError ? "border-danger" : "";
|
||||
|
||||
return (
|
||||
<React.Fragment
|
||||
key={`${path.key}-${hop.fromNodeId}-${hop.toNodeId}-${index}`}
|
||||
>
|
||||
{index === 0 ? (
|
||||
<TopologyNodeChip
|
||||
isAlternative={path.alternativeNodeIndex === 0}
|
||||
name={hop.fromNodeName}
|
||||
/>
|
||||
) : null}
|
||||
<div
|
||||
className="relative mx-1 flex min-w-[70px] shrink-0 flex-col items-center justify-center"
|
||||
title={hop.errorMessage}
|
||||
>
|
||||
<span
|
||||
className={`mb-1 text-[10px] font-mono leading-none ${colorClass}`}
|
||||
>
|
||||
{hop.latency >= 0 && !hop.errorMessage
|
||||
? `${hop.latency.toFixed(0)}ms`
|
||||
: "超时"}
|
||||
</span>
|
||||
<div
|
||||
className={`relative flex h-[2px] w-full items-center justify-end bg-default-200 ${hop.latency < 0 || hop.errorMessage ? "!bg-danger" : ""}`}
|
||||
>
|
||||
<ArrowRight
|
||||
className={`absolute -right-2 z-10 h-3.5 w-3.5 rounded-full bg-background p-[1px] ${colorClass}`}
|
||||
/>
|
||||
</div>
|
||||
<span
|
||||
className={`mt-1.5 text-[10px] font-mono leading-none ${hop.errorMessage || hop.loss > 0 ? "text-warning" : "text-default-400"}`}
|
||||
>
|
||||
{hop.errorMessage ? "探测失败" : `${hop.loss.toFixed(0)}% 丢包`}
|
||||
</span>
|
||||
</div>
|
||||
<TopologyNodeChip
|
||||
borderColor={borderColor}
|
||||
isAlternative={path.alternativeNodeIndex === index + 1}
|
||||
name={hop.toNodeName}
|
||||
/>
|
||||
</React.Fragment>
|
||||
);
|
||||
})}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
function TopologyNodeChip({
|
||||
name,
|
||||
isAlternative = false,
|
||||
borderColor = "",
|
||||
}: {
|
||||
name: string;
|
||||
isAlternative?: boolean;
|
||||
borderColor?: string;
|
||||
}) {
|
||||
return (
|
||||
<Chip
|
||||
className={`shrink-0 font-mono shadow-sm ${borderColor}`}
|
||||
size="sm"
|
||||
variant="flat"
|
||||
>
|
||||
<span className="flex items-center gap-1.5">
|
||||
<span>{name}</span>
|
||||
{isAlternative ? (
|
||||
<span className="rounded-full bg-warning/20 px-1.5 py-0.5 text-[9px] font-semibold leading-none text-warning">
|
||||
备选
|
||||
</span>
|
||||
) : null}
|
||||
</span>
|
||||
</Chip>
|
||||
);
|
||||
}
|
||||
|
||||
const ForwardingChainTopology = React.memo(function ForwardingChainTopology({
|
||||
hopsStr,
|
||||
}: {
|
||||
hopsStr?: string;
|
||||
}) {
|
||||
const details = useMemo(
|
||||
() => parseTunnelQualityChainDetails(hopsStr),
|
||||
[hopsStr],
|
||||
);
|
||||
const topologyPaths = useMemo(
|
||||
() => buildTunnelTopologyPaths(details),
|
||||
[details],
|
||||
);
|
||||
|
||||
if (topologyPaths.length === 0) return null;
|
||||
|
||||
return (
|
||||
<Card className="border border-divider/60 shadow-sm transition-shadow bg-gradient-to-br from-background to-default-50/50 mt-4">
|
||||
@@ -485,62 +723,24 @@ function ForwardingChainTopology({ hopsStr }: { hopsStr?: string }) {
|
||||
</h3>
|
||||
</CardHeader>
|
||||
<CardBody className="py-2 px-4 pb-4">
|
||||
<div className="flex items-center overflow-x-auto pb-2 py-2">
|
||||
{hops.map((hop, index) => {
|
||||
const hasError = hop.latency < 0 || hop.loss > 0;
|
||||
const colorClass =
|
||||
hop.latency < 0
|
||||
? "text-danger"
|
||||
: hop.loss > 0
|
||||
? "text-warning"
|
||||
: "text-success";
|
||||
const borderColor = hasError ? "border-danger" : "";
|
||||
|
||||
return (
|
||||
<React.Fragment key={index}>
|
||||
{index === 0 && (
|
||||
<Chip
|
||||
className="shrink-0 font-mono shadow-sm"
|
||||
size="sm"
|
||||
variant="flat"
|
||||
>
|
||||
{hop.fromNodeName}
|
||||
</Chip>
|
||||
)}
|
||||
<div className="flex flex-col items-center justify-center min-w-[70px] mx-1 shrink-0 relative">
|
||||
<span
|
||||
className={`text-[10px] font-mono leading-none mb-1 ${colorClass}`}
|
||||
>
|
||||
{hop.latency >= 0 ? `${hop.latency.toFixed(0)}ms` : "超时"}
|
||||
</span>
|
||||
<div
|
||||
className={`h-[2px] w-full relative flex items-center justify-end bg-default-200 ${hop.latency < 0 ? "!bg-danger" : ""}`}
|
||||
>
|
||||
<ArrowRight
|
||||
className={`w-3.5 h-3.5 absolute -right-2 ${colorClass} bg-background rounded-full p-[1px] z-10`}
|
||||
/>
|
||||
</div>
|
||||
<span
|
||||
className={`text-[10px] font-mono leading-none mt-1.5 ${hop.loss > 0 ? "text-warning" : "text-default-400"}`}
|
||||
>
|
||||
{hop.loss.toFixed(0)}% 丢包
|
||||
</span>
|
||||
</div>
|
||||
<Chip
|
||||
className={`shrink-0 font-mono shadow-sm ${borderColor}`}
|
||||
size="sm"
|
||||
variant="flat"
|
||||
>
|
||||
{hop.toNodeName}
|
||||
</Chip>
|
||||
</React.Fragment>
|
||||
);
|
||||
})}
|
||||
<div className="max-h-80 space-y-1 overflow-auto pb-1">
|
||||
{topologyPaths.map((path, index) => (
|
||||
<div
|
||||
key={path.key}
|
||||
className={
|
||||
index === 0
|
||||
? "overflow-x-auto"
|
||||
: "overflow-x-auto border-t border-dashed border-divider/60"
|
||||
}
|
||||
>
|
||||
<TunnelTopologyPathRow path={path} />
|
||||
</div>
|
||||
))}
|
||||
</div>
|
||||
</CardBody>
|
||||
</Card>
|
||||
);
|
||||
}
|
||||
});
|
||||
|
||||
export function TunnelMonitorView({
|
||||
viewMode = "grid",
|
||||
@@ -562,6 +762,9 @@ export function TunnelMonitorView({
|
||||
const qualityTimerRef = useRef<number | null>(null);
|
||||
const [monitorTunnelQualityEnabled, setMonitorTunnelQualityEnabled] =
|
||||
useState(true);
|
||||
const [tunnelQualityIntervalSec, setTunnelQualityIntervalSec] = useState(
|
||||
DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
);
|
||||
|
||||
// Detail view state
|
||||
const [detailTunnelId, setDetailTunnelId] = useState<number | null>(null);
|
||||
@@ -618,26 +821,28 @@ export function TunnelMonitorView({
|
||||
}
|
||||
}, []);
|
||||
|
||||
const loadMonitorTunnelQualityEnabled = useCallback(async () => {
|
||||
try {
|
||||
const response = await getConfigByName(
|
||||
MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY,
|
||||
);
|
||||
const loadTunnelQualityConfig = useCallback(async () => {
|
||||
const [enabledResponse, intervalResponse] = await Promise.all([
|
||||
getConfigByName(MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY).catch(
|
||||
() => null,
|
||||
),
|
||||
getConfigByName(TUNNEL_QUALITY_INTERVAL_CONFIG_KEY).catch(() => null),
|
||||
]);
|
||||
|
||||
setMonitorTunnelQualityEnabled(
|
||||
typeof response.data?.value === "string"
|
||||
? response.data.value === "true"
|
||||
: true,
|
||||
);
|
||||
} catch {
|
||||
setMonitorTunnelQualityEnabled(true);
|
||||
}
|
||||
setMonitorTunnelQualityEnabled(
|
||||
typeof enabledResponse?.data?.value === "string"
|
||||
? enabledResponse.data.value === "true"
|
||||
: true,
|
||||
);
|
||||
setTunnelQualityIntervalSec(
|
||||
parseTunnelQualityIntervalSeconds(intervalResponse?.data?.value),
|
||||
);
|
||||
}, []);
|
||||
|
||||
useEffect(() => {
|
||||
void loadTunnels();
|
||||
void loadMonitorTunnelQualityEnabled();
|
||||
}, [loadMonitorTunnelQualityEnabled, loadTunnels]);
|
||||
void loadTunnelQualityConfig();
|
||||
}, [loadTunnelQualityConfig, loadTunnels]);
|
||||
|
||||
useEffect(() => {
|
||||
const timer = window.setInterval(() => {
|
||||
@@ -649,8 +854,11 @@ export function TunnelMonitorView({
|
||||
|
||||
useEffect(() => {
|
||||
const handleMonitorTunnelQualityEnabledChanged = (event: Event) => {
|
||||
const enabled = (event as CustomEvent<{ enabled?: boolean }>).detail
|
||||
?.enabled;
|
||||
const detail = (
|
||||
event as CustomEvent<{ enabled?: boolean; intervalSec?: number }>
|
||||
).detail;
|
||||
const enabled = detail?.enabled;
|
||||
const intervalSec = detail?.intervalSec;
|
||||
|
||||
if (typeof enabled === "boolean") {
|
||||
setMonitorTunnelQualityEnabled(enabled);
|
||||
@@ -658,7 +866,12 @@ export function TunnelMonitorView({
|
||||
setQualityLoading(false);
|
||||
}
|
||||
} else {
|
||||
void loadMonitorTunnelQualityEnabled();
|
||||
void loadTunnelQualityConfig();
|
||||
}
|
||||
if (typeof intervalSec === "number") {
|
||||
setTunnelQualityIntervalSec(
|
||||
parseTunnelQualityIntervalSeconds(String(intervalSec)),
|
||||
);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -673,7 +886,7 @@ export function TunnelMonitorView({
|
||||
handleMonitorTunnelQualityEnabledChanged as EventListener,
|
||||
);
|
||||
};
|
||||
}, [loadMonitorTunnelQualityEnabled]);
|
||||
}, [loadTunnelQualityConfig]);
|
||||
|
||||
useEffect(() => {
|
||||
if (tunnels.length > 0 && !initialHistoryFetched.current) {
|
||||
@@ -726,7 +939,7 @@ export function TunnelMonitorView({
|
||||
}
|
||||
}, [tunnels]);
|
||||
|
||||
// --- Load quality snapshots (auto-polling every 10s) ---
|
||||
// --- Load quality snapshots using the configured probe interval ---
|
||||
const loadQuality = useCallback(async (options?: { silent?: boolean }) => {
|
||||
const silent = options?.silent ?? false;
|
||||
|
||||
@@ -788,7 +1001,7 @@ export function TunnelMonitorView({
|
||||
|
||||
qualityTimerRef.current = window.setInterval(() => {
|
||||
void loadQuality({ silent: true });
|
||||
}, QUALITY_POLL_INTERVAL);
|
||||
}, tunnelQualityIntervalSec * 1000);
|
||||
|
||||
return () => {
|
||||
if (qualityTimerRef.current) {
|
||||
@@ -796,7 +1009,7 @@ export function TunnelMonitorView({
|
||||
qualityTimerRef.current = null;
|
||||
}
|
||||
};
|
||||
}, [loadQuality, monitorTunnelQualityEnabled]);
|
||||
}, [loadQuality, monitorTunnelQualityEnabled, tunnelQualityIntervalSec]);
|
||||
|
||||
// --- Load quality history for detail chart ---
|
||||
const loadQualityHistory = useCallback(
|
||||
@@ -1065,7 +1278,7 @@ export function TunnelMonitorView({
|
||||
{monitorTunnelQualityEnabled ? (
|
||||
<>
|
||||
<LiveDot />
|
||||
<span>自动探测中(每秒测试,30秒上报)</span>
|
||||
<span>{`自动探测中(${tunnelQualityIntervalLabel(tunnelQualityIntervalSec)}测试)`}</span>
|
||||
</>
|
||||
) : (
|
||||
<>
|
||||
@@ -1129,7 +1342,7 @@ export function TunnelMonitorView({
|
||||
{monitorTunnelQualityEnabled ? (
|
||||
<>
|
||||
<LiveDot />
|
||||
<span>每秒探测 · 更新于 {lastQualityUpdate}</span>
|
||||
<span>{`${tunnelQualityIntervalLabel(tunnelQualityIntervalSec)}探测 · 更新于 ${lastQualityUpdate}`}</span>
|
||||
</>
|
||||
) : (
|
||||
<>
|
||||
|
||||
@@ -49,6 +49,41 @@ function useModalContext() {
|
||||
return React.useContext(ModalContext);
|
||||
}
|
||||
|
||||
interface ScrollPosition {
|
||||
element: HTMLElement | null;
|
||||
left: number;
|
||||
top: number;
|
||||
}
|
||||
|
||||
function captureScrollPositions(): ScrollPosition[] {
|
||||
const positions: ScrollPosition[] = [
|
||||
{ element: null, left: window.scrollX, top: window.scrollY },
|
||||
];
|
||||
|
||||
for (const element of Array.from(
|
||||
document.querySelectorAll<HTMLElement>("main, [data-scroll-container]"),
|
||||
)) {
|
||||
positions.push({
|
||||
element,
|
||||
left: element.scrollLeft,
|
||||
top: element.scrollTop,
|
||||
});
|
||||
}
|
||||
|
||||
return positions;
|
||||
}
|
||||
|
||||
function restoreScrollPositions(positions: ScrollPosition[]) {
|
||||
for (const position of positions) {
|
||||
if (position.element) {
|
||||
position.element.scrollLeft = position.left;
|
||||
position.element.scrollTop = position.top;
|
||||
} else {
|
||||
window.scrollTo(position.left, position.top);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type ModalSize = "sm" | "md" | "lg" | "xl" | "2xl" | "4xl" | "full";
|
||||
|
||||
function mapSize(size: ModalSize | undefined) {
|
||||
@@ -97,6 +132,46 @@ export function Modal({
|
||||
scrollBehavior,
|
||||
size,
|
||||
}: ModalProps) {
|
||||
const previousScrollPositionsRef = React.useRef<ScrollPosition[] | null>(
|
||||
null,
|
||||
);
|
||||
|
||||
// Radix focus management and scroll locking can move an ancestor scroll
|
||||
// container when a modal is opened from a card/grid item. Capture the
|
||||
// current positions before the open render and restore them after focus
|
||||
// settles so opening a modal never changes the page position.
|
||||
React.useLayoutEffect(() => {
|
||||
return () => {
|
||||
if (!isOpen) {
|
||||
previousScrollPositionsRef.current = captureScrollPositions();
|
||||
}
|
||||
};
|
||||
}, [isOpen]);
|
||||
|
||||
React.useLayoutEffect(() => {
|
||||
const positions = previousScrollPositionsRef.current;
|
||||
|
||||
if (!isOpen || !positions) {
|
||||
return;
|
||||
}
|
||||
|
||||
restoreScrollPositions(positions);
|
||||
let nestedFrame = 0;
|
||||
const frame = window.requestAnimationFrame(() => {
|
||||
restoreScrollPositions(positions);
|
||||
nestedFrame = window.requestAnimationFrame(() =>
|
||||
restoreScrollPositions(positions),
|
||||
);
|
||||
});
|
||||
|
||||
previousScrollPositionsRef.current = null;
|
||||
|
||||
return () => {
|
||||
window.cancelAnimationFrame(frame);
|
||||
window.cancelAnimationFrame(nestedFrame);
|
||||
};
|
||||
}, [isOpen]);
|
||||
|
||||
const handleOpenChange = (open: boolean) => {
|
||||
onOpenChange?.(open);
|
||||
if (!open) {
|
||||
|
||||
Reference in New Issue
Block a user