mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 15:46:38 +08:00
Compare commits
33 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 91d79b6b3a | |||
| fc7df6bd64 | |||
| 25dfb84324 | |||
| 4ebd6703fe | |||
| 1f53a39784 | |||
| 5ebd4c2a91 | |||
| 6c93d829c6 | |||
| 5d22d4cb06 | |||
| e5cd5af550 | |||
| cdcdfd8ff0 | |||
| 791773fd62 | |||
| 13764b4615 | |||
| 4c882d907b | |||
| 6033e39466 | |||
| 0f3242bf11 | |||
| d97d91801d | |||
| 727ef56c67 | |||
| a40150b136 | |||
| a923ec4785 | |||
| 42c5492c1d | |||
| 869d726b7a | |||
| 55a931510b | |||
| a259dd83b2 | |||
| cc0b8de2e1 | |||
| 58ef260755 | |||
| 521fe79b15 | |||
| 90012725cc | |||
| 615d9e67eb | |||
| a131b70613 | |||
| cbed4eab23 | |||
| c2745dcd56 | |||
| ad4109594a | |||
| efc8c75dcb |
@@ -1,6 +1,6 @@
|
||||
# FLVX
|
||||
|
||||
> **联系我们**: [Telegram群组](https://t.me/flvxpanel)
|
||||
> **联系我们**: [Telegram群组](https://t.me/flvxchannel)
|
||||
|
||||
|
||||
## 特性
|
||||
|
||||
@@ -15,10 +15,15 @@ services:
|
||||
JWT_SECRET: ${JWT_SECRET}
|
||||
SERVER_ADDR: :6365
|
||||
TZ: Asia/Shanghai
|
||||
FLUX_VERSION: ${FLUX_VERSION:-dev}
|
||||
PANEL_DEPLOY_DIR: /opt/flvx-panel
|
||||
PANEL_BACKEND_CONTAINER: flux-panel-backend
|
||||
ports:
|
||||
- "${BACKEND_PORT}:6365"
|
||||
volumes:
|
||||
- sqlite_data:/app/data
|
||||
- /var/run/docker.sock:/var/run/docker.sock
|
||||
- ./:/opt/flvx-panel
|
||||
networks:
|
||||
- gost-network
|
||||
stop_grace_period: 30s
|
||||
|
||||
@@ -15,10 +15,15 @@ services:
|
||||
JWT_SECRET: ${JWT_SECRET}
|
||||
SERVER_ADDR: :6365
|
||||
TZ: Asia/Shanghai
|
||||
FLUX_VERSION: ${FLUX_VERSION:-dev}
|
||||
PANEL_DEPLOY_DIR: /opt/flvx-panel
|
||||
PANEL_BACKEND_CONTAINER: flux-panel-backend
|
||||
ports:
|
||||
- "${BACKEND_PORT}:6365"
|
||||
volumes:
|
||||
- sqlite_data:/app/data
|
||||
- /var/run/docker.sock:/var/run/docker.sock
|
||||
- ./:/opt/flvx-panel
|
||||
networks:
|
||||
- gost-network
|
||||
stop_grace_period: 30s
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,180 @@
|
||||
# Custom Best-Exit Probe Target Design
|
||||
|
||||
Date: 2026-05-01
|
||||
Status: Approved design
|
||||
|
||||
## Goal
|
||||
|
||||
Allow each tunnel to define the TCP target used for exit-side quality probing instead of always probing `www.bing.com:443`.
|
||||
|
||||
The custom target must be used consistently by:
|
||||
|
||||
- `best` exit scoring: each exit probes the configured target to measure exit-to-public quality.
|
||||
- Tunnel quality monitoring: the existing exit-side quality check probes the same configured target.
|
||||
|
||||
If a tunnel does not configure a target, behavior remains compatible with today: `www.bing.com:443`.
|
||||
|
||||
## Non-Goals
|
||||
|
||||
- Do not add HTTP/HTTPS request probing in this phase. The probe remains TCP host/port measurement.
|
||||
- Do not add a global default target setting in this phase.
|
||||
- Do not require existing tunnels to be edited or migrated manually.
|
||||
- Do not change the `best` switching thresholds, confirmation rounds, cooldowns, or runtime chain ordering semantics.
|
||||
- Do not add frontend test infrastructure.
|
||||
|
||||
## User-Facing Behavior
|
||||
|
||||
Each tunnel form gets a compact quality target section:
|
||||
|
||||
- Host input, placeholder `www.bing.com`.
|
||||
- Port input, placeholder `443`.
|
||||
- Helper text: this target is used for tunnel quality detection and `best` optimal-exit scoring; leaving it empty uses `www.bing.com:443`.
|
||||
|
||||
Tunnel list/get responses include the configured target so edit forms can round-trip it. The UI displays the effective target near quality/best-exit information as `测试目标:host:port`.
|
||||
|
||||
## Data Model
|
||||
|
||||
Add nullable/default-compatible fields to `model.Tunnel`:
|
||||
|
||||
- `ProbeTargetHost string` mapped to `probe_target_host`, `type:text`, default `''`.
|
||||
- `ProbeTargetPort int` mapped to `probe_target_port`, default `0`.
|
||||
|
||||
Effective target resolution:
|
||||
|
||||
- If `ProbeTargetHost` is non-empty and `ProbeTargetPort` is valid, use it.
|
||||
- Otherwise use `www.bing.com:443`.
|
||||
|
||||
The existing `TunnelQuality` persisted fields `exit_to_bing_latency` and `exit_to_bing_loss` remain unchanged for compatibility. They will semantically mean exit-to-configured-test-target after this change. API/UI labels should avoid saying `Bing` for new displays.
|
||||
|
||||
## Validation
|
||||
|
||||
On create/update:
|
||||
|
||||
- Empty host and empty/zero port are allowed and mean default target.
|
||||
- If either host or port is set, validate both as a pair.
|
||||
- Host is trimmed and must not contain URL scheme, path, query, or whitespace.
|
||||
- Host can be a domain, IPv4, or IPv6 literal. Bracketed IPv6 input should be normalized by removing surrounding brackets.
|
||||
- Port must be an integer from `1` to `65535`.
|
||||
- Do not perform network probing during save; external network failures must not block configuration changes.
|
||||
|
||||
Errors should be specific, for example:
|
||||
|
||||
- `测试目标 Host 不能为空`
|
||||
- `测试目标端口必须是 1-65535`
|
||||
- `测试目标 Host 不能包含协议或路径`
|
||||
|
||||
## Backend Flow
|
||||
|
||||
Introduce a small value/helper near the tunnel quality and best-exit code:
|
||||
|
||||
```go
|
||||
type tunnelProbeTarget struct {
|
||||
Host string
|
||||
Port int
|
||||
}
|
||||
```
|
||||
|
||||
Helpers:
|
||||
|
||||
- `defaultTunnelProbeTarget() tunnelProbeTarget` returns `www.bing.com:443`.
|
||||
- `normalizeTunnelProbeTarget(host string, port int) (tunnelProbeTarget, bool, error)` validates user input; the boolean indicates whether the user explicitly configured a target.
|
||||
- `effectiveTunnelProbeTarget(tunnel *model.Tunnel) tunnelProbeTarget` returns configured target or default.
|
||||
|
||||
Use the effective target in `tunnelQualityProber.probeTunnel`:
|
||||
|
||||
- Type 1 and unknown tunnel fallback probes entry node to effective target instead of hardcoded Bing.
|
||||
- Type 2 probes the selected/current exit node to effective target instead of hardcoded Bing.
|
||||
- `probeBestExitOwners` receives the effective target and passes it into best-exit owner scoring.
|
||||
|
||||
Use the effective target in `evaluateBestExitOwner`:
|
||||
|
||||
- Owner-to-exit measurement stays unchanged.
|
||||
- Exit-to-public measurement probes `target.Host:target.Port` instead of `bestExitPublicTargetHost:bestExitPublicTargetPort`.
|
||||
- The per-round public probe cache key must include node ID plus target host and port so future extensions cannot reuse measurements across different targets.
|
||||
|
||||
## API Shape
|
||||
|
||||
Tunnel list/get data includes:
|
||||
|
||||
```json
|
||||
{
|
||||
"probeTargetHost": "example.com",
|
||||
"probeTargetPort": 443
|
||||
}
|
||||
```
|
||||
|
||||
For old/default tunnels, return empty host and `0` to represent `use default`. The edit form must preserve default-as-empty unless the user explicitly saves a custom target.
|
||||
|
||||
Quality monitoring response includes effective target display metadata:
|
||||
|
||||
```json
|
||||
{
|
||||
"probeTargetHost": "www.bing.com",
|
||||
"probeTargetPort": 443
|
||||
}
|
||||
```
|
||||
|
||||
Existing `exitToBingLatency` and `exitToBingLoss` keys stay to avoid breaking frontend and external consumers.
|
||||
|
||||
## Frontend Flow
|
||||
|
||||
Extend `ChainTunnel` only if needed for node-level data; the target belongs to the tunnel, so `Tunnel` and `TunnelForm` get:
|
||||
|
||||
- `probeTargetHost?: string`
|
||||
- `probeTargetPort?: number`
|
||||
|
||||
On edit:
|
||||
|
||||
- Populate form fields from tunnel response.
|
||||
- Empty or zero means default target.
|
||||
|
||||
On submit:
|
||||
|
||||
- Trim host.
|
||||
- Convert blank port to `0`.
|
||||
- Send `probeTargetHost` and `probeTargetPort` with create/update payload.
|
||||
|
||||
Display:
|
||||
|
||||
- In the form helper, show default target behavior.
|
||||
- In quality/best-exit display areas, avoid `Bing` wording; prefer `测试目标` or the concrete `host:port`.
|
||||
|
||||
## Error Handling
|
||||
|
||||
- Invalid target input returns a normal API error envelope with a specific message.
|
||||
- Probe failures use existing quality error paths and best-exit scoring failure entries.
|
||||
- If all exit-to-target probes fail, best-exit behavior remains the same as today when all Bing probes fail: no valid best decision is applied from that round.
|
||||
|
||||
## Testing
|
||||
|
||||
Backend tests:
|
||||
|
||||
- Normalize default target when host/port are empty.
|
||||
- Reject partial host/port configuration and invalid port ranges.
|
||||
- Reject host values with URL scheme/path/whitespace.
|
||||
- Create/update tunnel persists `probeTargetHost` and `probeTargetPort`.
|
||||
- `ListTunnels` returns target fields.
|
||||
- `tunnelQualityProber` uses configured target instead of `www.bing.com:443`.
|
||||
- `best` scoring uses configured target for exit-to-target probes.
|
||||
- Empty target preserves old default `www.bing.com:443` behavior.
|
||||
|
||||
Frontend verification:
|
||||
|
||||
- `pnpm run build` passes.
|
||||
- Manual UI check: create/edit tunnel with blank target and custom target, confirm payload and round-trip display.
|
||||
|
||||
## Rollout And Compatibility
|
||||
|
||||
- Existing tunnels continue using `www.bing.com:443` because empty target resolves to default.
|
||||
- SQLite/PostgreSQL schema changes are handled by existing auto-migration.
|
||||
- Historical `TunnelQuality` rows keep existing columns and are not rewritten.
|
||||
- No runtime agent change is required; the panel already performs these quality probes through existing node ping APIs.
|
||||
|
||||
## Open Decisions
|
||||
|
||||
None. User-approved decisions:
|
||||
|
||||
- Per-tunnel fields are `host + port`.
|
||||
- The target applies to both `best` scoring and tunnel quality monitoring.
|
||||
- Probe type remains TCP host/port.
|
||||
- Empty target defaults to `www.bing.com:443`.
|
||||
@@ -9,10 +9,14 @@ ARG TARGETOS
|
||||
ARG TARGETARCH
|
||||
RUN CGO_ENABLED=0 GOOS=${TARGETOS:-linux} env ${TARGETARCH:+GOARCH=${TARGETARCH}} go build -o /out/paneld ./cmd/paneld
|
||||
|
||||
FROM docker:27-cli AS dockercli
|
||||
|
||||
FROM debian:bookworm-slim
|
||||
WORKDIR /app
|
||||
RUN apt-get update && apt-get install -y --no-install-recommends ca-certificates wget && rm -rf /var/lib/apt/lists/*
|
||||
COPY --from=builder /out/paneld /app/paneld
|
||||
COPY --from=dockercli /usr/local/bin/docker /usr/local/bin/docker
|
||||
COPY --from=dockercli /usr/local/libexec/docker/cli-plugins/docker-compose /usr/local/libexec/docker/cli-plugins/docker-compose
|
||||
|
||||
ENV SERVER_ADDR=:6365
|
||||
EXPOSE 6365
|
||||
|
||||
@@ -978,6 +978,7 @@ func (h *Handler) prepareTunnelDiagnosis(tunnelID int64) (string, string, []diag
|
||||
|
||||
ipPreference := h.repo.GetTunnelIPPreference(tunnelID)
|
||||
protocol := strings.ToLower(strings.TrimSpace(tunnel.Protocol))
|
||||
probeTarget := effectiveTunnelProbeTargetValues(tunnel.ProbeTargetHost, tunnel.ProbeTargetPort)
|
||||
inNodes, chainHops, outNodes := splitChainNodeGroups(chainRows)
|
||||
workItems := make([]diagnosisWorkItem, 0, len(chainRows)*2)
|
||||
|
||||
@@ -987,8 +988,8 @@ func (h *Handler) prepareTunnelDiagnosis(tunnelID int64) (string, string, []diag
|
||||
description := fmt.Sprintf("入口(%s)->外网", inNode.NodeName)
|
||||
workItems = append(workItems, diagnosisWorkItem{
|
||||
fromNodeID: inNode.NodeID,
|
||||
targetIP: "www.bing.com",
|
||||
targetPort: 443,
|
||||
targetIP: probeTarget.Host,
|
||||
targetPort: probeTarget.Port,
|
||||
description: description,
|
||||
protocol: "tcp",
|
||||
metadata: map[string]interface{}{
|
||||
@@ -1079,8 +1080,8 @@ func (h *Handler) prepareTunnelDiagnosis(tunnelID int64) (string, string, []diag
|
||||
description := fmt.Sprintf("出口(%s)->外网", outNode.NodeName)
|
||||
workItems = append(workItems, diagnosisWorkItem{
|
||||
fromNodeID: outNode.NodeID,
|
||||
targetIP: "www.bing.com",
|
||||
targetPort: 443,
|
||||
targetIP: probeTarget.Host,
|
||||
targetPort: probeTarget.Port,
|
||||
description: description,
|
||||
protocol: "tcp",
|
||||
metadata: map[string]interface{}{
|
||||
@@ -1093,8 +1094,8 @@ func (h *Handler) prepareTunnelDiagnosis(tunnelID int64) (string, string, []diag
|
||||
description := fmt.Sprintf("入口(%s)->外网", inNode.NodeName)
|
||||
workItems = append(workItems, diagnosisWorkItem{
|
||||
fromNodeID: inNode.NodeID,
|
||||
targetIP: "www.bing.com",
|
||||
targetPort: 443,
|
||||
targetIP: probeTarget.Host,
|
||||
targetPort: probeTarget.Port,
|
||||
description: description,
|
||||
protocol: "tcp",
|
||||
metadata: map[string]interface{}{
|
||||
|
||||
@@ -46,6 +46,7 @@ type Handler struct {
|
||||
jobsWG sync.WaitGroup
|
||||
|
||||
upgradeMu sync.Mutex
|
||||
systemUpgradeMu sync.Mutex
|
||||
pendingUpgradeRedeploy map[int64]struct{}
|
||||
nodeOnlineRedeployAt map[int64]time.Time
|
||||
nodeOnlineRedeployQueued map[int64]struct{}
|
||||
@@ -156,6 +157,9 @@ func (h *Handler) Register(mux *http.ServeMux) {
|
||||
mux.HandleFunc("/api/v1/config/update", h.updateConfigs)
|
||||
mux.HandleFunc("/api/v1/config/update-single", h.updateSingleConfig)
|
||||
mux.HandleFunc("/api/v1/system/storage", h.storageSummary)
|
||||
mux.HandleFunc("/api/v1/system/version", h.systemVersion)
|
||||
mux.HandleFunc("/api/v1/system/check-updates", h.systemCheckUpdates)
|
||||
mux.HandleFunc("/api/v1/system/upgrade", h.systemUpgrade)
|
||||
mux.HandleFunc("/api/v1/license/activate", h.licenseActivate)
|
||||
mux.HandleFunc("/api/v1/backup/export", h.backupExport)
|
||||
mux.HandleFunc("/api/v1/backup/import", h.backupImport)
|
||||
|
||||
@@ -234,8 +234,22 @@ func (h *Handler) monitorTunnelQualityHandler(w http.ResponseWriter, r *http.Req
|
||||
return
|
||||
}
|
||||
|
||||
targetsByTunnelID := map[int64]tunnelProbeTarget{}
|
||||
if tunnels, listErr := h.repo.ListTunnels(); listErr == nil {
|
||||
for _, item := range tunnels {
|
||||
id := asInt64(item["id"], 0)
|
||||
if id > 0 {
|
||||
targetsByTunnelID[id] = effectiveTunnelProbeTargetValues(asString(item["probeTargetHost"]), asInt(item["probeTargetPort"], 0))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
snapshots := make([]tunnelQualitySnapshot, 0, len(qualities))
|
||||
for _, q := range qualities {
|
||||
target := targetsByTunnelID[q.TunnelID]
|
||||
if target.Host == "" {
|
||||
target = defaultTunnelProbeTarget()
|
||||
}
|
||||
snapshots = append(snapshots, tunnelQualitySnapshot{
|
||||
TunnelID: q.TunnelID,
|
||||
EntryToExitLatency: q.EntryToExitLatency,
|
||||
@@ -246,6 +260,8 @@ func (h *Handler) monitorTunnelQualityHandler(w http.ResponseWriter, r *http.Req
|
||||
ErrorMessage: q.ErrorMessage,
|
||||
Timestamp: q.Timestamp,
|
||||
ChainDetails: q.ChainDetails,
|
||||
ProbeTargetHost: target.Host,
|
||||
ProbeTargetPort: target.Port,
|
||||
})
|
||||
}
|
||||
response.WriteJSON(w, response.OK(snapshots))
|
||||
|
||||
@@ -607,6 +607,17 @@ func (h *Handler) tunnelCreate(w http.ResponseWriter, r *http.Request) {
|
||||
trafficRatio := asFloat(req["trafficRatio"], 1.0)
|
||||
inIP := asString(req["inIp"])
|
||||
ipPreference := asString(req["ipPreference"])
|
||||
probeTarget, probeTargetConfigured, err := parseTunnelProbeTargetFromRequest(req)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.ErrDefault(err.Error()))
|
||||
return
|
||||
}
|
||||
probeTargetHost := ""
|
||||
probeTargetPort := 0
|
||||
if probeTargetConfigured {
|
||||
probeTargetHost = probeTarget.Host
|
||||
probeTargetPort = probeTarget.Port
|
||||
}
|
||||
now := time.Now().UnixMilli()
|
||||
inx := h.repo.NextIndex("tunnel")
|
||||
localDomain := h.federationLocalDomain()
|
||||
@@ -685,17 +696,19 @@ func (h *Handler) tunnelCreate(w http.ResponseWriter, r *http.Request) {
|
||||
tunnelProtocol = strings.TrimSpace(runtimeState.InNodes[0].Protocol)
|
||||
}
|
||||
tunnel := model.Tunnel{
|
||||
Name: name,
|
||||
TrafficRatio: trafficRatio,
|
||||
Type: typeVal,
|
||||
Protocol: tunnelProtocol,
|
||||
Flow: flow,
|
||||
CreatedTime: now,
|
||||
UpdatedTime: now,
|
||||
Status: status,
|
||||
InIP: tunnelInIP,
|
||||
Inx: inx,
|
||||
IPPreference: ipPreference,
|
||||
Name: name,
|
||||
TrafficRatio: trafficRatio,
|
||||
Type: typeVal,
|
||||
Protocol: tunnelProtocol,
|
||||
Flow: flow,
|
||||
CreatedTime: now,
|
||||
UpdatedTime: now,
|
||||
Status: status,
|
||||
InIP: tunnelInIP,
|
||||
Inx: inx,
|
||||
IPPreference: ipPreference,
|
||||
ProbeTargetHost: probeTargetHost,
|
||||
ProbeTargetPort: probeTargetPort,
|
||||
}
|
||||
if err := tx.Create(&tunnel).Error; err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
@@ -871,9 +884,30 @@ func (h *Handler) tunnelUpdate(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.ErrDefault("隧道ID不能为空"))
|
||||
return
|
||||
}
|
||||
oldEntryNodeIDs, _ := h.tunnelEntryNodeIDs(id)
|
||||
typeVal := asInt(req["type"], 1)
|
||||
ipPreference := asString(req["ipPreference"])
|
||||
_, hasProbeTargetHost := req["probeTargetHost"]
|
||||
_, hasProbeTargetPort := req["probeTargetPort"]
|
||||
probeTargetFieldsPresent := hasProbeTargetHost || hasProbeTargetPort
|
||||
probeTargetHost := ""
|
||||
probeTargetPort := 0
|
||||
if probeTargetFieldsPresent {
|
||||
probeTarget, probeTargetConfigured, err := parseTunnelProbeTargetFromRequest(req)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.ErrDefault(err.Error()))
|
||||
return
|
||||
}
|
||||
if probeTargetConfigured {
|
||||
probeTargetHost = probeTarget.Host
|
||||
probeTargetPort = probeTarget.Port
|
||||
}
|
||||
}
|
||||
oldEntryNodeIDs, _ := h.tunnelEntryNodeIDs(id)
|
||||
oldTunnel, _ := h.getTunnelRecord(id)
|
||||
if !probeTargetFieldsPresent && oldTunnel != nil {
|
||||
probeTargetHost = oldTunnel.ProbeTargetHost
|
||||
probeTargetPort = oldTunnel.ProbeTargetPort
|
||||
}
|
||||
oldChainRows, _ := h.listChainNodesForTunnel(id)
|
||||
if oldTunnel != nil && oldTunnel.Type == 2 && typeVal != 2 {
|
||||
h.cleanupTunnelRuntime(id)
|
||||
@@ -881,7 +915,6 @@ func (h *Handler) tunnelUpdate(w http.ResponseWriter, r *http.Request) {
|
||||
h.cleanupFederationRuntime(id)
|
||||
|
||||
now := time.Now().UnixMilli()
|
||||
ipPreference := asString(req["ipPreference"])
|
||||
localDomain := h.federationLocalDomain()
|
||||
|
||||
runtimeState, err := h.prepareTunnelCreateState(h.repo.DB(), req, typeVal, id)
|
||||
@@ -928,6 +961,8 @@ func (h *Handler) tunnelUpdate(w http.ResponseWriter, r *http.Request) {
|
||||
inIp,
|
||||
ipPreference,
|
||||
updateProtocol,
|
||||
probeTargetHost,
|
||||
probeTargetPort,
|
||||
now,
|
||||
); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
|
||||
@@ -0,0 +1,592 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go-backend/internal/http/response"
|
||||
)
|
||||
|
||||
const (
|
||||
panelDeployDirEnv = "PANEL_DEPLOY_DIR"
|
||||
panelBackendContainerEnv = "PANEL_BACKEND_CONTAINER"
|
||||
defaultPanelDeployDir = "/opt/flvx-panel"
|
||||
defaultPanelBackendName = "flux-panel-backend"
|
||||
dockerSocketPath = "/var/run/docker.sock"
|
||||
maxSystemUpgradeComposeAssetBytes = 1 << 20
|
||||
systemUpgradeMessage = "升级 helper 已启动,面板服务将短暂重启"
|
||||
systemUpgradeConflictError = "已有面板升级任务执行中"
|
||||
)
|
||||
|
||||
var safeBackendContainerPattern = regexp.MustCompile(`^[A-Za-z0-9_.-]+$`)
|
||||
var enableIPv6ComposePattern = regexp.MustCompile(`(?im)^\s*enable_ipv6\s*:\s*['"]?true['"]?\s*(?:#.*)?$`)
|
||||
var systemUpgradeReleaseBaseURL = githubHTMLBase
|
||||
|
||||
type systemUpgradeExecutor struct {
|
||||
deployDir string
|
||||
backendContainer string
|
||||
}
|
||||
|
||||
type systemUpgradeCapabilityData struct {
|
||||
Capable bool `json:"capable"`
|
||||
Reasons []string `json:"reasons"`
|
||||
DeployDir string `json:"deployDir"`
|
||||
BackendContainer string `json:"backendContainer"`
|
||||
}
|
||||
|
||||
type systemUpgradeReleaseData struct {
|
||||
Version string `json:"version"`
|
||||
Name string `json:"name"`
|
||||
PublishedAt string `json:"publishedAt"`
|
||||
Prerelease bool `json:"prerelease"`
|
||||
Channel string `json:"channel"`
|
||||
}
|
||||
|
||||
type systemUpgradeVersionData struct {
|
||||
CurrentVersion string `json:"currentVersion"`
|
||||
LatestVersion string `json:"latestVersion"`
|
||||
HasUpdate bool `json:"hasUpdate"`
|
||||
Channel string `json:"channel"`
|
||||
Reason string `json:"reason,omitempty"`
|
||||
Capability systemUpgradeCapabilityData `json:"capability"`
|
||||
}
|
||||
|
||||
type systemUpgradeCheckData struct {
|
||||
CurrentVersion string `json:"currentVersion"`
|
||||
LatestVersion string `json:"latestVersion"`
|
||||
HasUpdate bool `json:"hasUpdate"`
|
||||
Channel string `json:"channel"`
|
||||
Capability systemUpgradeCapabilityData `json:"capability"`
|
||||
Releases []systemUpgradeReleaseData `json:"releases"`
|
||||
}
|
||||
|
||||
type systemUpgradeRunData struct {
|
||||
Version string `json:"version"`
|
||||
Channel string `json:"channel"`
|
||||
ComposeAsset string `json:"composeAsset"`
|
||||
HelperContainer string `json:"helperContainer"`
|
||||
BackendImageID string `json:"backendImageId"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
type systemUpgradeRequest struct {
|
||||
Version string `json:"version"`
|
||||
Channel string `json:"channel"`
|
||||
}
|
||||
|
||||
func newSystemUpgradeExecutor() *systemUpgradeExecutor {
|
||||
deployDir := strings.TrimSpace(os.Getenv(panelDeployDirEnv))
|
||||
if deployDir == "" {
|
||||
deployDir = defaultPanelDeployDir
|
||||
}
|
||||
backendContainer := strings.TrimSpace(os.Getenv(panelBackendContainerEnv))
|
||||
if backendContainer == "" {
|
||||
backendContainer = defaultPanelBackendName
|
||||
}
|
||||
return &systemUpgradeExecutor{deployDir: deployDir, backendContainer: backendContainer}
|
||||
}
|
||||
|
||||
func currentPanelVersion() string {
|
||||
version := strings.TrimSpace(os.Getenv("FLUX_VERSION"))
|
||||
if version == "" {
|
||||
return "dev"
|
||||
}
|
||||
return version
|
||||
}
|
||||
|
||||
func validateBackendContainerName(value string) error {
|
||||
if value == "" {
|
||||
return fmt.Errorf("backend container name is empty")
|
||||
}
|
||||
if !safeBackendContainerPattern.MatchString(value) {
|
||||
return fmt.Errorf("unsafe backend container name: %s", value)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateUpgradeVersion(value string) error {
|
||||
if strings.TrimSpace(value) == "" {
|
||||
return fmt.Errorf("upgrade version is empty")
|
||||
}
|
||||
for _, r := range value {
|
||||
if r < 0x20 || r == 0x7f {
|
||||
return fmt.Errorf("unsafe upgrade version: contains control character")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) composePath() string {
|
||||
return filepath.Join(e.deployDir, "docker-compose.yml")
|
||||
}
|
||||
func (e *systemUpgradeExecutor) envPath() string { return filepath.Join(e.deployDir, ".env") }
|
||||
|
||||
func (e *systemUpgradeExecutor) capability(ctx context.Context) systemUpgradeCapabilityData {
|
||||
reasons := make([]string, 0)
|
||||
if !filepath.IsAbs(e.deployDir) {
|
||||
reasons = append(reasons, "部署目录必须是绝对路径")
|
||||
}
|
||||
if err := validateBackendContainerName(e.backendContainer); err != nil {
|
||||
reasons = append(reasons, err.Error())
|
||||
}
|
||||
if out, err := exec.CommandContext(ctx, "docker", "--version").CombinedOutput(); err != nil {
|
||||
reasons = append(reasons, fmt.Sprintf("docker CLI不可用: %v: %s", err, strings.TrimSpace(string(out))))
|
||||
}
|
||||
if info, err := os.Stat(dockerSocketPath); err != nil {
|
||||
reasons = append(reasons, "docker socket不可用: "+err.Error())
|
||||
} else if info.IsDir() {
|
||||
reasons = append(reasons, "docker socket路径不是文件")
|
||||
}
|
||||
if info, err := os.Stat(e.composePath()); err != nil {
|
||||
reasons = append(reasons, "部署docker-compose.yml不可用: "+err.Error())
|
||||
} else if info.IsDir() {
|
||||
reasons = append(reasons, "部署docker-compose.yml不是文件")
|
||||
}
|
||||
if info, err := os.Stat(e.envPath()); err != nil {
|
||||
reasons = append(reasons, "部署.env不可用: "+err.Error())
|
||||
} else if info.IsDir() {
|
||||
reasons = append(reasons, "部署.env不是文件")
|
||||
}
|
||||
if out, err := exec.CommandContext(ctx, "docker", "compose", "version").CombinedOutput(); err != nil {
|
||||
reasons = append(reasons, fmt.Sprintf("docker compose不可用: %v: %s", err, strings.TrimSpace(string(out))))
|
||||
}
|
||||
if _, err := e.currentBackendImage(ctx); err != nil {
|
||||
reasons = append(reasons, err.Error())
|
||||
}
|
||||
|
||||
return systemUpgradeCapabilityData{
|
||||
Capable: len(reasons) == 0,
|
||||
Reasons: reasons,
|
||||
DeployDir: e.deployDir,
|
||||
BackendContainer: e.backendContainer,
|
||||
}
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) selectComposeAsset(current []byte) string {
|
||||
if enableIPv6ComposePattern.Match(current) {
|
||||
return "docker-compose-v6.yml"
|
||||
}
|
||||
return "docker-compose-v4.yml"
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) helperScript() string {
|
||||
return `set -eu
|
||||
LOGFILE="$PANEL_DEPLOY_DIR/upgrade.log"
|
||||
log() { echo "[$(date '+%Y-%m-%d %H:%M:%S')] $*" | tee -a "$LOGFILE"; }
|
||||
|
||||
cd "$PANEL_DEPLOY_DIR"
|
||||
echo "" > "$LOGFILE"
|
||||
log "开始面板升级"
|
||||
log "工作目录: $(pwd)"
|
||||
|
||||
if [ ! -f docker-compose.yml ]; then
|
||||
log "错误: docker-compose.yml 不存在"
|
||||
exit 1
|
||||
fi
|
||||
if [ ! -f .env ]; then
|
||||
log "错误: .env 不存在"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
log "拉取新镜像..."
|
||||
if ! docker compose pull backend frontend 2>&1 | tee -a "$LOGFILE"; then
|
||||
log "错误: 拉取镜像失败"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
log "等待旧容器释放资源..."
|
||||
sleep 3
|
||||
|
||||
log "重启服务(force-recreate)..."
|
||||
if ! docker compose up -d --force-recreate --remove-orphans backend frontend 2>&1 | tee -a "$LOGFILE"; then
|
||||
log "错误: 重启服务失败"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
log "升级完成"
|
||||
`
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) buildHelperRunArgs(imageID, helperName string) ([]string, error) {
|
||||
if err := validateBackendContainerName(e.backendContainer); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return []string{
|
||||
"run", "-d", "--rm", "--name", helperName,
|
||||
"--volumes-from", e.backendContainer,
|
||||
"-v", dockerSocketPath + ":" + dockerSocketPath,
|
||||
"-e", panelDeployDirEnv + "=" + e.deployDir,
|
||||
"--entrypoint", "/bin/sh", imageID,
|
||||
"-c", e.helperScript(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) updateEnvVersion(envPath, version string) error {
|
||||
if err := validateUpgradeVersion(version); err != nil {
|
||||
return err
|
||||
}
|
||||
mode, err := fileModeOrDefault(envPath, 0o600)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
data, err := os.ReadFile(envPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
lines := strings.Split(string(data), "\n")
|
||||
replaced := false
|
||||
for i, line := range lines {
|
||||
if strings.HasPrefix(line, "FLUX_VERSION=") {
|
||||
lines[i] = "FLUX_VERSION=" + version
|
||||
replaced = true
|
||||
}
|
||||
}
|
||||
if !replaced {
|
||||
trimmed := strings.TrimRight(strings.Join(lines, "\n"), "\n")
|
||||
if trimmed == "" {
|
||||
trimmed = "FLUX_VERSION=" + version
|
||||
} else {
|
||||
trimmed += "\nFLUX_VERSION=" + version
|
||||
}
|
||||
return writeFileWithMode(envPath, []byte(trimmed+"\n"), mode)
|
||||
}
|
||||
content := strings.TrimRight(strings.Join(lines, "\n"), "\n") + "\n"
|
||||
return writeFileWithMode(envPath, []byte(content), mode)
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) backupFile(path string) (string, error) {
|
||||
mode, err := fileModeOrDefault(path, 0o600)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
data, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
backupPath := path + ".upgrade.bak"
|
||||
if err := writeFileWithMode(backupPath, data, mode); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return backupPath, nil
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) restoreBackup(path string) error {
|
||||
backupPath := path + ".upgrade.bak"
|
||||
mode, err := fileModeOrDefault(backupPath, 0o600)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
data, err := os.ReadFile(backupPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return writeFileWithMode(path, data, mode)
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) restoreUpgradeBackups(paths ...string) error {
|
||||
var errs []string
|
||||
for _, path := range paths {
|
||||
if err := e.restoreBackup(path); err != nil {
|
||||
errs = append(errs, fmt.Sprintf("%s: %v", path, err))
|
||||
}
|
||||
}
|
||||
if len(errs) > 0 {
|
||||
return fmt.Errorf("%s", strings.Join(errs, "; "))
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) replaceCompose(path string, data []byte) error {
|
||||
if len(bytes.TrimSpace(data)) == 0 {
|
||||
return fmt.Errorf("compose asset is empty")
|
||||
}
|
||||
mode, err := fileModeOrDefault(path, 0o644)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return writeFileWithMode(path, data, mode)
|
||||
}
|
||||
|
||||
func fileModeOrDefault(path string, fallback os.FileMode) (os.FileMode, error) {
|
||||
info, err := os.Stat(path)
|
||||
if err != nil {
|
||||
if os.IsNotExist(err) {
|
||||
return fallback, nil
|
||||
}
|
||||
return 0, err
|
||||
}
|
||||
return info.Mode().Perm(), nil
|
||||
}
|
||||
|
||||
func writeFileWithMode(path string, data []byte, mode os.FileMode) error {
|
||||
if err := os.WriteFile(path, data, mode); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.Chmod(path, mode)
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) currentBackendImage(ctx context.Context) (string, error) {
|
||||
if err := validateBackendContainerName(e.backendContainer); err != nil {
|
||||
return "", err
|
||||
}
|
||||
out, err := exec.CommandContext(ctx, "docker", "inspect", "-f", "{{.Image}}", e.backendContainer).CombinedOutput()
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("inspect backend image failed: %v: %s", err, strings.TrimSpace(string(out)))
|
||||
}
|
||||
imageID := strings.TrimSpace(string(out))
|
||||
if imageID == "" {
|
||||
return "", fmt.Errorf("backend image id is empty")
|
||||
}
|
||||
return imageID, nil
|
||||
}
|
||||
|
||||
func (e *systemUpgradeExecutor) startHelper(ctx context.Context, imageID, helperName string) (string, error) {
|
||||
args, err := e.buildHelperRunArgs(imageID, helperName)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
out, err := exec.CommandContext(ctx, "docker", args...).CombinedOutput()
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("start helper failed: %v: %s", err, strings.TrimSpace(string(out)))
|
||||
}
|
||||
containerID := strings.TrimSpace(string(out))
|
||||
if containerID == "" {
|
||||
containerID = helperName
|
||||
}
|
||||
return containerID, nil
|
||||
}
|
||||
|
||||
func (h *Handler) downloadReleaseAsset(version, filename string) ([]byte, error) {
|
||||
url := fmt.Sprintf("%s/%s/releases/download/%s/%s", strings.TrimRight(systemUpgradeReleaseBaseURL, "/"), githubRepo, version, filename)
|
||||
client := &http.Client{Timeout: 60 * time.Second}
|
||||
resp, err := client.Get(url)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("下载%s失败: %v", filename, err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
body, _ := io.ReadAll(io.LimitReader(resp.Body, 1024))
|
||||
return nil, fmt.Errorf("下载%s返回 %d: %s", filename, resp.StatusCode, strings.TrimSpace(string(body)))
|
||||
}
|
||||
body, err := io.ReadAll(io.LimitReader(resp.Body, maxSystemUpgradeComposeAssetBytes+1))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("读取%s失败: %v", filename, err)
|
||||
}
|
||||
if len(body) > maxSystemUpgradeComposeAssetBytes {
|
||||
return nil, fmt.Errorf("下载%s过大", filename)
|
||||
}
|
||||
if len(bytes.TrimSpace(body)) == 0 {
|
||||
return nil, fmt.Errorf("下载%s内容为空", filename)
|
||||
}
|
||||
return body, nil
|
||||
}
|
||||
|
||||
func releasesForChannel(releases []githubRelease, channel string) []systemUpgradeReleaseData {
|
||||
channel = normalizeReleaseChannel(channel)
|
||||
items := make([]systemUpgradeReleaseData, 0, len(releases))
|
||||
for _, r := range releases {
|
||||
if r.Draft {
|
||||
continue
|
||||
}
|
||||
tag := strings.TrimSpace(r.TagName)
|
||||
if tag == "" {
|
||||
continue
|
||||
}
|
||||
itemChannel := releaseChannelFromTag(tag)
|
||||
if itemChannel != channel {
|
||||
continue
|
||||
}
|
||||
items = append(items, systemUpgradeReleaseData{
|
||||
Version: tag,
|
||||
Name: r.Name,
|
||||
PublishedAt: r.PublishedAt,
|
||||
Prerelease: itemChannel == releaseChannelDev,
|
||||
Channel: itemChannel,
|
||||
})
|
||||
}
|
||||
return items
|
||||
}
|
||||
|
||||
func decodeSystemUpgradeRequest(r *http.Request, req *systemUpgradeRequest) error {
|
||||
defer r.Body.Close()
|
||||
body, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(bytes.TrimSpace(body)) == 0 {
|
||||
return nil
|
||||
}
|
||||
decoder := json.NewDecoder(bytes.NewReader(body))
|
||||
decoder.DisallowUnknownFields()
|
||||
return decoder.Decode(req)
|
||||
}
|
||||
|
||||
func systemUpgradeVersionResponse(current, channel, latest string, lookupErr error, capability systemUpgradeCapabilityData) systemUpgradeVersionData {
|
||||
data := systemUpgradeVersionData{
|
||||
CurrentVersion: current,
|
||||
LatestVersion: latest,
|
||||
HasUpdate: latest != "" && latest != current,
|
||||
Channel: channel,
|
||||
Capability: capability,
|
||||
}
|
||||
if lookupErr != nil {
|
||||
data.LatestVersion = ""
|
||||
data.HasUpdate = false
|
||||
data.Reason = lookupErr.Error()
|
||||
}
|
||||
return data
|
||||
}
|
||||
|
||||
func (h *Handler) systemVersion(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("请求失败"))
|
||||
return
|
||||
}
|
||||
|
||||
channel := releaseChannelStable
|
||||
current := currentPanelVersion()
|
||||
exec := newSystemUpgradeExecutor()
|
||||
capability := exec.capability(r.Context())
|
||||
latest, err := resolveLatestReleaseByChannel(channel)
|
||||
response.WriteJSON(w, response.OK(systemUpgradeVersionResponse(current, channel, latest, err, capability)))
|
||||
}
|
||||
|
||||
func (h *Handler) systemCheckUpdates(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("请求失败"))
|
||||
return
|
||||
}
|
||||
|
||||
var req systemUpgradeRequest
|
||||
if err := decodeSystemUpgradeRequest(r, &req); err != nil {
|
||||
response.WriteJSON(w, response.ErrDefault("请求参数错误"))
|
||||
return
|
||||
}
|
||||
channel := normalizeReleaseChannel(req.Channel)
|
||||
current := currentPanelVersion()
|
||||
exec := newSystemUpgradeExecutor()
|
||||
capability := exec.capability(r.Context())
|
||||
|
||||
githubReleases, err := fetchGitHubReleases(50)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("获取版本列表失败: %v", err)))
|
||||
return
|
||||
}
|
||||
releases := releasesForChannel(githubReleases, channel)
|
||||
latest := ""
|
||||
if len(releases) > 0 {
|
||||
latest = releases[0].Version
|
||||
}
|
||||
response.WriteJSON(w, response.OK(systemUpgradeCheckData{
|
||||
CurrentVersion: current,
|
||||
LatestVersion: latest,
|
||||
HasUpdate: latest != "" && latest != current,
|
||||
Channel: channel,
|
||||
Capability: capability,
|
||||
Releases: releases,
|
||||
}))
|
||||
}
|
||||
|
||||
func (h *Handler) systemUpgrade(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("请求失败"))
|
||||
return
|
||||
}
|
||||
if !h.systemUpgradeMu.TryLock() {
|
||||
response.WriteJSON(w, response.ErrDefault(systemUpgradeConflictError))
|
||||
return
|
||||
}
|
||||
defer h.systemUpgradeMu.Unlock()
|
||||
|
||||
var req systemUpgradeRequest
|
||||
if err := decodeSystemUpgradeRequest(r, &req); err != nil {
|
||||
response.WriteJSON(w, response.ErrDefault("请求参数错误"))
|
||||
return
|
||||
}
|
||||
channel := normalizeReleaseChannel(req.Channel)
|
||||
version := strings.TrimSpace(req.Version)
|
||||
if version == "" {
|
||||
var err error
|
||||
version, err = resolveLatestReleaseByChannel(channel)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("获取最新%s失败: %v", releaseChannelLabel(channel), err)))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
exec := newSystemUpgradeExecutor()
|
||||
capability := exec.capability(r.Context())
|
||||
if !capability.Capable {
|
||||
response.WriteJSON(w, response.ErrDefault("当前环境不支持面板自升级: "+strings.Join(capability.Reasons, "; ")))
|
||||
return
|
||||
}
|
||||
imageID, err := exec.currentBackendImage(r.Context())
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
composePath := exec.composePath()
|
||||
envPath := exec.envPath()
|
||||
composeData, err := os.ReadFile(composePath)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, "读取compose失败: "+err.Error()))
|
||||
return
|
||||
}
|
||||
composeAsset := exec.selectComposeAsset(composeData)
|
||||
newCompose, err := h.downloadReleaseAsset(version, composeAsset)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
if _, err := exec.backupFile(composePath); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, "备份compose失败: "+err.Error()))
|
||||
return
|
||||
}
|
||||
if _, err := exec.backupFile(envPath); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, "备份.env失败: "+err.Error()))
|
||||
return
|
||||
}
|
||||
if err := exec.replaceCompose(composePath, newCompose); err != nil {
|
||||
if restoreErr := exec.restoreUpgradeBackups(composePath, envPath); restoreErr != nil {
|
||||
err = fmt.Errorf("%v; 回滚失败: %v", err, restoreErr)
|
||||
}
|
||||
response.WriteJSON(w, response.Err(-2, "替换compose失败: "+err.Error()))
|
||||
return
|
||||
}
|
||||
if err := exec.updateEnvVersion(envPath, version); err != nil {
|
||||
if restoreErr := exec.restoreUpgradeBackups(composePath, envPath); restoreErr != nil {
|
||||
err = fmt.Errorf("%v; 回滚失败: %v", err, restoreErr)
|
||||
}
|
||||
response.WriteJSON(w, response.Err(-2, "更新版本配置失败: "+err.Error()))
|
||||
return
|
||||
}
|
||||
helperName := fmt.Sprintf("flvx-upgrade-helper-%d", time.Now().Unix())
|
||||
helperContainer, err := exec.startHelper(r.Context(), imageID, helperName)
|
||||
if err != nil {
|
||||
if restoreErr := exec.restoreUpgradeBackups(composePath, envPath); restoreErr != nil {
|
||||
err = fmt.Errorf("%v; 回滚失败: %v", err, restoreErr)
|
||||
}
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
response.WriteJSON(w, response.OK(systemUpgradeRunData{
|
||||
Version: version,
|
||||
Channel: channel,
|
||||
ComposeAsset: composeAsset,
|
||||
HelperContainer: helperContainer,
|
||||
BackendImageID: imageID,
|
||||
Message: systemUpgradeMessage,
|
||||
}))
|
||||
}
|
||||
@@ -0,0 +1,398 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestSelectComposeAssetUsesIPv6Template(t *testing.T) {
|
||||
exec := &systemUpgradeExecutor{deployDir: "/opt/flvx-panel", backendContainer: "flux-panel-backend"}
|
||||
compose := []byte("networks:\n gost-network:\n enable_ipv6: true\n")
|
||||
|
||||
if got := exec.selectComposeAsset(compose); got != "docker-compose-v6.yml" {
|
||||
t.Fatalf("selectComposeAsset() = %q, want %q", got, "docker-compose-v6.yml")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDownloadReleaseAssetUsesDirectReleaseURL(t *testing.T) {
|
||||
var gotPath string
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
gotPath = r.URL.Path
|
||||
_, _ = w.Write([]byte("services:\n backend:\n image: test\n"))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
originalBase := systemUpgradeReleaseBaseURL
|
||||
systemUpgradeReleaseBaseURL = server.URL
|
||||
t.Cleanup(func() { systemUpgradeReleaseBaseURL = originalBase })
|
||||
|
||||
h := &Handler{}
|
||||
data, err := h.downloadReleaseAsset("2.1.9", "docker-compose-v4.yml")
|
||||
if err != nil {
|
||||
t.Fatalf("downloadReleaseAsset() error = %v", err)
|
||||
}
|
||||
if !strings.Contains(string(data), "backend") {
|
||||
t.Fatalf("downloadReleaseAsset() data = %q, want compose data", string(data))
|
||||
}
|
||||
|
||||
wantPath := "/" + githubRepo + "/releases/download/2.1.9/docker-compose-v4.yml"
|
||||
if gotPath != wantPath {
|
||||
t.Fatalf("download path = %q, want %q", gotPath, wantPath)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDownloadReleaseAssetRejectsOversizedBody(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
_, _ = w.Write(bytes.Repeat([]byte("a"), maxSystemUpgradeComposeAssetBytes+1))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
originalBase := systemUpgradeReleaseBaseURL
|
||||
systemUpgradeReleaseBaseURL = server.URL
|
||||
t.Cleanup(func() { systemUpgradeReleaseBaseURL = originalBase })
|
||||
|
||||
h := &Handler{}
|
||||
_, err := h.downloadReleaseAsset("2.1.9", "docker-compose-v4.yml")
|
||||
if err == nil || !strings.Contains(err.Error(), "过大") {
|
||||
t.Fatalf("downloadReleaseAsset() error = %v, want oversized error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSelectComposeAssetUsesIPv6TemplateForYAMLVariants(t *testing.T) {
|
||||
exec := &systemUpgradeExecutor{deployDir: "/opt/flvx-panel", backendContainer: "flux-panel-backend"}
|
||||
for _, compose := range [][]byte{
|
||||
[]byte("networks:\n gost-network:\n enable_ipv6:true\n"),
|
||||
[]byte("networks:\n gost-network:\n enable_ipv6: True\n"),
|
||||
[]byte("networks:\n gost-network:\n enable_ipv6: \"true\"\n"),
|
||||
[]byte("networks:\n gost-network:\n enable_ipv6: 'true'\n"),
|
||||
[]byte("networks:\n gost-network:\n enable_ipv6: true # comment\n"),
|
||||
} {
|
||||
if got := exec.selectComposeAsset(compose); got != "docker-compose-v6.yml" {
|
||||
t.Fatalf("selectComposeAsset(%q) = %q, want %q", string(compose), got, "docker-compose-v6.yml")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSelectComposeAssetFallsBackToIPv4Template(t *testing.T) {
|
||||
exec := &systemUpgradeExecutor{deployDir: "/opt/flvx-panel", backendContainer: "flux-panel-backend"}
|
||||
compose := []byte("services:\n backend:\n image: test\n")
|
||||
|
||||
if got := exec.selectComposeAsset(compose); got != "docker-compose-v4.yml" {
|
||||
t.Fatalf("selectComposeAsset() = %q, want %q", got, "docker-compose-v4.yml")
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateEnvVersionReplacesExistingValue(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
envPath := filepath.Join(dir, ".env")
|
||||
if err := os.WriteFile(envPath, []byte("FLUX_VERSION=2.1.8\nJWT_SECRET=test\n"), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
if err := exec.updateEnvVersion(envPath, "2.1.9"); err != nil {
|
||||
t.Fatalf("updateEnvVersion() error = %v", err)
|
||||
}
|
||||
|
||||
data, err := os.ReadFile(envPath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile() error = %v", err)
|
||||
}
|
||||
|
||||
want := "FLUX_VERSION=2.1.9\nJWT_SECRET=test\n"
|
||||
if string(data) != want {
|
||||
t.Fatalf("env content = %q, want %q", string(data), want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateEnvVersionAppendsMissingValue(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
envPath := filepath.Join(dir, ".env")
|
||||
if err := os.WriteFile(envPath, []byte("JWT_SECRET=test\n"), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
if err := exec.updateEnvVersion(envPath, "2.1.9"); err != nil {
|
||||
t.Fatalf("updateEnvVersion() error = %v", err)
|
||||
}
|
||||
|
||||
data, err := os.ReadFile(envPath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile() error = %v", err)
|
||||
}
|
||||
|
||||
want := "JWT_SECRET=test\nFLUX_VERSION=2.1.9\n"
|
||||
if string(data) != want {
|
||||
t.Fatalf("env content = %q, want %q", string(data), want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateEnvVersionRejectsUnsafeValue(t *testing.T) {
|
||||
for _, version := range []string{"", "2.1.9\nJWT_SECRET=bad", "2.1.9\rbad", "2.1.9\x00bad", "2.1.9\x1fbad"} {
|
||||
t.Run(version, func(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
envPath := filepath.Join(dir, ".env")
|
||||
original := []byte("JWT_SECRET=test\n")
|
||||
if err := os.WriteFile(envPath, original, 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
if err := exec.updateEnvVersion(envPath, version); err == nil {
|
||||
t.Fatal("expected unsafe version to fail validation")
|
||||
}
|
||||
|
||||
data, err := os.ReadFile(envPath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile() error = %v", err)
|
||||
}
|
||||
if string(data) != string(original) {
|
||||
t.Fatalf("env content changed to %q, want %q", string(data), string(original))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateEnvVersionAcceptsVersionLabels(t *testing.T) {
|
||||
for _, version := range []string{"2.1.9", "2.1.9-beta14", "v-test"} {
|
||||
t.Run(version, func(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
envPath := filepath.Join(dir, ".env")
|
||||
if err := os.WriteFile(envPath, []byte("JWT_SECRET=test\n"), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
if err := exec.updateEnvVersion(envPath, version); err != nil {
|
||||
t.Fatalf("updateEnvVersion() error = %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpdateEnvVersionPreservesFileMode(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
envPath := filepath.Join(dir, ".env")
|
||||
if err := os.WriteFile(envPath, []byte("FLUX_VERSION=2.1.8\nJWT_SECRET=test\n"), 0o600); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
if err := exec.updateEnvVersion(envPath, "2.1.9"); err != nil {
|
||||
t.Fatalf("updateEnvVersion() error = %v", err)
|
||||
}
|
||||
|
||||
info, err := os.Stat(envPath)
|
||||
if err != nil {
|
||||
t.Fatalf("Stat() error = %v", err)
|
||||
}
|
||||
if got := info.Mode().Perm(); got != 0o600 {
|
||||
t.Fatalf("env mode = %o, want 0600", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateBackendContainerNameRejectsUnsafeValue(t *testing.T) {
|
||||
if err := validateBackendContainerName("flux-panel-backend;rm -rf /"); err == nil {
|
||||
t.Fatal("expected unsafe container name to fail validation")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildHelperRunArgsUsesDetachedContainer(t *testing.T) {
|
||||
exec := &systemUpgradeExecutor{deployDir: "/opt/flvx-panel", backendContainer: "flux-panel-backend"}
|
||||
args, err := exec.buildHelperRunArgs("sha256:abc", "flvx-upgrade-helper")
|
||||
if err != nil {
|
||||
t.Fatalf("buildHelperRunArgs() error = %v", err)
|
||||
}
|
||||
want := []string{
|
||||
"run", "-d", "--rm", "--name", "flvx-upgrade-helper",
|
||||
"--volumes-from", "flux-panel-backend",
|
||||
"-v", "/var/run/docker.sock:/var/run/docker.sock",
|
||||
"-e", "PANEL_DEPLOY_DIR=/opt/flvx-panel",
|
||||
"--entrypoint", "/bin/sh", "sha256:abc",
|
||||
"-c", exec.helperScript(),
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(args, want) {
|
||||
t.Fatalf("buildHelperRunArgs() = %#v, want %#v", args, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildHelperRunArgsRejectsUnsafeBackendContainer(t *testing.T) {
|
||||
exec := &systemUpgradeExecutor{deployDir: "/opt/flvx-panel", backendContainer: "flux-panel-backend;rm -rf /"}
|
||||
if _, err := exec.buildHelperRunArgs("sha256:abc", "flvx-upgrade-helper"); err == nil {
|
||||
t.Fatal("expected unsafe backend container name to fail validation")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemVersionRejectsWrongMethod(t *testing.T) {
|
||||
h := &Handler{}
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/system/version", nil)
|
||||
rr := httptest.NewRecorder()
|
||||
|
||||
h.systemVersion(rr, req)
|
||||
|
||||
if !strings.Contains(rr.Body.String(), "请求失败") {
|
||||
t.Fatalf("expected wrong-method response, got %s", rr.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemUpgradeRejectsConcurrentRequests(t *testing.T) {
|
||||
h := &Handler{}
|
||||
h.systemUpgradeMu.Lock()
|
||||
defer h.systemUpgradeMu.Unlock()
|
||||
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/v1/system/upgrade", strings.NewReader(`{"channel":"stable"}`))
|
||||
rr := httptest.NewRecorder()
|
||||
|
||||
h.systemUpgrade(rr, req)
|
||||
|
||||
if !strings.Contains(rr.Body.String(), systemUpgradeConflictError) {
|
||||
t.Fatalf("expected conflict message, got %s", rr.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemUpgradeFailsFastBeforeMutatingFiles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
composePath := filepath.Join(dir, "docker-compose.yml")
|
||||
envPath := filepath.Join(dir, ".env")
|
||||
if err := os.WriteFile(composePath, []byte("services:\n backend:\n image: test\n"), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() compose error = %v", err)
|
||||
}
|
||||
if err := os.WriteFile(envPath, []byte("FLUX_VERSION=2.1.8\nJWT_SECRET=test\n"), 0o600); err != nil {
|
||||
t.Fatalf("WriteFile() env error = %v", err)
|
||||
}
|
||||
|
||||
fakeDockerDir := t.TempDir()
|
||||
fakeDockerPath := filepath.Join(fakeDockerDir, "docker")
|
||||
fakeDockerScript := "#!/bin/sh\ncase \"$1\" in\n --version)\n echo 'Docker version 27.0.0'\n exit 0\n ;;&\n compose)\n if [ \"$2\" = version ]; then\n echo 'Docker Compose version v2.33.0'\n exit 0\n fi\n exit 0\n ;;&\n inspect)\n echo 'No such object: flux-panel-backend' >&2\n exit 1\n ;;&\n *)\n exit 0\n ;;&\n esac\n"
|
||||
if err := os.WriteFile(fakeDockerPath, []byte(fakeDockerScript), 0o755); err != nil {
|
||||
t.Fatalf("WriteFile() fake docker error = %v", err)
|
||||
}
|
||||
t.Setenv("PATH", fakeDockerDir+string(os.PathListSeparator)+os.Getenv("PATH"))
|
||||
t.Setenv(panelDeployDirEnv, dir)
|
||||
t.Setenv(panelBackendContainerEnv, "flux-panel-backend")
|
||||
|
||||
h := &Handler{}
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/v1/system/upgrade", strings.NewReader(`{"channel":"stable"}`))
|
||||
rr := httptest.NewRecorder()
|
||||
|
||||
h.systemUpgrade(rr, req)
|
||||
|
||||
if !strings.Contains(rr.Body.String(), "当前环境不支持面板自升级") {
|
||||
t.Fatalf("expected fail-fast capability error, got %s", rr.Body.String())
|
||||
}
|
||||
if _, err := os.Stat(composePath + ".upgrade.bak"); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected no compose backup, got err=%v", err)
|
||||
}
|
||||
if _, err := os.Stat(envPath + ".upgrade.bak"); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected no env backup, got err=%v", err)
|
||||
}
|
||||
composeData, err := os.ReadFile(composePath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile() compose error = %v", err)
|
||||
}
|
||||
if string(composeData) != "services:\n backend:\n image: test\n" {
|
||||
t.Fatalf("compose mutated unexpectedly: %q", string(composeData))
|
||||
}
|
||||
envData, err := os.ReadFile(envPath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile() env error = %v", err)
|
||||
}
|
||||
if string(envData) != "FLUX_VERSION=2.1.8\nJWT_SECRET=test\n" {
|
||||
t.Fatalf("env mutated unexpectedly: %q", string(envData))
|
||||
}
|
||||
}
|
||||
|
||||
func TestUpgradeBackupUsesStablePathAndRestoreRestoresOriginal(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, "docker-compose.yml")
|
||||
if err := os.WriteFile(path, []byte("original"), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
backupPath, err := exec.backupFile(path)
|
||||
if err != nil {
|
||||
t.Fatalf("backupFile() error = %v", err)
|
||||
}
|
||||
if backupPath != path+".upgrade.bak" {
|
||||
t.Fatalf("backup path = %q, want %q", backupPath, path+".upgrade.bak")
|
||||
}
|
||||
if err := os.WriteFile(path, []byte("mutated"), 0o644); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
if err := exec.restoreBackup(path); err != nil {
|
||||
t.Fatalf("restoreBackup() error = %v", err)
|
||||
}
|
||||
data, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile() error = %v", err)
|
||||
}
|
||||
if string(data) != "original" {
|
||||
t.Fatalf("restored content = %q, want original", string(data))
|
||||
}
|
||||
}
|
||||
|
||||
func TestRestoreBackupPreservesOriginalFileMode(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, ".env")
|
||||
if err := os.WriteFile(path, []byte("FLUX_VERSION=2.1.8\nJWT_SECRET=test\n"), 0o600); err != nil {
|
||||
t.Fatalf("WriteFile() error = %v", err)
|
||||
}
|
||||
|
||||
exec := &systemUpgradeExecutor{deployDir: dir, backendContainer: "flux-panel-backend"}
|
||||
if _, err := exec.backupFile(path); err != nil {
|
||||
t.Fatalf("backupFile() error = %v", err)
|
||||
}
|
||||
if err := os.Remove(path); err != nil {
|
||||
t.Fatalf("Remove() error = %v", err)
|
||||
}
|
||||
if err := exec.restoreBackup(path); err != nil {
|
||||
t.Fatalf("restoreBackup() error = %v", err)
|
||||
}
|
||||
|
||||
info, err := os.Stat(path)
|
||||
if err != nil {
|
||||
t.Fatalf("Stat() error = %v", err)
|
||||
}
|
||||
if got := info.Mode().Perm(); got != 0o600 {
|
||||
t.Fatalf("restored mode = %o, want 0600", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDecodeSystemUpgradeRequestRejectsTruncatedJSON(t *testing.T) {
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/v1/system/check-updates", strings.NewReader(`{"channel":"stable"`))
|
||||
var payload systemUpgradeRequest
|
||||
|
||||
if err := decodeSystemUpgradeRequest(req, &payload); err == nil {
|
||||
t.Fatal("expected truncated JSON to be rejected")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDecodeSystemUpgradeRequestAllowsEmptyBody(t *testing.T) {
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/v1/system/check-updates", strings.NewReader(""))
|
||||
var payload systemUpgradeRequest
|
||||
|
||||
if err := decodeSystemUpgradeRequest(req, &payload); err != nil {
|
||||
t.Fatalf("expected empty body to be accepted, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSystemUpgradeVersionDataSurfacesLookupFailureReason(t *testing.T) {
|
||||
data, err := json.Marshal(systemUpgradeVersionData{Reason: "GitHub unavailable"})
|
||||
if err != nil {
|
||||
t.Fatalf("Marshal() error = %v", err)
|
||||
}
|
||||
if !strings.Contains(string(data), `"reason":"GitHub unavailable"`) {
|
||||
t.Fatalf("expected reason field in JSON, got %s", string(data))
|
||||
}
|
||||
}
|
||||
@@ -57,6 +57,12 @@ type bestExitProbeResult struct {
|
||||
err error
|
||||
}
|
||||
|
||||
type bestExitProbeCacheKey struct {
|
||||
NodeID int64
|
||||
Host string
|
||||
Port int
|
||||
}
|
||||
|
||||
type bestExitDecision struct {
|
||||
AppliedExitNodeID int64
|
||||
PendingExitNodeID int64
|
||||
@@ -138,7 +144,7 @@ func sortBestExitScores(scores []bestExitCandidateScore) {
|
||||
})
|
||||
}
|
||||
|
||||
func evaluateBestExitOwner(owner chainNodeRecord, exits []chainNodeRecord, nodes map[int64]*nodeRecord, ipPreference string, options diagnosisExecOptions, ping bestExitProbeFunc) []bestExitCandidateScore {
|
||||
func evaluateBestExitOwner(owner chainNodeRecord, exits []chainNodeRecord, nodes map[int64]*nodeRecord, ipPreference string, options diagnosisExecOptions, target tunnelProbeTarget, ping bestExitProbeFunc) []bestExitCandidateScore {
|
||||
scores := make([]bestExitCandidateScore, 0, len(exits))
|
||||
if owner.NodeID <= 0 || len(exits) == 0 || ping == nil {
|
||||
return scores
|
||||
@@ -160,7 +166,7 @@ func evaluateBestExitOwner(owner chainNodeRecord, exits []chainNodeRecord, nodes
|
||||
scores = append(scores, failedBestExitCandidate(owner.NodeID, exit, ownerErr.Error()))
|
||||
continue
|
||||
}
|
||||
publicLatency, publicLoss, publicErr := ping(exit.NodeID, bestExitPublicTargetHost, bestExitPublicTargetPort, options)
|
||||
publicLatency, publicLoss, publicErr := ping(exit.NodeID, target.Host, target.Port, options)
|
||||
if publicErr != nil {
|
||||
scores = append(scores, failedBestExitCandidate(owner.NodeID, exit, publicErr.Error()))
|
||||
continue
|
||||
@@ -193,17 +199,15 @@ func resolveBestExitProbeTarget(fromNode, targetNode *nodeRecord, preferredPort
|
||||
}
|
||||
|
||||
func newBestExitRoundPinger(base bestExitProbeFunc) bestExitProbeFunc {
|
||||
cache := make(map[int64]bestExitProbeResult)
|
||||
cache := make(map[bestExitProbeCacheKey]bestExitProbeResult)
|
||||
return func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
if ip == bestExitPublicTargetHost && port == bestExitPublicTargetPort {
|
||||
if cached, ok := cache[nodeID]; ok {
|
||||
return cached.latency, cached.loss, cached.err
|
||||
}
|
||||
lat, loss, err := base(nodeID, ip, port, options)
|
||||
cache[nodeID] = bestExitProbeResult{latency: lat, loss: loss, err: err}
|
||||
return lat, loss, err
|
||||
key := bestExitProbeCacheKey{NodeID: nodeID, Host: ip, Port: port}
|
||||
if cached, ok := cache[key]; ok {
|
||||
return cached.latency, cached.loss, cached.err
|
||||
}
|
||||
return base(nodeID, ip, port, options)
|
||||
lat, loss, err := base(nodeID, ip, port, options)
|
||||
cache[key] = bestExitProbeResult{latency: lat, loss: loss, err: err}
|
||||
return lat, loss, err
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -2,6 +2,9 @@ package handler
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"slices"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -189,7 +192,7 @@ func TestBestExitEnsureAppliedDoesNotOverrideExistingAppliedExit(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestBestExitRoundPingerCachesPublicProbeOnly(t *testing.T) {
|
||||
func TestBestExitRoundPingerCachesByNodeHostAndPort(t *testing.T) {
|
||||
publicCalls := 0
|
||||
ownerCalls := 0
|
||||
pinger := newBestExitRoundPinger(func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
@@ -220,8 +223,8 @@ func TestBestExitRoundPingerCachesPublicProbeOnly(t *testing.T) {
|
||||
if _, _, err := pinger(10, "10.0.0.30", 30030, diagnosisExecOptions{}); err != nil {
|
||||
t.Fatalf("unexpected repeated owner ping err=%v", err)
|
||||
}
|
||||
if ownerCalls != 2 {
|
||||
t.Fatalf("expected owner-to-exit probes not cached, got %d calls", ownerCalls)
|
||||
if ownerCalls != 1 {
|
||||
t.Fatalf("expected owner-to-exit probes cached by target, got %d calls", ownerCalls)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -375,7 +378,7 @@ func TestEvaluateBestExitOwnerScoresAllCandidates(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, pinger)
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, defaultTunnelProbeTarget(), pinger)
|
||||
if len(scores) != 2 {
|
||||
t.Fatalf("expected two scores, got %+v", scores)
|
||||
}
|
||||
@@ -384,6 +387,34 @@ func TestEvaluateBestExitOwnerScoresAllCandidates(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
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"},
|
||||
}
|
||||
target := tunnelProbeTarget{Host: "speed.example.com", Port: 8443}
|
||||
var calls []string
|
||||
ping := func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
calls = append(calls, fmt.Sprintf("%d|%s|%d", nodeID, ip, port))
|
||||
return 10, 0, nil
|
||||
}
|
||||
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, target, ping)
|
||||
if len(scores) != 1 || !scores[0].Success {
|
||||
t.Fatalf("expected successful score, got %+v", scores)
|
||||
}
|
||||
if !slices.Contains(calls, "30|speed.example.com|8443") {
|
||||
t.Fatalf("expected exit public probe to use configured target, calls=%+v", calls)
|
||||
}
|
||||
for _, call := range calls {
|
||||
if strings.Contains(call, defaultTunnelProbeTargetHost) {
|
||||
t.Fatalf("did not expect default target call when custom target configured: %+v", calls)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluateBestExitOwnerMarksCandidateFailedWhenOwnerToExitFails(t *testing.T) {
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30030}}
|
||||
@@ -395,7 +426,7 @@ func TestEvaluateBestExitOwnerMarksCandidateFailedWhenOwnerToExitFails(t *testin
|
||||
return 0, 100, errBestExitProbeForTest
|
||||
}
|
||||
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, pinger)
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, defaultTunnelProbeTarget(), pinger)
|
||||
if len(scores) != 1 || scores[0].Success {
|
||||
t.Fatalf("expected failed candidate, got %+v", scores)
|
||||
}
|
||||
@@ -413,7 +444,7 @@ func TestEvaluateBestExitOwnerMarksCandidateFailedWhenTargetResolutionFails(t *t
|
||||
return 0, 100, nil
|
||||
}
|
||||
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "v4", diagnosisExecOptions{}, pinger)
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "v4", diagnosisExecOptions{}, defaultTunnelProbeTarget(), pinger)
|
||||
if len(scores) != 1 || scores[0].Success {
|
||||
t.Fatalf("expected failed candidate, got %+v", scores)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/netip"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"go-backend/internal/store/model"
|
||||
)
|
||||
|
||||
const (
|
||||
defaultTunnelProbeTargetHost = "www.bing.com"
|
||||
defaultTunnelProbeTargetPort = 443
|
||||
)
|
||||
|
||||
type tunnelProbeTarget struct {
|
||||
Host string
|
||||
Port int
|
||||
}
|
||||
|
||||
func defaultTunnelProbeTarget() tunnelProbeTarget {
|
||||
return tunnelProbeTarget{Host: defaultTunnelProbeTargetHost, Port: defaultTunnelProbeTargetPort}
|
||||
}
|
||||
|
||||
func normalizeTunnelProbeTarget(host string, port int) (tunnelProbeTarget, bool, error) {
|
||||
host = strings.TrimSpace(host)
|
||||
if host == "" && port == 0 {
|
||||
return defaultTunnelProbeTarget(), false, nil
|
||||
}
|
||||
if host == "" {
|
||||
return tunnelProbeTarget{}, false, errors.New("测试目标 Host 不能为空")
|
||||
}
|
||||
if port <= 0 || port > 65535 {
|
||||
return tunnelProbeTarget{}, false, errors.New("测试目标端口必须是 1-65535")
|
||||
}
|
||||
if strings.Contains(host, "://") || strings.ContainsAny(host, "/?#") || strings.ContainsAny(host, " \t\r\n") || isTunnelProbeTargetSchemeLikeHost(host) {
|
||||
return tunnelProbeTarget{}, false, errors.New("测试目标 Host 不能包含协议或路径")
|
||||
}
|
||||
if normalized, ok := normalizeTunnelProbeTargetHost(host); ok {
|
||||
host = normalized
|
||||
} else {
|
||||
return tunnelProbeTarget{}, false, errors.New("测试目标 Host 格式无效")
|
||||
}
|
||||
|
||||
return tunnelProbeTarget{Host: host, Port: port}, true, nil
|
||||
}
|
||||
|
||||
func normalizeTunnelProbeTargetHost(host string) (string, bool) {
|
||||
if strings.HasPrefix(host, "[") || strings.HasSuffix(host, "]") {
|
||||
if !strings.HasPrefix(host, "[") || !strings.HasSuffix(host, "]") {
|
||||
return "", false
|
||||
}
|
||||
inner := strings.TrimPrefix(strings.TrimSuffix(host, "]"), "[")
|
||||
addr, err := netip.ParseAddr(inner)
|
||||
if err != nil || !addr.Is6() {
|
||||
return "", false
|
||||
}
|
||||
return inner, true
|
||||
}
|
||||
|
||||
if addr, err := netip.ParseAddr(host); err == nil {
|
||||
return addr.String(), true
|
||||
}
|
||||
if strings.Contains(host, ":") || isTunnelProbeTargetIPv4Like(host) {
|
||||
return "", false
|
||||
}
|
||||
if !isValidTunnelProbeTargetHost(host) {
|
||||
return "", false
|
||||
}
|
||||
return host, true
|
||||
}
|
||||
|
||||
func isValidTunnelProbeTargetHost(host string) bool {
|
||||
if host == "" || len(host) > 253 {
|
||||
return false
|
||||
}
|
||||
for _, label := range strings.Split(host, ".") {
|
||||
if len(label) == 0 || len(label) > 63 || label[0] == '-' || label[len(label)-1] == '-' {
|
||||
return false
|
||||
}
|
||||
for _, r := range label {
|
||||
if !isASCIILetter(r) && !isASCIIDigit(r) && r != '-' {
|
||||
return false
|
||||
}
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func isTunnelProbeTargetIPv4Like(host string) bool {
|
||||
if host == "" {
|
||||
return false
|
||||
}
|
||||
for _, r := range host {
|
||||
if !isASCIIDigit(r) && r != '.' {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return strings.Contains(host, ".")
|
||||
}
|
||||
|
||||
func isTunnelProbeTargetSchemeLikeHost(host string) bool {
|
||||
if _, err := netip.ParseAddr(host); err == nil {
|
||||
return false
|
||||
}
|
||||
|
||||
colon := strings.IndexByte(host, ':')
|
||||
if colon <= 0 {
|
||||
return false
|
||||
}
|
||||
for i, r := range host[:colon] {
|
||||
if i == 0 {
|
||||
if !isASCIILetter(r) {
|
||||
return false
|
||||
}
|
||||
continue
|
||||
}
|
||||
if !isASCIILetter(r) && !isASCIIDigit(r) && r != '+' && r != '-' && r != '.' {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func isASCIILetter(r rune) bool {
|
||||
return (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z')
|
||||
}
|
||||
|
||||
func isASCIIDigit(r rune) bool {
|
||||
return r >= '0' && r <= '9'
|
||||
}
|
||||
|
||||
func parseTunnelProbeTargetFromRequest(req map[string]interface{}) (tunnelProbeTarget, bool, error) {
|
||||
if req == nil {
|
||||
return defaultTunnelProbeTarget(), false, nil
|
||||
}
|
||||
rawHost, hasHost := req["probeTargetHost"]
|
||||
rawPort, hasPort := req["probeTargetPort"]
|
||||
if !hasHost && !hasPort {
|
||||
return defaultTunnelProbeTarget(), false, nil
|
||||
}
|
||||
host, err := parseTunnelProbeTargetHostValue(rawHost)
|
||||
if err != nil {
|
||||
return tunnelProbeTarget{}, false, err
|
||||
}
|
||||
port, err := parseTunnelProbeTargetPortValue(rawPort)
|
||||
if err != nil {
|
||||
return tunnelProbeTarget{}, false, err
|
||||
}
|
||||
return normalizeTunnelProbeTarget(host, port)
|
||||
}
|
||||
|
||||
func parseTunnelProbeTargetHostValue(raw interface{}) (string, error) {
|
||||
if raw == nil {
|
||||
return "", nil
|
||||
}
|
||||
host, ok := raw.(string)
|
||||
if !ok {
|
||||
return "", errors.New("测试目标 Host 格式无效")
|
||||
}
|
||||
if host != strings.TrimSpace(host) {
|
||||
return "", errors.New("测试目标 Host 不能包含协议或路径")
|
||||
}
|
||||
return host, nil
|
||||
}
|
||||
|
||||
func parseTunnelProbeTargetPortValue(raw interface{}) (int, error) {
|
||||
if raw == nil {
|
||||
return 0, nil
|
||||
}
|
||||
switch v := raw.(type) {
|
||||
case float64:
|
||||
if v != float64(int64(v)) {
|
||||
return 0, errors.New("测试目标端口必须是整数")
|
||||
}
|
||||
return int(v), nil
|
||||
case string:
|
||||
if v == "" {
|
||||
return 0, nil
|
||||
}
|
||||
if v != strings.TrimSpace(v) {
|
||||
return 0, errors.New("测试目标端口必须是整数")
|
||||
}
|
||||
port, err := strconv.Atoi(v)
|
||||
if err != nil {
|
||||
return 0, errors.New("测试目标端口必须是整数")
|
||||
}
|
||||
return port, nil
|
||||
case int:
|
||||
return v, nil
|
||||
case int32:
|
||||
return int(v), nil
|
||||
case int64:
|
||||
return int(v), nil
|
||||
default:
|
||||
return 0, errors.New("测试目标端口必须是整数")
|
||||
}
|
||||
}
|
||||
|
||||
func effectiveTunnelProbeTarget(tunnel *model.Tunnel) tunnelProbeTarget {
|
||||
if tunnel == nil {
|
||||
return defaultTunnelProbeTarget()
|
||||
}
|
||||
return effectiveTunnelProbeTargetValues(tunnel.ProbeTargetHost, tunnel.ProbeTargetPort)
|
||||
}
|
||||
|
||||
func effectiveTunnelProbeTargetValues(host string, port int) tunnelProbeTarget {
|
||||
target, configured, err := normalizeTunnelProbeTarget(host, port)
|
||||
if err != nil || !configured {
|
||||
return defaultTunnelProbeTarget()
|
||||
}
|
||||
return target
|
||||
}
|
||||
|
||||
func formatTunnelProbeTarget(target tunnelProbeTarget) string {
|
||||
if addr, err := netip.ParseAddr(target.Host); err == nil && addr.Is6() {
|
||||
return fmt.Sprintf("[%s]:%d", target.Host, target.Port)
|
||||
}
|
||||
return fmt.Sprintf("%s:%d", target.Host, target.Port)
|
||||
}
|
||||
@@ -0,0 +1,305 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"go-backend/internal/store/repo"
|
||||
)
|
||||
|
||||
func TestTunnelCreatePersistsProbeTargetAndListReturnsConfiguredValue(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
body := bytes.NewReader([]byte(`{
|
||||
"name":"custom-target",
|
||||
"type":1,
|
||||
"flow":1,
|
||||
"trafficRatio":1,
|
||||
"status":1,
|
||||
"inNodeId":[{"nodeId":10,"protocol":"tls"}],
|
||||
"probeTargetHost":"speed.example.com",
|
||||
"probeTargetPort":8443
|
||||
}`))
|
||||
|
||||
res := httptest.NewRecorder()
|
||||
h.tunnelCreate(res, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/create", body))
|
||||
assertProbeTargetSuccess(t, res)
|
||||
|
||||
listRes := httptest.NewRecorder()
|
||||
h.tunnelList(listRes, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/list", nil))
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
Data []map[string]any `json:"data"`
|
||||
}
|
||||
decodeProbeTargetResponse(t, listRes, &payload)
|
||||
if payload.Code != 0 {
|
||||
t.Fatalf("expected success, got code %d", payload.Code)
|
||||
}
|
||||
item := payload.Data[0]
|
||||
if item["probeTargetHost"] != "speed.example.com" || item["probeTargetPort"] != float64(8443) {
|
||||
t.Fatalf("unexpected probe target in list response: %+v", item)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelUpdatePersistsDefaultProbeTargetAsEmpty(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 77, "existing", "old.example.com", 9443)
|
||||
body := bytes.NewReader([]byte(`{
|
||||
"id":77,
|
||||
"name":"existing",
|
||||
"type":1,
|
||||
"flow":1,
|
||||
"trafficRatio":1,
|
||||
"status":1,
|
||||
"inNodeId":[{"nodeId":10,"protocol":"tls"}],
|
||||
"probeTargetHost":"",
|
||||
"probeTargetPort":0
|
||||
}`))
|
||||
|
||||
res := httptest.NewRecorder()
|
||||
h.tunnelUpdate(res, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/update", body))
|
||||
assertProbeTargetSuccess(t, res)
|
||||
|
||||
items, err := h.repo.ListTunnels()
|
||||
if err != nil {
|
||||
t.Fatalf("list tunnels: %v", err)
|
||||
}
|
||||
item := findProbeTargetTunnelItem(t, items, 77)
|
||||
if item["probeTargetHost"] != "" || item["probeTargetPort"] != 0 {
|
||||
t.Fatalf("expected default target to round-trip as empty/0, got %+v", item)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelUpdateWithoutProbeTargetFieldsPreservesExistingTarget(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 79, "existing", "old.example.com", 9443)
|
||||
body := bytes.NewReader([]byte(`{
|
||||
"id":79,
|
||||
"name":"existing",
|
||||
"type":1,
|
||||
"flow":1,
|
||||
"trafficRatio":1,
|
||||
"status":1,
|
||||
"inNodeId":[{"nodeId":10,"protocol":"tls"}]
|
||||
}`))
|
||||
|
||||
res := httptest.NewRecorder()
|
||||
h.tunnelUpdate(res, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/update", body))
|
||||
assertProbeTargetSuccess(t, res)
|
||||
|
||||
items, err := h.repo.ListTunnels()
|
||||
if err != nil {
|
||||
t.Fatalf("list tunnels: %v", err)
|
||||
}
|
||||
item := findProbeTargetTunnelItem(t, items, 79)
|
||||
if item["probeTargetHost"] != "old.example.com" || item["probeTargetPort"] != 9443 {
|
||||
t.Fatalf("expected omitted probe target fields to preserve existing target, got %+v", item)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelUpdateRejectsInvalidProbeTargetWithoutClearingExistingTarget(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
probeFields string
|
||||
}{
|
||||
{name: "non numeric port", probeFields: `,"probeTargetPort":"abc"`},
|
||||
{name: "fractional port", probeFields: `,"probeTargetPort":443.5`},
|
||||
{name: "whitespace host", probeFields: `,"probeTargetHost":" "`},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 80, "existing", "old.example.com", 9443)
|
||||
body := bytes.NewReader([]byte(`{
|
||||
"id":80,
|
||||
"name":"existing",
|
||||
"type":1,
|
||||
"flow":1,
|
||||
"trafficRatio":1,
|
||||
"status":1,
|
||||
"inNodeId":[{"nodeId":10,"protocol":"tls"}]
|
||||
` + tt.probeFields + `}`))
|
||||
|
||||
res := httptest.NewRecorder()
|
||||
h.tunnelUpdate(res, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/update", body))
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
Msg string `json:"msg"`
|
||||
}
|
||||
decodeProbeTargetResponse(t, res, &payload)
|
||||
if payload.Code == 0 || payload.Msg == "" {
|
||||
t.Fatalf("expected validation failure, got %+v", payload)
|
||||
}
|
||||
|
||||
items, err := h.repo.ListTunnels()
|
||||
if err != nil {
|
||||
t.Fatalf("list tunnels: %v", err)
|
||||
}
|
||||
item := findProbeTargetTunnelItem(t, items, 80)
|
||||
if item["probeTargetHost"] != "old.example.com" || item["probeTargetPort"] != 9443 {
|
||||
t.Fatalf("expected invalid probe target to preserve existing target, got %+v", item)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelCreateRejectsInvalidProbeTarget(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
body := bytes.NewReader([]byte(`{
|
||||
"name":"bad-target",
|
||||
"type":1,
|
||||
"flow":1,
|
||||
"trafficRatio":1,
|
||||
"status":1,
|
||||
"inNodeId":[{"nodeId":10,"protocol":"tls"}],
|
||||
"probeTargetHost":"https://example.com",
|
||||
"probeTargetPort":443
|
||||
}`))
|
||||
|
||||
res := httptest.NewRecorder()
|
||||
h.tunnelCreate(res, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/create", body))
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
Msg string `json:"msg"`
|
||||
}
|
||||
decodeProbeTargetResponse(t, res, &payload)
|
||||
if payload.Code == 0 || payload.Msg == "" {
|
||||
t.Fatalf("expected validation failure, got %+v", payload)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelUpdateInvalidProbeTargetDoesNotCleanFederationBindings(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 88, "existing", "old.example.com", 9443)
|
||||
seedProbeTargetFederationBinding(t, h, 88)
|
||||
body := bytes.NewReader([]byte(`{
|
||||
"id":88,
|
||||
"name":"existing",
|
||||
"type":1,
|
||||
"flow":1,
|
||||
"trafficRatio":1,
|
||||
"status":1,
|
||||
"inNodeId":[{"nodeId":10,"protocol":"tls"}],
|
||||
"probeTargetHost":"https://example.com",
|
||||
"probeTargetPort":443
|
||||
}`))
|
||||
|
||||
res := httptest.NewRecorder()
|
||||
h.tunnelUpdate(res, httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/update", body))
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
Msg string `json:"msg"`
|
||||
}
|
||||
decodeProbeTargetResponse(t, res, &payload)
|
||||
if payload.Code == 0 || payload.Msg == "" {
|
||||
t.Fatalf("expected validation failure, got %+v", payload)
|
||||
}
|
||||
|
||||
bindings, err := h.repo.ListActiveFederationTunnelBindingsByTunnel(88)
|
||||
if err != nil {
|
||||
t.Fatalf("list federation bindings: %v", err)
|
||||
}
|
||||
if len(bindings) != 1 {
|
||||
t.Fatalf("expected federation binding to remain after invalid update, got %d", len(bindings))
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelDiagnosisUsesConfiguredProbeTarget(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 90, "diagnosis-target", "speed.example.com", 8443)
|
||||
|
||||
_, _, workItems, err := h.prepareTunnelDiagnosis(90)
|
||||
if err != nil {
|
||||
t.Fatalf("prepare tunnel diagnosis: %v", err)
|
||||
}
|
||||
if len(workItems) != 1 {
|
||||
t.Fatalf("expected one diagnosis item, got %d", len(workItems))
|
||||
}
|
||||
if workItems[0].targetIP != "speed.example.com" || workItems[0].targetPort != 8443 {
|
||||
t.Fatalf("expected custom diagnosis target speed.example.com:8443, got %s:%d", workItems[0].targetIP, workItems[0].targetPort)
|
||||
}
|
||||
}
|
||||
|
||||
func setupProbeTargetTunnelHandler(t *testing.T) *Handler {
|
||||
t.Helper()
|
||||
r, err := repo.Open(filepath.Join(t.TempDir(), "panel.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = r.Close() })
|
||||
h := New(r, "secret")
|
||||
now := time.Now().UnixMilli()
|
||||
if err := r.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(10, 'entry-a', 'entry-secret', '10.0.0.1', '10.0.0.1', '', '30000-30010', '', 'v1', 1, 1, 1, ?, ?, 1, '[::]', '[::]', 0)
|
||||
`, now, now).Error; err != nil {
|
||||
t.Fatalf("insert node: %v", err)
|
||||
}
|
||||
return h
|
||||
}
|
||||
|
||||
func seedProbeTargetTunnel(t *testing.T, h *Handler, id int64, name string, host string, port 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, 1, 'tls', 1, ?, ?, 1, ?, '', ?, ?)
|
||||
`, id, name, now, now, id, host, port).Error; err != nil {
|
||||
t.Fatalf("insert 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, 'round', 1, 'tls')
|
||||
`, id).Error; err != nil {
|
||||
t.Fatalf("insert chain: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func seedProbeTargetFederationBinding(t *testing.T, h *Handler, tunnelID int64) {
|
||||
t.Helper()
|
||||
now := time.Now().UnixMilli()
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO federation_tunnel_binding(tunnel_id, node_id, chain_type, hop_inx, remote_url, resource_key, remote_binding_id, allocated_port, status, created_time, updated_time)
|
||||
VALUES(?, 10, 1, 0, 'http://peer.example', ?, 'remote-binding', 30001, 1, ?, ?)
|
||||
`, tunnelID, "probe-target-test-binding", now, now).Error; err != nil {
|
||||
t.Fatalf("insert federation binding: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func assertProbeTargetSuccess(t *testing.T, res *httptest.ResponseRecorder) {
|
||||
t.Helper()
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
Msg string `json:"msg"`
|
||||
}
|
||||
decodeProbeTargetResponse(t, res, &payload)
|
||||
if payload.Code != 0 {
|
||||
t.Fatalf("expected success, got %+v", payload)
|
||||
}
|
||||
}
|
||||
|
||||
func decodeProbeTargetResponse(t *testing.T, res *httptest.ResponseRecorder, v any) {
|
||||
t.Helper()
|
||||
if res.Code != http.StatusOK {
|
||||
t.Fatalf("expected HTTP %d, got %d", http.StatusOK, res.Code)
|
||||
}
|
||||
if err := json.NewDecoder(res.Body).Decode(v); err != nil {
|
||||
t.Fatalf("decode response: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func findProbeTargetTunnelItem(t *testing.T, items []map[string]interface{}, id int64) map[string]interface{} {
|
||||
t.Helper()
|
||||
for _, item := range items {
|
||||
if asInt64(item["id"], 0) == id {
|
||||
return item
|
||||
}
|
||||
}
|
||||
t.Fatalf("tunnel %d not found: %+v", id, items)
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package handler
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestNormalizeTunnelProbeTargetDefaultsWhenEmpty(t *testing.T) {
|
||||
target, configured, err := normalizeTunnelProbeTarget("", 0)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if configured {
|
||||
t.Fatalf("expected empty input to be default, not configured")
|
||||
}
|
||||
if target.Host != defaultTunnelProbeTargetHost || target.Port != defaultTunnelProbeTargetPort {
|
||||
t.Fatalf("unexpected default target: %+v", target)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelProbeTargetAcceptsHostPortAndIPv6(t *testing.T) {
|
||||
target, configured, err := normalizeTunnelProbeTarget(" [2001:db8::1] ", 8443)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if !configured {
|
||||
t.Fatalf("expected explicit target")
|
||||
}
|
||||
if target.Host != "2001:db8::1" || target.Port != 8443 {
|
||||
t.Fatalf("unexpected normalized target: %+v", target)
|
||||
}
|
||||
if got := formatTunnelProbeTarget(target); got != "[2001:db8::1]:8443" {
|
||||
t.Fatalf("unexpected formatted target: %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelProbeTargetRejectsPartialAndInvalidInputs(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
host string
|
||||
port int
|
||||
}{
|
||||
{name: "missing host", host: "", port: 443},
|
||||
{name: "missing port", host: "example.com", port: 0},
|
||||
{name: "port too high", host: "example.com", port: 70000},
|
||||
{name: "scheme", host: "https://example.com", port: 443},
|
||||
{name: "path", host: "example.com/ping", port: 443},
|
||||
{name: "space", host: "example .com", port: 443},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if _, _, err := normalizeTunnelProbeTarget(tt.host, tt.port); err == nil {
|
||||
t.Fatalf("expected validation error")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelProbeTargetRejectsSchemePrefixButAllowsIPv6(t *testing.T) {
|
||||
for _, host := range []string{"https:example.com", "mailto:ops@example.com"} {
|
||||
if _, _, err := normalizeTunnelProbeTarget(host, 443); err == nil {
|
||||
t.Fatalf("expected scheme-like host %q to be rejected", host)
|
||||
}
|
||||
}
|
||||
|
||||
for _, host := range []string{"2001:db8::1", "[2001:db8::1]"} {
|
||||
target, configured, err := normalizeTunnelProbeTarget(host, 443)
|
||||
if err != nil {
|
||||
t.Fatalf("expected IPv6 host %q to be accepted: %v", host, err)
|
||||
}
|
||||
if !configured || target.Host != "2001:db8::1" {
|
||||
t.Fatalf("unexpected IPv6 normalization for %q: %+v configured=%v", host, target, configured)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelProbeTargetValidatesHostShape(t *testing.T) {
|
||||
validHosts := []string{
|
||||
"example.com",
|
||||
"localhost",
|
||||
"api-1.example.co.uk",
|
||||
"192.0.2.10",
|
||||
"2001:db8::1",
|
||||
"[2001:db8::1]",
|
||||
}
|
||||
for _, host := range validHosts {
|
||||
if _, _, err := normalizeTunnelProbeTarget(host, 443); err != nil {
|
||||
t.Fatalf("expected valid host %q: %v", host, err)
|
||||
}
|
||||
}
|
||||
|
||||
invalidHosts := []string{
|
||||
"1:2:3",
|
||||
"[2001:db8::1",
|
||||
"2001:db8::1]",
|
||||
"[example.com]",
|
||||
"example..com",
|
||||
"-example.com",
|
||||
"example-.com",
|
||||
"exa_mple.com",
|
||||
"999.1.1.1",
|
||||
}
|
||||
for _, host := range invalidHosts {
|
||||
if _, _, err := normalizeTunnelProbeTarget(host, 443); err == nil {
|
||||
t.Fatalf("expected invalid host %q to be rejected", host)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseTunnelProbeTargetFromRequest(t *testing.T) {
|
||||
req := map[string]interface{}{
|
||||
"probeTargetHost": "speed.example.com",
|
||||
"probeTargetPort": float64(1443),
|
||||
}
|
||||
target, configured, err := parseTunnelProbeTargetFromRequest(req)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if !configured || target.Host != "speed.example.com" || target.Port != 1443 {
|
||||
t.Fatalf("unexpected request target: %+v configured=%v", target, configured)
|
||||
}
|
||||
}
|
||||
@@ -42,6 +42,8 @@ type tunnelQualitySnapshot struct {
|
||||
ErrorMessage string `json:"errorMessage,omitempty"`
|
||||
Timestamp int64 `json:"timestamp"`
|
||||
ChainDetails string `json:"chainDetails,omitempty"`
|
||||
ProbeTargetHost string `json:"probeTargetHost,omitempty"`
|
||||
ProbeTargetPort int `json:"probeTargetPort,omitempty"`
|
||||
|
||||
// internal fields for db reporting
|
||||
lastDBWrite int64 `json:"-"`
|
||||
@@ -57,6 +59,7 @@ type tunnelQualityProber struct {
|
||||
interval time.Duration
|
||||
lastPrune int64
|
||||
probing int32 // atomic flag: 1 = probeAll running, 0 = idle
|
||||
probeNode bestExitProbeFunc
|
||||
}
|
||||
|
||||
// newTunnelQualityProber creates a new prober (not yet running).
|
||||
@@ -226,6 +229,9 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
p.storeResult(snap)
|
||||
return
|
||||
}
|
||||
probeTarget := effectiveTunnelProbeTargetValues(tunnel.ProbeTargetHost, tunnel.ProbeTargetPort)
|
||||
snap.ProbeTargetHost = probeTarget.Host
|
||||
snap.ProbeTargetPort = probeTarget.Port
|
||||
|
||||
chainRows, err := h.listChainNodesForTunnel(tunnelID)
|
||||
if err != nil || len(chainRows) == 0 {
|
||||
@@ -242,13 +248,13 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
pingTimeoutMS: tunnelQualityPingTimeoutMs,
|
||||
timeoutMessage: "探测超时",
|
||||
}
|
||||
p.probeBestExitOwners(tunnelID, inNodes, midNodesGrouped, outNodes, ipPreference, options)
|
||||
p.probeBestExitOwners(tunnelID, inNodes, midNodesGrouped, outNodes, ipPreference, options, probeTarget)
|
||||
|
||||
switch tunnel.Type {
|
||||
case 1:
|
||||
// Port forwarding: entry → Bing only
|
||||
// Port forwarding: entry → public probe target only.
|
||||
if len(inNodes) > 0 {
|
||||
lat, loss, err := p.tcpPingNode(inNodes[0].NodeID, "www.bing.com", 443, options)
|
||||
lat, loss, err := p.pingNode(inNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -310,7 +316,7 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
hop.TargetIP = targetIP
|
||||
hop.TargetPort = targetPort
|
||||
|
||||
lat, loss, err := p.tcpPingNode(source.NodeID, targetIP, targetPort, options)
|
||||
lat, loss, err := p.pingNode(source.NodeID, targetIP, targetPort, options)
|
||||
if err == nil {
|
||||
hop.Latency = lat
|
||||
hop.Loss = loss
|
||||
@@ -346,7 +352,7 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
|
||||
// Exit → Bing
|
||||
if len(outNodes) > 0 {
|
||||
lat, loss, err := p.tcpPingNode(outNodes[0].NodeID, "www.bing.com", 443, options)
|
||||
lat, loss, err := p.pingNode(outNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -360,9 +366,9 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
|
||||
snap.Success = probeOK
|
||||
default:
|
||||
// Unknown type: entry → Bing
|
||||
// Unknown type: entry → public probe target.
|
||||
if len(inNodes) > 0 {
|
||||
lat, loss, err := p.tcpPingNode(inNodes[0].NodeID, "www.bing.com", 443, options)
|
||||
lat, loss, err := p.pingNode(inNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -376,7 +382,7 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
p.storeResult(snap)
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chainNodeRecord, chainHops [][]chainNodeRecord, outNodes []chainNodeRecord, ipPreference string, options diagnosisExecOptions) {
|
||||
func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chainNodeRecord, chainHops [][]chainNodeRecord, outNodes []chainNodeRecord, ipPreference string, options diagnosisExecOptions, probeTarget tunnelProbeTarget) {
|
||||
if p == nil || p.handler == nil || p.handler.bestExit == nil || len(outNodes) <= 1 {
|
||||
return
|
||||
}
|
||||
@@ -400,14 +406,14 @@ func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chai
|
||||
}
|
||||
// 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.tcpPingNode)
|
||||
roundPinger := newBestExitRoundPinger(p.pingNode)
|
||||
for _, owner := range owners {
|
||||
if nodeMap[owner.NodeID] == nil {
|
||||
continue
|
||||
}
|
||||
key := bestExitOwnerKey{TunnelID: tunnelID, OwnerNodeID: owner.NodeID}
|
||||
p.handler.bestExit.ensureApplied(key, outNodes[0].NodeID, time.Now())
|
||||
scores := evaluateBestExitOwner(owner, outNodes, nodeMap, ipPreference, options, roundPinger)
|
||||
scores := evaluateBestExitOwner(owner, outNodes, nodeMap, ipPreference, options, probeTarget, roundPinger)
|
||||
decision := p.handler.bestExit.observeScores(key, scores, time.Now())
|
||||
if decision.Switch {
|
||||
now := time.Now()
|
||||
@@ -421,6 +427,13 @@ func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chai
|
||||
}
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) pingNode(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
if p != nil && p.probeNode != nil {
|
||||
return p.probeNode(nodeID, ip, port, options)
|
||||
}
|
||||
return p.tcpPingNode(nodeID, ip, port, options)
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) tcpPingNode(nodeID int64, ip string, port int, options diagnosisExecOptions) (latency float64, loss float64, err error) {
|
||||
h := p.handler
|
||||
if h == nil {
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"slices"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestTunnelQualityProberUsesConfiguredProbeTarget(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 77, "quality-target", "speed.example.com", 8443)
|
||||
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(30, 'exit-a', 'exit-secret', '10.0.0.30', '10.0.0.30', '', '30000-30010', '', 'v1', 1, 1, 1, ?, ?, 1, '[::]', '[::]', 0)
|
||||
`, time.Now().UnixMilli(), time.Now().UnixMilli()).Error; err != nil {
|
||||
t.Fatalf("insert exit node: %v", err)
|
||||
}
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol)
|
||||
VALUES(77, '3', 30, 30001, 'round', 1, 'tls')
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert exit chain: %v", err)
|
||||
}
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
var calls []string
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
calls = append(calls, fmt.Sprintf("%d|%s|%d", nodeID, ip, port))
|
||||
return 10, 0, nil
|
||||
}
|
||||
p.probeTunnel(77)
|
||||
|
||||
if !slices.Contains(calls, "10|speed.example.com|8443") {
|
||||
t.Fatalf("expected type 1 public probe from entry to configured target, calls=%+v", calls)
|
||||
}
|
||||
if slices.Contains(calls, "30|speed.example.com|8443") {
|
||||
t.Fatalf("did not expect type 1 public probe from exit node, calls=%+v", calls)
|
||||
}
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 {
|
||||
t.Fatalf("expected one quality snapshot, got %+v", snaps)
|
||||
}
|
||||
if snaps[0].ProbeTargetHost != "speed.example.com" || snaps[0].ProbeTargetPort != 8443 {
|
||||
t.Fatalf("unexpected snapshot target metadata: %+v", snaps[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberStoresProbeTargetWhenChainIncomplete(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 78, "quality-target-incomplete", "speed.example.com", 8443)
|
||||
if err := h.repo.DB().Exec(`DELETE FROM chain_tunnel WHERE tunnel_id = ?`, 78).Error; err != nil {
|
||||
t.Fatalf("delete chain rows: %v", err)
|
||||
}
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
p.probeTunnel(78)
|
||||
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 {
|
||||
t.Fatalf("expected one quality snapshot, got %+v", snaps)
|
||||
}
|
||||
if snaps[0].ErrorMessage == "" {
|
||||
t.Fatalf("expected incomplete chain error, got %+v", snaps[0])
|
||||
}
|
||||
if snaps[0].ProbeTargetHost != "speed.example.com" || snaps[0].ProbeTargetPort != 8443 {
|
||||
t.Fatalf("unexpected snapshot target metadata: %+v", snaps[0])
|
||||
}
|
||||
}
|
||||
@@ -119,18 +119,20 @@ type StatisticsFlow struct {
|
||||
func (StatisticsFlow) TableName() string { return "statistics_flow" }
|
||||
|
||||
type Tunnel struct {
|
||||
ID int64 `gorm:"primaryKey;autoIncrement"`
|
||||
Name string `gorm:"type:varchar(100);not null"`
|
||||
TrafficRatio float64 `gorm:"column:traffic_ratio;not null;default:1.0"`
|
||||
Type int `gorm:"not null"`
|
||||
Protocol string `gorm:"type:varchar(10);not null;default:'tls'"`
|
||||
Flow int64 `gorm:"not null"`
|
||||
CreatedTime int64 `gorm:"column:created_time;not null"`
|
||||
UpdatedTime int64 `gorm:"column:updated_time;not null"`
|
||||
Status int `gorm:"not null"`
|
||||
InIP sql.NullString `gorm:"column:in_ip;type:text"`
|
||||
Inx int `gorm:"not null;default:0"`
|
||||
IPPreference string `gorm:"column:ip_preference;type:varchar(10);not null;default:''"`
|
||||
ID int64 `gorm:"primaryKey;autoIncrement"`
|
||||
Name string `gorm:"type:varchar(100);not null"`
|
||||
TrafficRatio float64 `gorm:"column:traffic_ratio;not null;default:1.0"`
|
||||
Type int `gorm:"not null"`
|
||||
Protocol string `gorm:"type:varchar(10);not null;default:'tls'"`
|
||||
Flow int64 `gorm:"not null"`
|
||||
CreatedTime int64 `gorm:"column:created_time;not null"`
|
||||
UpdatedTime int64 `gorm:"column:updated_time;not null"`
|
||||
Status int `gorm:"not null"`
|
||||
InIP sql.NullString `gorm:"column:in_ip;type:text"`
|
||||
Inx int `gorm:"not null;default:0"`
|
||||
IPPreference string `gorm:"column:ip_preference;type:varchar(10);not null;default:''"`
|
||||
ProbeTargetHost string `gorm:"column:probe_target_host;type:text;not null;default:''"`
|
||||
ProbeTargetPort int `gorm:"column:probe_target_port;not null;default:0"`
|
||||
}
|
||||
|
||||
func (Tunnel) TableName() string { return "tunnel" }
|
||||
@@ -403,19 +405,21 @@ type NodeBackup struct {
|
||||
}
|
||||
|
||||
type TunnelBackup struct {
|
||||
ID int64 `json:"id"`
|
||||
Name string `json:"name"`
|
||||
TrafficRatio float64 `json:"trafficRatio"`
|
||||
Type int `json:"type"`
|
||||
Protocol string `json:"protocol"`
|
||||
Flow int64 `json:"flow"`
|
||||
CreatedTime int64 `json:"createdTime"`
|
||||
UpdatedTime int64 `json:"updatedTime"`
|
||||
Status int `json:"status"`
|
||||
InIP string `json:"inIp,omitempty"`
|
||||
Inx int `json:"inx"`
|
||||
IPPreference string `json:"ipPreference,omitempty"`
|
||||
ChainTunnels []ChainTunnelBackup `json:"chainTunnels,omitempty"`
|
||||
ID int64 `json:"id"`
|
||||
Name string `json:"name"`
|
||||
TrafficRatio float64 `json:"trafficRatio"`
|
||||
Type int `json:"type"`
|
||||
Protocol string `json:"protocol"`
|
||||
Flow int64 `json:"flow"`
|
||||
CreatedTime int64 `json:"createdTime"`
|
||||
UpdatedTime int64 `json:"updatedTime"`
|
||||
Status int `json:"status"`
|
||||
InIP string `json:"inIp,omitempty"`
|
||||
Inx int `json:"inx"`
|
||||
IPPreference string `json:"ipPreference,omitempty"`
|
||||
ProbeTargetHost string `json:"probeTargetHost,omitempty"`
|
||||
ProbeTargetPort int `json:"probeTargetPort,omitempty"`
|
||||
ChainTunnels []ChainTunnelBackup `json:"chainTunnels,omitempty"`
|
||||
}
|
||||
|
||||
type ChainTunnelBackup struct {
|
||||
@@ -553,12 +557,14 @@ type ForwardRecord struct {
|
||||
|
||||
// TunnelRecord is a minimal tunnel view used by control plane.
|
||||
type TunnelRecord struct {
|
||||
ID int64
|
||||
Type int
|
||||
Status int
|
||||
Flow int64
|
||||
TrafficRatio float64
|
||||
Protocol string
|
||||
ID int64
|
||||
Type int
|
||||
Status int
|
||||
Flow int64
|
||||
TrafficRatio float64
|
||||
Protocol string
|
||||
ProbeTargetHost string
|
||||
ProbeTargetPort int
|
||||
}
|
||||
|
||||
type UserQuotaView struct {
|
||||
|
||||
@@ -304,6 +304,7 @@ func autoMigrateAll(db *gorm.DB) error {
|
||||
m := db.Migrator()
|
||||
hasNode := m.HasTable(&model.Node{})
|
||||
hasTunnel := m.HasTable(&model.Tunnel{})
|
||||
hasForward := m.HasTable(&model.Forward{})
|
||||
|
||||
for _, item := range models {
|
||||
if hasNode {
|
||||
@@ -316,6 +317,11 @@ func autoMigrateAll(db *gorm.DB) error {
|
||||
continue
|
||||
}
|
||||
}
|
||||
if hasForward {
|
||||
if _, ok := item.(*model.Forward); ok {
|
||||
continue
|
||||
}
|
||||
}
|
||||
if err := db.AutoMigrate(item); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -385,7 +391,7 @@ func prepareSQLiteLegacyColumns(db *gorm.DB) error {
|
||||
}
|
||||
|
||||
if m.HasTable(&model.Tunnel{}) {
|
||||
for _, field := range []string{"Inx", "IPPreference"} {
|
||||
for _, field := range []string{"Inx", "IPPreference", "ProbeTargetHost", "ProbeTargetPort"} {
|
||||
if m.HasColumn(&model.Tunnel{}, field) {
|
||||
continue
|
||||
}
|
||||
@@ -396,7 +402,7 @@ func prepareSQLiteLegacyColumns(db *gorm.DB) error {
|
||||
}
|
||||
|
||||
if m.HasTable(&model.Forward{}) {
|
||||
for _, field := range []string{"ProxyProtocol"} {
|
||||
for _, field := range []string{"MaxConn", "IPMaxConn", "IPSpeedID", "ProxyProtocol"} {
|
||||
if m.HasColumn(&model.Forward{}, field) {
|
||||
continue
|
||||
}
|
||||
@@ -1141,11 +1147,13 @@ func (r *Repository) ListTunnels() ([]map[string]interface{}, error) {
|
||||
"id": t.ID, "inx": t.Inx, "name": t.Name,
|
||||
"type": t.Type, "flow": t.Flow, "trafficRatio": t.TrafficRatio,
|
||||
"status": t.Status, "createdTime": t.CreatedTime,
|
||||
"inIp": nullableString(t.InIP),
|
||||
"ipPreference": t.IPPreference,
|
||||
"inNodeId": make([]map[string]interface{}, 0),
|
||||
"outNodeId": make([]map[string]interface{}, 0),
|
||||
"chainNodes": make([][]map[string]interface{}, 0),
|
||||
"inIp": nullableString(t.InIP),
|
||||
"ipPreference": t.IPPreference,
|
||||
"probeTargetHost": t.ProbeTargetHost,
|
||||
"probeTargetPort": t.ProbeTargetPort,
|
||||
"inNodeId": make([]map[string]interface{}, 0),
|
||||
"outNodeId": make([]map[string]interface{}, 0),
|
||||
"chainNodes": make([][]map[string]interface{}, 0),
|
||||
}
|
||||
orderedIDs = append(orderedIDs, t.ID)
|
||||
}
|
||||
@@ -2061,6 +2069,7 @@ func (r *Repository) exportTunnels() ([]model.TunnelBackup, error) {
|
||||
Type: t.Type, Protocol: t.Protocol, Flow: t.Flow,
|
||||
CreatedTime: t.CreatedTime, UpdatedTime: t.UpdatedTime,
|
||||
Status: t.Status, Inx: t.Inx, IPPreference: t.IPPreference,
|
||||
ProbeTargetHost: t.ProbeTargetHost, ProbeTargetPort: t.ProbeTargetPort,
|
||||
}
|
||||
if t.InIP.Valid {
|
||||
b.InIP = t.InIP.String
|
||||
@@ -2456,23 +2465,25 @@ func importTunnels(tx *gorm.DB, tunnels []model.TunnelBackup, now int64) (int, e
|
||||
count := 0
|
||||
for _, t := range tunnels {
|
||||
item := model.Tunnel{
|
||||
ID: t.ID,
|
||||
Name: t.Name,
|
||||
TrafficRatio: t.TrafficRatio,
|
||||
Type: t.Type,
|
||||
Protocol: t.Protocol,
|
||||
Flow: t.Flow,
|
||||
CreatedTime: t.CreatedTime,
|
||||
UpdatedTime: now,
|
||||
Status: t.Status,
|
||||
InIP: sql.NullString{String: t.InIP, Valid: true},
|
||||
Inx: t.Inx,
|
||||
IPPreference: t.IPPreference,
|
||||
ID: t.ID,
|
||||
Name: t.Name,
|
||||
TrafficRatio: t.TrafficRatio,
|
||||
Type: t.Type,
|
||||
Protocol: t.Protocol,
|
||||
Flow: t.Flow,
|
||||
CreatedTime: t.CreatedTime,
|
||||
UpdatedTime: now,
|
||||
Status: t.Status,
|
||||
InIP: sql.NullString{String: t.InIP, Valid: true},
|
||||
Inx: t.Inx,
|
||||
IPPreference: t.IPPreference,
|
||||
ProbeTargetHost: t.ProbeTargetHost,
|
||||
ProbeTargetPort: t.ProbeTargetPort,
|
||||
}
|
||||
err := tx.Clauses(clause.OnConflict{
|
||||
Columns: []clause.Column{{Name: "id"}},
|
||||
DoUpdates: clause.AssignmentColumns([]string{
|
||||
"name", "traffic_ratio", "type", "protocol", "flow", "updated_time", "status", "in_ip", "inx", "ip_preference",
|
||||
"name", "traffic_ratio", "type", "protocol", "flow", "updated_time", "status", "in_ip", "inx", "ip_preference", "probe_target_host", "probe_target_port",
|
||||
}),
|
||||
}).Create(&item).Error
|
||||
if err != nil {
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
package repo
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestBackupRoundTripsTunnelProbeTarget(t *testing.T) {
|
||||
source, err := Open(filepath.Join(t.TempDir(), "source.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open source repo: %v", err)
|
||||
}
|
||||
defer source.Close()
|
||||
|
||||
now := time.Now().UnixMilli()
|
||||
if err := source.DB().Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx, probe_target_host, probe_target_port)
|
||||
VALUES(20, 'backup-target', 1, 2, 'tls', 1, ?, ?, 1, '', 1, 'speed.example.com', 8443)
|
||||
`, now, now).Error; err != nil {
|
||||
t.Fatalf("insert source tunnel: %v", err)
|
||||
}
|
||||
|
||||
backup, err := source.ExportAll()
|
||||
if err != nil {
|
||||
t.Fatalf("export backup: %v", err)
|
||||
}
|
||||
if len(backup.Tunnels) != 1 {
|
||||
t.Fatalf("expected one exported tunnel, got %d", len(backup.Tunnels))
|
||||
}
|
||||
if backup.Tunnels[0].ProbeTargetHost != "speed.example.com" || backup.Tunnels[0].ProbeTargetPort != 8443 {
|
||||
t.Fatalf("unexpected exported probe target: %+v", backup.Tunnels[0])
|
||||
}
|
||||
|
||||
dest, err := Open(filepath.Join(t.TempDir(), "dest.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open dest repo: %v", err)
|
||||
}
|
||||
defer dest.Close()
|
||||
|
||||
result, err := dest.Import(backup, []string{"tunnels"})
|
||||
if err != nil {
|
||||
t.Fatalf("import backup: %v", err)
|
||||
}
|
||||
if result.TunnelsImported != 1 {
|
||||
t.Fatalf("expected one imported tunnel, got %d", result.TunnelsImported)
|
||||
}
|
||||
|
||||
items, err := dest.ListTunnels()
|
||||
if err != nil {
|
||||
t.Fatalf("list imported tunnels: %v", err)
|
||||
}
|
||||
if len(items) != 1 {
|
||||
t.Fatalf("expected one imported tunnel item, got %d", len(items))
|
||||
}
|
||||
if items[0]["probeTargetHost"] != "speed.example.com" || items[0]["probeTargetPort"] != 8443 {
|
||||
t.Fatalf("unexpected imported probe target: %+v", items[0])
|
||||
}
|
||||
}
|
||||
@@ -253,12 +253,14 @@ func (r *Repository) GetTunnelRecord(tunnelID int64) (*model.TunnelRecord, error
|
||||
return nil, err
|
||||
}
|
||||
tr := model.TunnelRecord{
|
||||
ID: t.ID,
|
||||
Type: t.Type,
|
||||
Status: t.Status,
|
||||
Flow: t.Flow,
|
||||
TrafficRatio: t.TrafficRatio,
|
||||
Protocol: t.Protocol,
|
||||
ID: t.ID,
|
||||
Type: t.Type,
|
||||
Status: t.Status,
|
||||
Flow: t.Flow,
|
||||
TrafficRatio: t.TrafficRatio,
|
||||
Protocol: t.Protocol,
|
||||
ProbeTargetHost: t.ProbeTargetHost,
|
||||
ProbeTargetPort: t.ProbeTargetPort,
|
||||
}
|
||||
if tr.Flow <= 0 {
|
||||
tr.Flow = 1
|
||||
|
||||
@@ -111,6 +111,33 @@ func TestGetFlowUploadForwardMetasKeepsForwardsWhenTunnelRowMissing(t *testing.T
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetTunnelRecordIncludesProbeTarget(t *testing.T) {
|
||||
r, err := Open(filepath.Join(t.TempDir(), "tunnel-record-probe-target.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open repo: %v", err)
|
||||
}
|
||||
defer r.Close()
|
||||
|
||||
now := time.Now().UnixMilli()
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx, probe_target_host, probe_target_port)
|
||||
VALUES(1, 't1', 1, 2, 'tls', 1, ?, ?, 1, NULL, 0, 'speed.example.com', 8443)
|
||||
`, now, now).Error; err != nil {
|
||||
t.Fatalf("insert tunnel: %v", err)
|
||||
}
|
||||
|
||||
record, err := r.GetTunnelRecord(1)
|
||||
if err != nil {
|
||||
t.Fatalf("get tunnel record: %v", err)
|
||||
}
|
||||
if record == nil {
|
||||
t.Fatalf("expected tunnel record")
|
||||
}
|
||||
if record.ProbeTargetHost != "speed.example.com" || record.ProbeTargetPort != 8443 {
|
||||
t.Fatalf("unexpected probe target on record: %#v", record)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAddUserQuotaUsageBatchReturnsNormalizedViews(t *testing.T) {
|
||||
r, err := Open(filepath.Join(t.TempDir(), "quota-batch.db"))
|
||||
if err != nil {
|
||||
|
||||
@@ -3,6 +3,7 @@ package repo
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
@@ -58,6 +59,126 @@ func TestPrepareSQLiteLegacyColumnsAddsNodeMetadataColumns(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenBackfillsSQLiteLegacyTunnelProbeTargetColumns(t *testing.T) {
|
||||
dbPath := filepath.Join(t.TempDir(), "legacy.db")
|
||||
db, err := gorm.Open(gsqlite.Open(dbPath), &gorm.Config{
|
||||
Logger: logger.Default.LogMode(logger.Silent),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("open legacy sqlite: %v", err)
|
||||
}
|
||||
|
||||
if err := db.Exec(`
|
||||
CREATE TABLE tunnel (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name VARCHAR(100) NOT NULL,
|
||||
traffic_ratio REAL NOT NULL DEFAULT 1.0,
|
||||
type INTEGER NOT NULL,
|
||||
protocol VARCHAR(10) NOT NULL DEFAULT 'tls',
|
||||
flow INTEGER NOT NULL,
|
||||
created_time INTEGER NOT NULL,
|
||||
updated_time INTEGER NOT NULL,
|
||||
status INTEGER NOT NULL,
|
||||
in_ip TEXT
|
||||
)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("create legacy tunnel table: %v", err)
|
||||
}
|
||||
if err := db.Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip)
|
||||
VALUES(1, 'legacy-tunnel', 1, 1, 'tls', 1, 1, 1, 1, '')
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert legacy tunnel: %v", err)
|
||||
}
|
||||
if sqlDB, _ := db.DB(); sqlDB != nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
|
||||
r, err := Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("open migrated sqlite: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = r.Close() })
|
||||
|
||||
m := r.DB().Migrator()
|
||||
for _, field := range []string{"ProbeTargetHost", "ProbeTargetPort"} {
|
||||
if !m.HasColumn(&model.Tunnel{}, field) {
|
||||
t.Fatalf("expected tunnel.%s column to exist", field)
|
||||
}
|
||||
}
|
||||
|
||||
var host string
|
||||
var port int
|
||||
if err := r.DB().Raw(`SELECT probe_target_host, probe_target_port FROM tunnel WHERE id = 1`).Row().Scan(&host, &port); err != nil {
|
||||
t.Fatalf("query probe target defaults: %v", err)
|
||||
}
|
||||
if host != "" || port != 0 {
|
||||
t.Fatalf("expected default probe target empty/0, got %q/%d", host, port)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenBackfillsSQLiteLegacyForwardColumns(t *testing.T) {
|
||||
dbPath := filepath.Join(t.TempDir(), "legacy-forward.db")
|
||||
db, err := gorm.Open(gsqlite.Open(dbPath), &gorm.Config{
|
||||
Logger: logger.Default.LogMode(logger.Silent),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("open legacy sqlite: %v", err)
|
||||
}
|
||||
|
||||
if err := db.Exec(`
|
||||
CREATE TABLE forward (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
user_id INTEGER NOT NULL,
|
||||
user_name VARCHAR(100) NOT NULL,
|
||||
name VARCHAR(100) NOT NULL,
|
||||
tunnel_id INTEGER NOT NULL,
|
||||
remote_addr TEXT NOT NULL,
|
||||
strategy VARCHAR(100) NOT NULL DEFAULT 'fifo',
|
||||
in_flow INTEGER NOT NULL DEFAULT 0,
|
||||
out_flow INTEGER NOT NULL DEFAULT 0,
|
||||
created_time INTEGER NOT NULL,
|
||||
updated_time INTEGER NOT NULL,
|
||||
status INTEGER NOT NULL,
|
||||
inx INTEGER NOT NULL DEFAULT 0,
|
||||
speed_id INTEGER
|
||||
)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("create legacy forward table: %v", err)
|
||||
}
|
||||
if err := db.Exec(`
|
||||
INSERT INTO forward(id, user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx, speed_id)
|
||||
VALUES(1, 2, 'legacy-user', 'legacy-forward', 3, '127.0.0.1:9000', 'fifo', 0, 0, 1, 1, 1, 0, NULL)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert legacy forward: %v", err)
|
||||
}
|
||||
if sqlDB, _ := db.DB(); sqlDB != nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
|
||||
r, err := Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("open migrated sqlite: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = r.Close() })
|
||||
|
||||
m := r.DB().Migrator()
|
||||
for _, field := range []string{"MaxConn", "IPMaxConn", "IPSpeedID", "ProxyProtocol"} {
|
||||
if !m.HasColumn(&model.Forward{}, field) {
|
||||
t.Fatalf("expected forward.%s column to exist", field)
|
||||
}
|
||||
}
|
||||
|
||||
var maxConn, ipMaxConn, proxyProtocol int
|
||||
var ipSpeedID sql.NullInt64
|
||||
if err := r.DB().Raw(`SELECT max_conn, ip_max_conn, ip_speed_id, proxy_protocol FROM forward WHERE id = 1`).Row().Scan(&maxConn, &ipMaxConn, &ipSpeedID, &proxyProtocol); err != nil {
|
||||
t.Fatalf("query forward defaults: %v", err)
|
||||
}
|
||||
if maxConn != 0 || ipMaxConn != 0 || ipSpeedID.Valid || proxyProtocol != 0 {
|
||||
t.Fatalf("expected default forward columns 0/0/NULL/0, got max_conn=%d ip_max_conn=%d ip_speed_id=%+v proxy_protocol=%d", maxConn, ipMaxConn, ipSpeedID, proxyProtocol)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateSchemaRunsPostgresIDRepairEvenAtCurrentVersion(t *testing.T) {
|
||||
db, err := gorm.Open(gsqlite.Open(":memory:"), &gorm.Config{
|
||||
Logger: logger.Default.LogMode(logger.Silent),
|
||||
|
||||
@@ -397,22 +397,24 @@ func (r *Repository) UpdateTunnelOrder(tunnelID int64, inx int, now int64) {
|
||||
Updates(map[string]interface{}{"inx": inx, "updated_time": now}).Error
|
||||
}
|
||||
|
||||
func (r *Repository) UpdateTunnelTx(tx *gorm.DB, tunnelID int64, name string, typeVal int, flow int64, trafficRatio float64, status int, inIP, ipPreference string, protocol string, now int64) error {
|
||||
func (r *Repository) UpdateTunnelTx(tx *gorm.DB, tunnelID int64, name string, typeVal int, flow int64, trafficRatio float64, status int, inIP, ipPreference string, protocol string, probeTargetHost string, probeTargetPort int, now int64) error {
|
||||
if tx == nil {
|
||||
return errors.New("database unavailable")
|
||||
}
|
||||
return tx.Model(&model.Tunnel{}).
|
||||
Where("id = ?", tunnelID).
|
||||
Updates(map[string]interface{}{
|
||||
"name": name,
|
||||
"type": typeVal,
|
||||
"flow": flow,
|
||||
"traffic_ratio": trafficRatio,
|
||||
"status": status,
|
||||
"in_ip": nullStringFromInterface(inIP),
|
||||
"ip_preference": ipPreference,
|
||||
"protocol": protocol,
|
||||
"updated_time": now,
|
||||
"name": name,
|
||||
"type": typeVal,
|
||||
"flow": flow,
|
||||
"traffic_ratio": trafficRatio,
|
||||
"status": status,
|
||||
"in_ip": nullStringFromInterface(inIP),
|
||||
"ip_preference": ipPreference,
|
||||
"protocol": protocol,
|
||||
"probe_target_host": probeTargetHost,
|
||||
"probe_target_port": probeTargetPort,
|
||||
"updated_time": now,
|
||||
}).Error
|
||||
}
|
||||
|
||||
@@ -1326,20 +1328,22 @@ func (r *Repository) BatchUpdateForwardStatus(ids []int64, status int) (int, int
|
||||
return s, f
|
||||
}
|
||||
|
||||
func (r *Repository) CreateTunnelTx(tx *gorm.DB, name string, trafficRatio float64, typeVal int, flow int64, now int64, status int, inIP interface{}, inx int, ipPreference string) (int64, error) {
|
||||
func (r *Repository) CreateTunnelTx(tx *gorm.DB, name string, trafficRatio float64, typeVal int, flow int64, now int64, status int, inIP interface{}, inx int, ipPreference string, probeTargetHost string, probeTargetPort int) (int64, error) {
|
||||
inIPVal := nullStringFromInterface(inIP)
|
||||
tunnel := model.Tunnel{
|
||||
Name: name,
|
||||
TrafficRatio: trafficRatio,
|
||||
Type: typeVal,
|
||||
Protocol: "tls",
|
||||
Flow: flow,
|
||||
CreatedTime: now,
|
||||
UpdatedTime: now,
|
||||
Status: status,
|
||||
InIP: inIPVal,
|
||||
Inx: inx,
|
||||
IPPreference: ipPreference,
|
||||
Name: name,
|
||||
TrafficRatio: trafficRatio,
|
||||
Type: typeVal,
|
||||
Protocol: "tls",
|
||||
Flow: flow,
|
||||
CreatedTime: now,
|
||||
UpdatedTime: now,
|
||||
Status: status,
|
||||
InIP: inIPVal,
|
||||
Inx: inx,
|
||||
IPPreference: ipPreference,
|
||||
ProbeTargetHost: probeTargetHost,
|
||||
ProbeTargetPort: probeTargetPort,
|
||||
}
|
||||
if err := tx.Create(&tunnel).Error; err != nil {
|
||||
return 0, err
|
||||
|
||||
@@ -158,6 +158,9 @@ type WebSocketReporter struct {
|
||||
addr string // 保存服务器地址
|
||||
secret string // 保存密钥
|
||||
version string // 保存版本号
|
||||
http int
|
||||
tls int
|
||||
socks int
|
||||
preferredWSScheme string
|
||||
conn *websocket.Conn
|
||||
curBackoff time.Duration // 当前重连退避间隔
|
||||
@@ -296,9 +299,9 @@ func (w *WebSocketReporter) connect() error {
|
||||
Socks int `json:"socks"`
|
||||
}
|
||||
|
||||
var cfg LocalConfig
|
||||
cfg := LocalConfig{Http: w.http, Tls: w.tls, Socks: w.socks}
|
||||
if b, err := os.ReadFile("config.json"); err == nil {
|
||||
json.Unmarshal(b, &cfg)
|
||||
_ = json.Unmarshal(b, &cfg)
|
||||
}
|
||||
|
||||
candidates := buildWebSocketCandidates(w.addr, w.secret, w.version, cfg.Http, cfg.Tls, cfg.Socks, w.preferredWSScheme)
|
||||
@@ -1369,7 +1372,7 @@ func (w *WebSocketReporter) handleUpgradeAgent(data interface{}) error {
|
||||
// 执行重启脚本
|
||||
// 使用 systemd-run 在独立的 transient unit 中运行重启脚本,
|
||||
// 避免 systemctl stop 杀死 flux_agent cgroup 内所有进程(包括此脚本自身)导致 mv 未执行。
|
||||
script := fmt.Sprintf("sleep 1 && systemctl stop flux_agent && mv %s %s && systemctl start flux_agent", tmpPath, binaryPath)
|
||||
script := buildAgentRestartScript(tmpPath, binaryPath)
|
||||
cmd := exec.Command("systemd-run", "--quiet", "/bin/sh", "-c", script)
|
||||
if err := cmd.Start(); err != nil {
|
||||
os.Remove(tmpPath)
|
||||
@@ -1403,6 +1406,14 @@ func (w *WebSocketReporter) handleRollbackAgent(data interface{}) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func buildAgentRestartScript(tmpPath, binaryPath string) string {
|
||||
return fmt.Sprintf(
|
||||
"sleep 1 && systemctl stop flux_agent && legacy_service='' && for service_file in /etc/systemd/system/gost.service /lib/systemd/system/gost.service /usr/lib/systemd/system/gost.service; do if [ -f \"$service_file\" ] && grep -Fq \"WorkingDirectory=/etc/gost\" \"$service_file\" && (grep -Fq \"ExecStart=/etc/gost/gost\" \"$service_file\" || (grep -Fq \"ExecStart=/usr/local/bin/gost\" \"$service_file\" && [ -f /etc/gost/config.json ] && [ -f /etc/gost/gost.json ])); then legacy_service=\"$service_file\"; break; fi; done && if [ -n \"$legacy_service\" ]; then (systemctl stop gost 2>/dev/null || true) && (systemctl disable gost 2>/dev/null || true) && rm -f /usr/local/bin/gost /etc/gost/gost \"$legacy_service\" && (systemctl daemon-reload 2>/dev/null || true); fi && mv %s %s && systemctl start flux_agent",
|
||||
tmpPath,
|
||||
binaryPath,
|
||||
)
|
||||
}
|
||||
|
||||
// updateLocalConfigJSON 将 http/tls/socks 写入工作目录下的 config.json
|
||||
func updateLocalConfigJSON(httpVal int, tlsVal int, socksVal int) error {
|
||||
path := "config.json"
|
||||
@@ -1650,13 +1661,16 @@ func StartWebSocketReporterWithConfig(addr string, secret string, http int, tls
|
||||
candidates := buildWebSocketCandidates(addr, secret, version, http, tls, socks, "")
|
||||
fullURL := candidates[0]
|
||||
|
||||
fmt.Printf("🔗 WebSocket连接URL: %s\n", fullURL)
|
||||
fmt.Printf("🔗 WebSocket连接URL: %s\n", sanitizeWebSocketURL(fullURL))
|
||||
|
||||
reporter := NewWebSocketReporter(fullURL, secret)
|
||||
// 保存 addr, secret, version 供重连时使用
|
||||
// 保存 addr, secret, version 和协议能力供重连时使用
|
||||
reporter.addr = addr
|
||||
reporter.secret = secret
|
||||
reporter.version = version
|
||||
reporter.http = http
|
||||
reporter.tls = tls
|
||||
reporter.socks = socks
|
||||
reporter.Start()
|
||||
return reporter
|
||||
}
|
||||
|
||||
@@ -1,15 +1,45 @@
|
||||
package socket
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"os"
|
||||
"runtime"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
)
|
||||
|
||||
func captureStdout(t *testing.T, fn func()) string {
|
||||
t.Helper()
|
||||
|
||||
orig := os.Stdout
|
||||
r, w, err := os.Pipe()
|
||||
if err != nil {
|
||||
t.Fatalf("create stdout pipe: %v", err)
|
||||
}
|
||||
os.Stdout = w
|
||||
defer func() {
|
||||
os.Stdout = orig
|
||||
_ = w.Close()
|
||||
_ = r.Close()
|
||||
}()
|
||||
|
||||
fn()
|
||||
|
||||
_ = w.Close()
|
||||
|
||||
var buf bytes.Buffer
|
||||
if _, err := io.Copy(&buf, r); err != nil {
|
||||
t.Fatalf("read stdout: %v", err)
|
||||
}
|
||||
return buf.String()
|
||||
}
|
||||
|
||||
func TestBuildWebSocketCandidatesSecureFirst(t *testing.T) {
|
||||
candidates := buildWebSocketCandidates("panel.example.com:443", "abc", "2.0.2", 1, 0, 1, "")
|
||||
|
||||
@@ -133,3 +163,103 @@ func TestFormatWebSocketDialErrorIncludesHTTPStatus(t *testing.T) {
|
||||
t.Fatalf("expected response body in message, got %s", msg)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentUpgradeRestartScriptStopsLegacyGostService(t *testing.T) {
|
||||
script := buildAgentRestartScript("/tmp/flux_agent.new", "/etc/flux_agent/flux_agent")
|
||||
|
||||
if !strings.Contains(script, "systemctl stop flux_agent") {
|
||||
t.Fatalf("expected script to stop flux_agent, got %s", script)
|
||||
}
|
||||
if !strings.Contains(script, "mv /tmp/flux_agent.new /etc/flux_agent/flux_agent") {
|
||||
t.Fatalf("expected script to replace the flux_agent binary, got %s", script)
|
||||
}
|
||||
if !strings.Contains(script, "systemctl stop gost") {
|
||||
t.Fatalf("expected script to stop the legacy gost service, got %s", script)
|
||||
}
|
||||
if !strings.Contains(script, "systemctl disable gost") {
|
||||
t.Fatalf("expected script to disable the legacy gost service, got %s", script)
|
||||
}
|
||||
if !strings.Contains(script, "rm -f /usr/local/bin/gost") {
|
||||
t.Fatalf("expected script to remove the legacy gost binary, got %s", script)
|
||||
}
|
||||
if !strings.Contains(script, "WorkingDirectory=/etc/gost") {
|
||||
t.Fatalf("expected script to scope cleanup to the legacy FLVX gost service definition, got %s", script)
|
||||
}
|
||||
if !strings.Contains(script, "systemctl start flux_agent") {
|
||||
t.Fatalf("expected script to restart flux_agent, got %s", script)
|
||||
}
|
||||
if strings.Contains(script, "systemctl stop flux_agent && systemctl stop gost 2>/dev/null || true") {
|
||||
t.Fatalf("expected legacy gost cleanup fallback to be scoped, got %s", script)
|
||||
}
|
||||
if runtime.GOARCH == "" {
|
||||
t.Fatalf("unexpected empty runtime arch")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStartWebSocketReporterWithConfigPreservesProtocolDefaultsWithoutConfigFile(t *testing.T) {
|
||||
origDial := wsDial
|
||||
defer func() { wsDial = origDial }()
|
||||
|
||||
origWD, err := os.Getwd()
|
||||
if err != nil {
|
||||
t.Fatalf("get working directory: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
_ = os.Chdir(origWD)
|
||||
})
|
||||
if err := os.Chdir(t.TempDir()); err != nil {
|
||||
t.Fatalf("change working directory: %v", err)
|
||||
}
|
||||
|
||||
urls := make(chan string, 1)
|
||||
wsDial = func(_ *websocket.Dialer, rawURL string) (*websocket.Conn, *http.Response, error) {
|
||||
select {
|
||||
case urls <- rawURL:
|
||||
default:
|
||||
}
|
||||
return nil, nil, errors.New("dial failed")
|
||||
}
|
||||
|
||||
reporter := StartWebSocketReporterWithConfig("panel.example.com:443", "abc", 1, 0, 1, "2.0.2")
|
||||
defer reporter.Stop()
|
||||
|
||||
select {
|
||||
case rawURL := <-urls:
|
||||
if !strings.Contains(rawURL, "http=1&tls=0&socks=1") {
|
||||
t.Fatalf("expected reconnect URL to preserve startup protocol values, got %s", rawURL)
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("timed out waiting for websocket dial")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStartWebSocketReporterWithConfigLogsSanitizedURL(t *testing.T) {
|
||||
origDial := wsDial
|
||||
defer func() { wsDial = origDial }()
|
||||
|
||||
ready := make(chan struct{}, 1)
|
||||
wsDial = func(_ *websocket.Dialer, rawURL string) (*websocket.Conn, *http.Response, error) {
|
||||
select {
|
||||
case ready <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
return nil, nil, errors.New("dial failed")
|
||||
}
|
||||
|
||||
output := captureStdout(t, func() {
|
||||
reporter := StartWebSocketReporterWithConfig("panel.example.com:443", "abc123", 1, 0, 1, "2.0.2")
|
||||
select {
|
||||
case <-ready:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("timed out waiting for websocket dial")
|
||||
}
|
||||
reporter.Stop()
|
||||
})
|
||||
|
||||
if strings.Contains(output, "secret=abc123") {
|
||||
t.Fatalf("expected logged websocket URL to mask the node secret, got %s", output)
|
||||
}
|
||||
if !strings.Contains(output, "secret=%2A%2A%2A") {
|
||||
t.Fatalf("expected logged websocket URL to include masked secret, got %s", output)
|
||||
}
|
||||
}
|
||||
|
||||
+79
-11
@@ -24,6 +24,11 @@ get_architecture() {
|
||||
|
||||
# 安装目录
|
||||
INSTALL_DIR="/etc/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"
|
||||
|
||||
# 镜像加速配置(可由面板传入或交互式询问)
|
||||
PROXY_ENABLED="${PROXY_ENABLED:-}"
|
||||
@@ -234,6 +239,69 @@ check_and_install_tcpkill() {
|
||||
return 0
|
||||
}
|
||||
|
||||
json_escape() {
|
||||
local value="$1"
|
||||
value=${value//\\/\\\\}
|
||||
value=${value//\"/\\\"}
|
||||
value=${value//$'\n'/\\n}
|
||||
value=${value//$'\r'/\\r}
|
||||
value=${value//$'\t'/\\t}
|
||||
printf '%s' "$value"
|
||||
}
|
||||
|
||||
write_flux_agent_config() {
|
||||
local path="$1"
|
||||
printf '{\n "addr": "%s",\n "secret": "%s"\n}\n' \
|
||||
"$(json_escape "$SERVER_ADDR")" \
|
||||
"$(json_escape "$SECRET")" > "$path"
|
||||
}
|
||||
|
||||
cleanup_legacy_gost_installation() {
|
||||
local matched_service_files=()
|
||||
local service_file=""
|
||||
local removed_service_file="0"
|
||||
|
||||
for service_file in "$LEGACY_GOST_SERVICE_FILE_ETC" "$LEGACY_GOST_SERVICE_FILE_LIB" "$LEGACY_GOST_SERVICE_FILE_USR_LIB"; do
|
||||
if [[ ! -f "$service_file" ]]; then
|
||||
continue
|
||||
fi
|
||||
if ! grep -Fq "WorkingDirectory=$LEGACY_GOST_CONFIG_DIR" "$service_file"; then
|
||||
continue
|
||||
fi
|
||||
if grep -Fq "ExecStart=$LEGACY_GOST_CONFIG_DIR/gost" "$service_file" || \
|
||||
(grep -Fq "ExecStart=$LEGACY_GOST_BINARY" "$service_file" && [[ -f "$LEGACY_GOST_CONFIG_DIR/config.json" && -f "$LEGACY_GOST_CONFIG_DIR/gost.json" ]]); then
|
||||
matched_service_files+=("$service_file")
|
||||
fi
|
||||
done
|
||||
|
||||
if [[ ${#matched_service_files[@]} -eq 0 ]]; then
|
||||
return 0
|
||||
fi
|
||||
|
||||
if systemctl list-units --full -all 2>/dev/null | grep -Fq "gost.service"; then
|
||||
systemctl stop gost 2>/dev/null || true
|
||||
systemctl disable gost 2>/dev/null || true
|
||||
fi
|
||||
|
||||
for service_file in "${matched_service_files[@]}"; do
|
||||
if [[ -f "$service_file" ]]; then
|
||||
rm -f "$service_file"
|
||||
removed_service_file="1"
|
||||
fi
|
||||
done
|
||||
|
||||
if [[ -f "$LEGACY_GOST_BINARY" ]]; then
|
||||
rm -f "$LEGACY_GOST_BINARY"
|
||||
fi
|
||||
if [[ -f "$LEGACY_GOST_CONFIG_DIR/gost" ]]; then
|
||||
rm -f "$LEGACY_GOST_CONFIG_DIR/gost"
|
||||
fi
|
||||
|
||||
if [[ "$removed_service_file" == "1" ]]; then
|
||||
systemctl daemon-reload 2>/dev/null || true
|
||||
fi
|
||||
}
|
||||
|
||||
|
||||
# 获取用户输入的配置参数
|
||||
get_config_params() {
|
||||
@@ -279,6 +347,8 @@ install_flux_agent() {
|
||||
|
||||
mkdir -p "$INSTALL_DIR"
|
||||
|
||||
local tmp_binary="$INSTALL_DIR/flux_agent.new"
|
||||
|
||||
# 停止并禁用已有服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
echo "🔍 检测到已存在的flux_agent服务"
|
||||
@@ -286,16 +356,17 @@ install_flux_agent() {
|
||||
systemctl disable flux_agent 2>/dev/null && echo "🚫 禁用自启"
|
||||
fi
|
||||
|
||||
# 删除旧文件
|
||||
[[ -f "$INSTALL_DIR/flux_agent" ]] && echo "🧹 删除旧文件 flux_agent" && rm -f "$INSTALL_DIR/flux_agent"
|
||||
|
||||
# 下载 flux_agent
|
||||
echo "⬇️ 下载 flux_agent 中..."
|
||||
curl -L "$DOWNLOAD_URL" -o "$INSTALL_DIR/flux_agent"
|
||||
if [[ ! -f "$INSTALL_DIR/flux_agent" || ! -s "$INSTALL_DIR/flux_agent" ]]; then
|
||||
rm -f "$tmp_binary"
|
||||
curl -L "$DOWNLOAD_URL" -o "$tmp_binary"
|
||||
if [[ ! -f "$tmp_binary" || ! -s "$tmp_binary" ]]; then
|
||||
rm -f "$tmp_binary"
|
||||
echo "❌ 下载失败,请检查网络或下载链接。"
|
||||
exit 1
|
||||
fi
|
||||
cleanup_legacy_gost_installation
|
||||
mv "$tmp_binary" "$INSTALL_DIR/flux_agent"
|
||||
chmod +x "$INSTALL_DIR/flux_agent"
|
||||
echo "✅ 下载完成"
|
||||
|
||||
@@ -305,12 +376,7 @@ install_flux_agent() {
|
||||
# 写入 config.json (安装时总是创建新的)
|
||||
CONFIG_FILE="$INSTALL_DIR/config.json"
|
||||
echo "📄 创建新配置: config.json"
|
||||
cat > "$CONFIG_FILE" <<EOF
|
||||
{
|
||||
"addr": "$SERVER_ADDR",
|
||||
"secret": "$SECRET"
|
||||
}
|
||||
EOF
|
||||
write_flux_agent_config "$CONFIG_FILE"
|
||||
|
||||
# 写入 gost.json
|
||||
GOST_CONFIG="$INSTALL_DIR/gost.json"
|
||||
@@ -380,11 +446,13 @@ update_flux_agent() {
|
||||
|
||||
# 先下载新版本
|
||||
echo "⬇️ 下载最新版本..."
|
||||
rm -f "$INSTALL_DIR/flux_agent.new"
|
||||
curl -L "$DOWNLOAD_URL" -o "$INSTALL_DIR/flux_agent.new"
|
||||
if [[ ! -f "$INSTALL_DIR/flux_agent.new" || ! -s "$INSTALL_DIR/flux_agent.new" ]]; then
|
||||
echo "❌ 下载失败。"
|
||||
return 1
|
||||
fi
|
||||
cleanup_legacy_gost_installation
|
||||
|
||||
# 停止服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
|
||||
@@ -98,6 +98,7 @@ EOF
|
||||
chmod +x "$INSTALL_DIR/flux_agent"
|
||||
|
||||
local ask_called="0"
|
||||
local cleanup_called="0"
|
||||
|
||||
ask_proxy_config() {
|
||||
ask_called="1"
|
||||
@@ -105,6 +106,10 @@ EOF
|
||||
DOWNLOAD_URL=""
|
||||
}
|
||||
|
||||
cleanup_legacy_gost_installation() {
|
||||
cleanup_called="1"
|
||||
}
|
||||
|
||||
check_and_install_tcpkill() { :; }
|
||||
|
||||
systemctl() {
|
||||
@@ -133,6 +138,7 @@ EOF
|
||||
update_flux_agent >/dev/null
|
||||
|
||||
assert_equals "1" "$ask_called" "update_flux_agent should ask for proxy config before downloading"
|
||||
assert_equals "1" "$cleanup_called" "update_flux_agent should clean up legacy gost before restarting the agent"
|
||||
assert_equals "$(build_download_url)" "$DOWNLOAD_URL" "update_flux_agent should honor the prompted proxy choice"
|
||||
)
|
||||
|
||||
@@ -155,6 +161,198 @@ test_update_flux_agent_skips_proxy_prompt_when_not_installed() (
|
||||
assert_equals "0" "$ask_called" "update_flux_agent should not prompt for proxy config when the agent is missing"
|
||||
)
|
||||
|
||||
test_install_flux_agent_preserves_legacy_gost_when_download_fails() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
echo "old version"
|
||||
EOF
|
||||
chmod +x "$INSTALL_DIR/flux_agent"
|
||||
SERVER_ADDR="panel.example.com:443"
|
||||
SECRET="secret"
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
|
||||
local cleanup_called="0"
|
||||
local rc="0"
|
||||
|
||||
ask_proxy_config() { :; }
|
||||
ensure_download_url_initialized() { :; }
|
||||
get_config_params() { :; }
|
||||
check_and_install_tcpkill() { :; }
|
||||
cleanup_legacy_gost_installation() {
|
||||
cleanup_called="1"
|
||||
}
|
||||
systemctl() { return 0; }
|
||||
curl() { return 0; }
|
||||
|
||||
( install_flux_agent >/dev/null ) || rc="$?"
|
||||
|
||||
assert_equals "1" "$rc" "install_flux_agent should fail when the download artifact is missing"
|
||||
assert_equals "0" "$cleanup_called" "install_flux_agent should preserve legacy gost when download fails"
|
||||
[[ -f "$INSTALL_DIR/flux_agent" ]] || fail "install_flux_agent should keep the existing flux_agent binary when download fails"
|
||||
)
|
||||
|
||||
test_update_flux_agent_preserves_legacy_gost_when_download_fails() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
echo "old version"
|
||||
EOF
|
||||
chmod +x "$INSTALL_DIR/flux_agent"
|
||||
cat > "$INSTALL_DIR/flux_agent.new" <<'EOF'
|
||||
#!/bin/bash
|
||||
echo "stale version"
|
||||
EOF
|
||||
chmod +x "$INSTALL_DIR/flux_agent.new"
|
||||
|
||||
local cleanup_called="0"
|
||||
local rc="0"
|
||||
|
||||
ask_proxy_config() {
|
||||
PROXY_ENABLED="false"
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
}
|
||||
check_and_install_tcpkill() { :; }
|
||||
cleanup_legacy_gost_installation() {
|
||||
cleanup_called="1"
|
||||
}
|
||||
systemctl() { return 0; }
|
||||
curl() { return 0; }
|
||||
|
||||
update_flux_agent >/dev/null || rc="$?"
|
||||
|
||||
assert_equals "1" "$rc" "update_flux_agent should fail when the download artifact is missing"
|
||||
assert_equals "0" "$cleanup_called" "update_flux_agent should preserve legacy gost when download fails"
|
||||
[[ ! -f "$INSTALL_DIR/flux_agent.new" ]] || fail "update_flux_agent should remove stale download artifacts before retrying"
|
||||
)
|
||||
|
||||
test_install_flux_agent_writes_json_safe_config() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
SERVER_ADDR='panel"addr'
|
||||
SECRET='sec\ret"1'
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
|
||||
ask_proxy_config() { :; }
|
||||
ensure_download_url_initialized() { :; }
|
||||
get_config_params() { :; }
|
||||
check_and_install_tcpkill() { :; }
|
||||
cleanup_legacy_gost_installation() { :; }
|
||||
systemctl() { return 0; }
|
||||
curl() {
|
||||
local output=""
|
||||
while [[ $# -gt 0 ]]; do
|
||||
if [[ "$1" == "-o" ]]; then
|
||||
output="$2"
|
||||
shift 2
|
||||
continue
|
||||
fi
|
||||
shift
|
||||
done
|
||||
|
||||
cat > "$output" <<'EOF'
|
||||
#!/bin/bash
|
||||
echo "new version"
|
||||
EOF
|
||||
chmod +x "$output"
|
||||
}
|
||||
|
||||
( install_flux_agent >/dev/null 2>/dev/null ) || true
|
||||
|
||||
local actual
|
||||
actual=$(<"$INSTALL_DIR/config.json")
|
||||
local expected=$'{\n "addr": "panel\\"addr",\n "secret": "sec\\\\ret\\"1"\n}'
|
||||
|
||||
assert_equals "$expected" "$actual" "install_flux_agent should JSON-escape config values"
|
||||
)
|
||||
|
||||
test_cleanup_legacy_gost_installation_removes_service_and_binary() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
LEGACY_GOST_BINARY=$(mktemp)
|
||||
LEGACY_GOST_SERVICE_FILE_ETC=$(mktemp)
|
||||
LEGACY_GOST_SERVICE_FILE_LIB=$(mktemp -u)
|
||||
LEGACY_GOST_SERVICE_FILE_USR_LIB=$(mktemp -u)
|
||||
LEGACY_GOST_CONFIG_DIR=$(mktemp -d)
|
||||
cat > "$LEGACY_GOST_SERVICE_FILE_ETC" <<EOF
|
||||
[Unit]
|
||||
Description=Gost Proxy Service
|
||||
|
||||
[Service]
|
||||
WorkingDirectory=$LEGACY_GOST_CONFIG_DIR
|
||||
ExecStart=$LEGACY_GOST_CONFIG_DIR/gost
|
||||
EOF
|
||||
: > "$LEGACY_GOST_CONFIG_DIR/config.json"
|
||||
: > "$LEGACY_GOST_CONFIG_DIR/gost.json"
|
||||
|
||||
local systemctl_calls=""
|
||||
|
||||
systemctl() {
|
||||
systemctl_calls+=$'\n'"$*"
|
||||
if [[ "$1" == "list-units" ]]; then
|
||||
printf 'gost.service loaded active running\n'
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
|
||||
cleanup_legacy_gost_installation >/dev/null
|
||||
|
||||
if [[ -e "$LEGACY_GOST_BINARY" ]]; then
|
||||
fail "cleanup_legacy_gost_installation should remove the legacy gost binary"
|
||||
fi
|
||||
if [[ -e "$LEGACY_GOST_SERVICE_FILE_ETC" ]]; then
|
||||
fail "cleanup_legacy_gost_installation should remove the legacy gost service file"
|
||||
fi
|
||||
[[ "$systemctl_calls" == *"stop gost"* ]] || fail "cleanup_legacy_gost_installation should stop the legacy gost service"
|
||||
[[ "$systemctl_calls" == *"disable gost"* ]] || fail "cleanup_legacy_gost_installation should disable the legacy gost service"
|
||||
[[ "$systemctl_calls" == *"daemon-reload"* ]] || fail "cleanup_legacy_gost_installation should reload systemd after removing the legacy service"
|
||||
)
|
||||
|
||||
test_cleanup_legacy_gost_installation_preserves_unrelated_gost() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
LEGACY_GOST_BINARY=$(mktemp)
|
||||
LEGACY_GOST_SERVICE_FILE_ETC=$(mktemp)
|
||||
LEGACY_GOST_SERVICE_FILE_LIB=$(mktemp -u)
|
||||
LEGACY_GOST_SERVICE_FILE_USR_LIB=$(mktemp -u)
|
||||
LEGACY_GOST_CONFIG_DIR=$(mktemp -d)
|
||||
cat > "$LEGACY_GOST_SERVICE_FILE_ETC" <<'EOF'
|
||||
[Unit]
|
||||
Description=Unrelated Gost Service
|
||||
|
||||
[Service]
|
||||
WorkingDirectory=/srv/custom-gost
|
||||
ExecStart=/usr/local/bin/gost -C /srv/custom-gost/gost.yaml
|
||||
EOF
|
||||
|
||||
local systemctl_calls=""
|
||||
|
||||
systemctl() {
|
||||
systemctl_calls+=$'\n'"$*"
|
||||
if [[ "$1" == "list-units" ]]; then
|
||||
printf 'gost.service loaded active running\n'
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
|
||||
cleanup_legacy_gost_installation >/dev/null
|
||||
|
||||
[[ -e "$LEGACY_GOST_BINARY" ]] || fail "cleanup_legacy_gost_installation should preserve unrelated gost binaries"
|
||||
[[ -e "$LEGACY_GOST_SERVICE_FILE_ETC" ]] || fail "cleanup_legacy_gost_installation should preserve unrelated gost service files"
|
||||
[[ "$systemctl_calls" != *"stop gost"* ]] || fail "cleanup_legacy_gost_installation should not stop unrelated gost services"
|
||||
[[ "$systemctl_calls" != *"disable gost"* ]] || fail "cleanup_legacy_gost_installation should not disable unrelated gost services"
|
||||
)
|
||||
|
||||
test_install_script_accepts_proxy_url_env_without_prompt() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
@@ -303,6 +501,11 @@ test_install_script_asks_for_proxy_config
|
||||
test_install_script_recomputes_download_url_after_prompt
|
||||
test_update_flux_agent_asks_for_proxy_config
|
||||
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_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
|
||||
test_panel_install_script_can_disable_proxy
|
||||
test_panel_install_script_recomputes_compose_urls_after_prompt
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
# 多阶段构建 - 构建阶段
|
||||
FROM node:20.19.0 AS builder
|
||||
FROM node:22-alpine AS builder
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
|
||||
@@ -42,6 +42,9 @@ import type {
|
||||
MonitorAccessApiData,
|
||||
TunnelQualityApiItem,
|
||||
StorageSummaryApiData,
|
||||
SystemUpgradeCheckApiData,
|
||||
SystemUpgradeRunApiData,
|
||||
SystemUpgradeVersionApiData,
|
||||
} from "./types";
|
||||
|
||||
import axios from "axios";
|
||||
@@ -258,6 +261,24 @@ export const updateConfig = (name: string, value: string) =>
|
||||
export const getStorageSummary = () =>
|
||||
Network.get<StorageSummaryApiData>("/system/storage");
|
||||
|
||||
export const getSystemUpgradeVersion = () =>
|
||||
Network.post<SystemUpgradeVersionApiData>("/system/version");
|
||||
|
||||
export const checkSystemUpgrade = (channel: ReleaseChannel = "stable") =>
|
||||
Network.post<SystemUpgradeCheckApiData>("/system/check-updates", {
|
||||
channel,
|
||||
});
|
||||
|
||||
export const runSystemUpgrade = (
|
||||
version?: string,
|
||||
channel: ReleaseChannel = "stable",
|
||||
) =>
|
||||
Network.post<SystemUpgradeRunApiData>(
|
||||
"/system/upgrade",
|
||||
{ version: version || "", channel },
|
||||
{ timeout: 60 * 1000 },
|
||||
);
|
||||
|
||||
export const activateLicense = (licenseKey: string) =>
|
||||
Network.post("/license/activate", { license_key: licenseKey });
|
||||
|
||||
|
||||
@@ -48,6 +48,8 @@ export interface TunnelApiItem {
|
||||
trafficRatio?: number;
|
||||
inIp?: string;
|
||||
ipPreference?: string;
|
||||
probeTargetHost?: string;
|
||||
probeTargetPort?: number;
|
||||
inNodeId?: TunnelChainNodePayload[];
|
||||
outNodeId?: TunnelChainNodePayload[];
|
||||
chainNodes?: TunnelChainNodePayload[][];
|
||||
@@ -198,6 +200,14 @@ export interface NodeReleaseApiItem {
|
||||
channel: "stable" | "dev";
|
||||
}
|
||||
|
||||
export interface SystemUpgradeReleaseApiItem {
|
||||
version: string;
|
||||
name: string;
|
||||
publishedAt: string;
|
||||
prerelease: boolean;
|
||||
channel: "stable" | "dev";
|
||||
}
|
||||
|
||||
export interface UserPackageInfoApiData {
|
||||
userInfo: {
|
||||
flow: number;
|
||||
@@ -328,6 +338,8 @@ export interface TunnelMutationPayload {
|
||||
trafficRatio?: number;
|
||||
inIp?: string;
|
||||
ipPreference?: string;
|
||||
probeTargetHost?: string;
|
||||
probeTargetPort?: number;
|
||||
inNodeId?: TunnelChainNodePayload[];
|
||||
outNodeId?: TunnelChainNodePayload[];
|
||||
chainNodes?: TunnelChainNodePayload[][];
|
||||
@@ -480,6 +492,35 @@ export interface StorageSummaryApiData {
|
||||
databaseSizeText: string;
|
||||
}
|
||||
|
||||
export interface SystemUpgradeCapabilityApiData {
|
||||
capable: boolean;
|
||||
reasons: string[];
|
||||
deployDir: string;
|
||||
backendContainer: string;
|
||||
}
|
||||
|
||||
export interface SystemUpgradeVersionApiData {
|
||||
currentVersion: string;
|
||||
latestVersion: string;
|
||||
hasUpdate: boolean;
|
||||
channel: "stable" | "dev";
|
||||
reason?: string;
|
||||
capability: SystemUpgradeCapabilityApiData;
|
||||
}
|
||||
|
||||
export interface SystemUpgradeCheckApiData extends SystemUpgradeVersionApiData {
|
||||
releases: SystemUpgradeReleaseApiItem[];
|
||||
}
|
||||
|
||||
export interface SystemUpgradeRunApiData {
|
||||
version: string;
|
||||
channel: "stable" | "dev";
|
||||
composeAsset: string;
|
||||
helperContainer: string;
|
||||
backendImageId: string;
|
||||
message: string;
|
||||
}
|
||||
|
||||
export interface MonitorNodeApiItem {
|
||||
id: number;
|
||||
inx: number;
|
||||
@@ -525,6 +566,8 @@ export interface TunnelQualityApiItem {
|
||||
exitToBingLatency: number;
|
||||
entryToExitLoss: number;
|
||||
exitToBingLoss: number;
|
||||
probeTargetHost?: string;
|
||||
probeTargetPort?: number;
|
||||
success: boolean;
|
||||
errorMessage?: string;
|
||||
timestamp: number;
|
||||
|
||||
@@ -1,3 +1,10 @@
|
||||
import type {
|
||||
SystemUpgradeCheckApiData,
|
||||
SystemUpgradeRunApiData,
|
||||
SystemUpgradeReleaseApiItem,
|
||||
SystemUpgradeVersionApiData,
|
||||
} from "@/api/types";
|
||||
|
||||
import { useState, useEffect, useRef } from "react";
|
||||
import { useNavigate } from "react-router-dom";
|
||||
import { AnimatePresence, motion } from "framer-motion";
|
||||
@@ -27,6 +34,9 @@ import {
|
||||
getAnnouncement,
|
||||
updateAnnouncement,
|
||||
getStorageSummary,
|
||||
getSystemUpgradeVersion,
|
||||
checkSystemUpgrade,
|
||||
runSystemUpgrade,
|
||||
type AnnouncementData,
|
||||
} from "@/api";
|
||||
import { BackIcon, SettingsIcon } from "@/components/icons";
|
||||
@@ -292,6 +302,19 @@ export default function ConfigPage() {
|
||||
const [updateChannel, setUpdateChannel] = useState<UpdateReleaseChannel>(
|
||||
getUpdateReleaseChannel(),
|
||||
);
|
||||
const [systemUpgradeInfo, setSystemUpgradeInfo] =
|
||||
useState<SystemUpgradeVersionApiData | null>(null);
|
||||
const [systemUpgradeChecking, setSystemUpgradeChecking] = useState(false);
|
||||
const [systemUpgradeExecuting, setSystemUpgradeExecuting] = useState(false);
|
||||
const [systemUpgradeLoading, setSystemUpgradeLoading] = useState(true);
|
||||
const [systemUpgradeModalOpen, setSystemUpgradeModalOpen] = useState(false);
|
||||
const [systemUpgradeReleases, setSystemUpgradeReleases] = useState<
|
||||
SystemUpgradeReleaseApiItem[]
|
||||
>([]);
|
||||
const [systemUpgradeCheckedChannel, setSystemUpgradeCheckedChannel] =
|
||||
useState<UpdateReleaseChannel | null>(null);
|
||||
const [systemUpgradeSelectedVersion, setSystemUpgradeSelectedVersion] =
|
||||
useState("");
|
||||
const [previewLoadFailed, setPreviewLoadFailed] = useState<
|
||||
Partial<Record<BrandPreviewKey, boolean>>
|
||||
>({});
|
||||
@@ -299,6 +322,26 @@ export default function ConfigPage() {
|
||||
Partial<Record<BrandPreviewKey, boolean>>
|
||||
>({});
|
||||
const [storageSummary, setStorageSummary] = useState("加载中...");
|
||||
const systemUpgradeReleasesMatchChannel =
|
||||
systemUpgradeCheckedChannel === updateChannel;
|
||||
const systemUpgradeHasConfirmedUpdate = Boolean(
|
||||
systemUpgradeInfo?.hasUpdate &&
|
||||
systemUpgradeReleasesMatchChannel &&
|
||||
systemUpgradeReleases.length > 0,
|
||||
);
|
||||
const canTriggerSystemUpgrade = Boolean(
|
||||
!systemUpgradeLoading &&
|
||||
!systemUpgradeChecking &&
|
||||
!systemUpgradeExecuting &&
|
||||
systemUpgradeInfo?.capability.capable !== false,
|
||||
);
|
||||
const canOpenSystemUpgradeModal = Boolean(
|
||||
systemUpgradeInfo?.capability.capable &&
|
||||
systemUpgradeHasConfirmedUpdate &&
|
||||
!systemUpgradeLoading &&
|
||||
!systemUpgradeChecking &&
|
||||
!systemUpgradeExecuting,
|
||||
);
|
||||
|
||||
const canGoBack =
|
||||
typeof window !== "undefined" &&
|
||||
@@ -373,11 +416,38 @@ export default function ConfigPage() {
|
||||
}
|
||||
};
|
||||
|
||||
const loadSystemUpgradeInfo = async (channel = updateChannel) => {
|
||||
setSystemUpgradeLoading(true);
|
||||
try {
|
||||
const response = await getSystemUpgradeVersion();
|
||||
|
||||
if (response.code === 0 && response.data) {
|
||||
setSystemUpgradeInfo({
|
||||
...response.data,
|
||||
channel,
|
||||
hasUpdate:
|
||||
response.data.channel === channel ? response.data.hasUpdate : false,
|
||||
latestVersion:
|
||||
response.data.channel === channel
|
||||
? response.data.latestVersion
|
||||
: "",
|
||||
});
|
||||
} else {
|
||||
setSystemUpgradeInfo(null);
|
||||
}
|
||||
} catch {
|
||||
setSystemUpgradeInfo(null);
|
||||
} finally {
|
||||
setSystemUpgradeLoading(false);
|
||||
}
|
||||
};
|
||||
|
||||
useEffect(() => {
|
||||
const timer = setTimeout(() => {
|
||||
loadConfigs(initialConfigs);
|
||||
loadAnnouncement();
|
||||
loadStorageSummary();
|
||||
void loadSystemUpgradeInfo();
|
||||
}, 100);
|
||||
|
||||
return () => clearTimeout(timer);
|
||||
@@ -417,11 +487,94 @@ export default function ConfigPage() {
|
||||
const handleUpdateChannelChange = (channel: UpdateReleaseChannel) => {
|
||||
setUpdateChannel(channel);
|
||||
setUpdateReleaseChannel(channel);
|
||||
setSystemUpgradeSelectedVersion("");
|
||||
setSystemUpgradeReleases([]);
|
||||
setSystemUpgradeCheckedChannel(null);
|
||||
void loadSystemUpgradeInfo(channel);
|
||||
toast.success(
|
||||
`更新通道已切换为${channel === "stable" ? "稳定版" : "开发版"}`,
|
||||
);
|
||||
};
|
||||
|
||||
const handleCheckSystemUpgrade = async () => {
|
||||
const channel = updateChannel;
|
||||
|
||||
setSystemUpgradeChecking(true);
|
||||
try {
|
||||
const response = await checkSystemUpgrade(channel);
|
||||
|
||||
if (response.code === 0 && response.data) {
|
||||
const data = response.data as SystemUpgradeCheckApiData;
|
||||
|
||||
setSystemUpgradeInfo(data);
|
||||
setSystemUpgradeReleases(data.releases || []);
|
||||
setSystemUpgradeCheckedChannel(channel);
|
||||
setSystemUpgradeSelectedVersion("");
|
||||
if (data.latestVersion && !data.hasUpdate) {
|
||||
toast.success("当前已是最新版本");
|
||||
|
||||
return false;
|
||||
}
|
||||
toast.success(
|
||||
data.latestVersion
|
||||
? `已检查到最新版本 ${data.latestVersion}`
|
||||
: "未获取到可用版本",
|
||||
);
|
||||
|
||||
return Boolean(
|
||||
data.capability.capable && data.hasUpdate && data.releases?.length,
|
||||
);
|
||||
} else {
|
||||
setSystemUpgradeReleases([]);
|
||||
setSystemUpgradeCheckedChannel(null);
|
||||
toast.error(response.msg || "检查更新失败");
|
||||
}
|
||||
} catch {
|
||||
setSystemUpgradeReleases([]);
|
||||
setSystemUpgradeCheckedChannel(null);
|
||||
toast.error("检查更新失败,请重试");
|
||||
} finally {
|
||||
setSystemUpgradeChecking(false);
|
||||
}
|
||||
|
||||
return false;
|
||||
};
|
||||
|
||||
const handleOpenSystemUpgradeModal = async () => {
|
||||
if (!canOpenSystemUpgradeModal) {
|
||||
const checked = await handleCheckSystemUpgrade();
|
||||
|
||||
if (!checked) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
setSystemUpgradeModalOpen(true);
|
||||
};
|
||||
|
||||
const handleConfirmSystemUpgrade = async () => {
|
||||
setSystemUpgradeExecuting(true);
|
||||
try {
|
||||
const response = await runSystemUpgrade(
|
||||
systemUpgradeSelectedVersion || undefined,
|
||||
updateChannel,
|
||||
);
|
||||
|
||||
if (response.code === 0 && response.data) {
|
||||
const data = response.data as SystemUpgradeRunApiData;
|
||||
|
||||
setSystemUpgradeModalOpen(false);
|
||||
setSystemUpgradeSelectedVersion("");
|
||||
toast.success(data.message || "升级已触发,请稍后刷新页面");
|
||||
} else {
|
||||
toast.error(response.msg || "面板升级失败");
|
||||
}
|
||||
} catch {
|
||||
toast.error("面板升级失败,请重试");
|
||||
} finally {
|
||||
setSystemUpgradeExecuting(false);
|
||||
}
|
||||
};
|
||||
|
||||
const handleActivateLicense = async () => {
|
||||
if (!licenseKeyInput.trim()) {
|
||||
toast.error("请输入有效的商业授权码");
|
||||
@@ -1329,6 +1482,7 @@ export default function ConfigPage() {
|
||||
</div>
|
||||
|
||||
<Select
|
||||
aria-label="更新通道"
|
||||
selectedKeys={[updateChannel]}
|
||||
size="md"
|
||||
variant="bordered"
|
||||
@@ -1367,6 +1521,164 @@ export default function ConfigPage() {
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<Divider className="my-2" />
|
||||
|
||||
<div className="space-y-4 rounded-xl border border-divider bg-default-50/60 p-4 dark:bg-default-100/10">
|
||||
<div className="space-y-1">
|
||||
<p className="text-sm font-medium text-gray-700 dark:text-gray-300">
|
||||
面板自升级
|
||||
</p>
|
||||
<p className="text-xs text-gray-500 dark:text-gray-400">
|
||||
检查当前版本、可用发布并在容器环境中触发面板升级。
|
||||
</p>
|
||||
</div>
|
||||
|
||||
{systemUpgradeLoading ? (
|
||||
<div className="flex items-center gap-2 rounded-lg border border-divider bg-background px-4 py-3 text-sm text-default-500">
|
||||
<Spinner size="sm" />
|
||||
正在加载升级状态...
|
||||
</div>
|
||||
) : (
|
||||
<div className="space-y-4 rounded-lg border border-divider bg-background px-4 py-4 text-sm text-default-700 dark:text-default-300">
|
||||
<div className="grid gap-3 md:grid-cols-2">
|
||||
<div>
|
||||
<p className="text-xs text-default-500">当前版本</p>
|
||||
<p className="mt-1 font-medium">
|
||||
{systemUpgradeInfo?.currentVersion || "未获取到版本信息"}
|
||||
</p>
|
||||
</div>
|
||||
<div>
|
||||
<p className="text-xs text-default-500">最新版本</p>
|
||||
<p className="mt-1 font-medium">
|
||||
{systemUpgradeInfo?.latestVersion || "未获取到可用版本"}
|
||||
</p>
|
||||
</div>
|
||||
<div>
|
||||
<p className="text-xs text-default-500">当前通道</p>
|
||||
<p className="mt-1 font-medium">
|
||||
{systemUpgradeInfo?.channel === "dev"
|
||||
? "开发版"
|
||||
: systemUpgradeInfo?.channel === "stable"
|
||||
? "稳定版"
|
||||
: updateChannel === "dev"
|
||||
? "开发版"
|
||||
: "稳定版"}
|
||||
</p>
|
||||
</div>
|
||||
<div>
|
||||
<p className="text-xs text-default-500">升级能力</p>
|
||||
<p className="mt-1 font-medium">
|
||||
{systemUpgradeInfo?.capability.capable
|
||||
? "可升级"
|
||||
: "当前不可升级"}
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div className="grid gap-3 md:grid-cols-2">
|
||||
<div>
|
||||
<p className="text-xs text-default-500">部署目录</p>
|
||||
<p className="mt-1 break-all font-medium">
|
||||
{systemUpgradeInfo?.capability.deployDir ||
|
||||
"未获取到部署目录"}
|
||||
</p>
|
||||
</div>
|
||||
<div>
|
||||
<p className="text-xs text-default-500">后端容器</p>
|
||||
<p className="mt-1 break-all font-medium">
|
||||
{systemUpgradeInfo?.capability.backendContainer ||
|
||||
"未获取到容器信息"}
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
{!systemUpgradeInfo?.capability.capable && (
|
||||
<div className="rounded-lg border border-warning-200 bg-warning-50 px-4 py-3 text-warning-800 dark:border-warning-900/40 dark:bg-warning-950/30 dark:text-warning-200">
|
||||
<p className="text-xs font-medium">当前无法升级</p>
|
||||
<ul className="mt-2 list-disc space-y-1 pl-4 text-xs">
|
||||
{(systemUpgradeInfo?.capability.reasons?.length
|
||||
? systemUpgradeInfo.capability.reasons
|
||||
: ["暂未获取到不可升级原因"]
|
||||
).map((reason) => (
|
||||
<li key={reason}>{reason}</li>
|
||||
))}
|
||||
</ul>
|
||||
</div>
|
||||
)}
|
||||
|
||||
<div className="space-y-3">
|
||||
<div className="flex flex-col gap-1">
|
||||
<p className="text-sm font-medium text-gray-700 dark:text-gray-300">
|
||||
可用发布版本
|
||||
</p>
|
||||
<p className="text-xs text-gray-500 dark:text-gray-400">
|
||||
选择指定版本后执行升级;留空则使用当前通道下最新可用版本。
|
||||
</p>
|
||||
</div>
|
||||
|
||||
<Select
|
||||
aria-label="目标版本"
|
||||
isDisabled={
|
||||
!systemUpgradeReleasesMatchChannel ||
|
||||
systemUpgradeReleases.length === 0 ||
|
||||
systemUpgradeExecuting
|
||||
}
|
||||
placeholder={
|
||||
systemUpgradeReleasesMatchChannel &&
|
||||
systemUpgradeReleases.length > 0
|
||||
? "留空时自动选择最新版本"
|
||||
: "请先检查当前通道更新"
|
||||
}
|
||||
selectedKeys={
|
||||
systemUpgradeSelectedVersion
|
||||
? [systemUpgradeSelectedVersion]
|
||||
: []
|
||||
}
|
||||
size="md"
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selected = Array.from(keys)[0] as
|
||||
| string
|
||||
| undefined;
|
||||
|
||||
setSystemUpgradeSelectedVersion(selected || "");
|
||||
}}
|
||||
>
|
||||
{(systemUpgradeReleasesMatchChannel
|
||||
? systemUpgradeReleases
|
||||
: []
|
||||
).map((release) => (
|
||||
<SelectItem
|
||||
key={release.version}
|
||||
description={release.publishedAt || "暂无发布时间"}
|
||||
>
|
||||
{release.name || release.version}
|
||||
</SelectItem>
|
||||
))}
|
||||
</Select>
|
||||
</div>
|
||||
</div>
|
||||
)}
|
||||
|
||||
<div className="flex flex-col gap-3 pt-1 sm:flex-row sm:justify-end">
|
||||
<Button
|
||||
isLoading={systemUpgradeChecking}
|
||||
variant="flat"
|
||||
onPress={handleCheckSystemUpgrade}
|
||||
>
|
||||
检查更新
|
||||
</Button>
|
||||
<Button
|
||||
color="primary"
|
||||
isDisabled={!canTriggerSystemUpgrade}
|
||||
isLoading={systemUpgradeExecuting}
|
||||
onPress={handleOpenSystemUpgradeModal}
|
||||
>
|
||||
立即升级
|
||||
</Button>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div className="flex justify-end pt-6 border-t border-divider/50 mt-4">
|
||||
<Button
|
||||
color="primary"
|
||||
@@ -1606,6 +1918,64 @@ export default function ConfigPage() {
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
|
||||
<Modal
|
||||
backdrop="blur"
|
||||
classNames={{
|
||||
base: "!w-[calc(100%-32px)] !mx-auto sm:!w-full rounded-2xl overflow-hidden",
|
||||
}}
|
||||
isOpen={systemUpgradeModalOpen}
|
||||
onOpenChange={(open) => {
|
||||
if (!systemUpgradeExecuting) {
|
||||
setSystemUpgradeModalOpen(open);
|
||||
}
|
||||
}}
|
||||
>
|
||||
<ModalContent>
|
||||
{(onClose) => (
|
||||
<>
|
||||
<ModalHeader>确认面板升级</ModalHeader>
|
||||
<ModalBody>
|
||||
<div className="space-y-3 text-sm text-default-700 dark:text-default-300">
|
||||
<p>
|
||||
升级过程需要访问 Docker
|
||||
Socket,并会在短时间内中断当前面板服务。
|
||||
</p>
|
||||
<p>
|
||||
请确认已经允许面板管理容器与宿主机 Docker
|
||||
交互,并且可以接受升级期间的临时不可用。
|
||||
</p>
|
||||
<div className="space-y-2 rounded-lg border border-warning-200 bg-warning-50 px-4 py-3 text-warning-800 dark:border-warning-900/40 dark:bg-warning-950/30 dark:text-warning-200">
|
||||
<p className="text-xs font-medium">升级前请确认</p>
|
||||
<ul className="list-disc space-y-1 pl-4 text-xs">
|
||||
<li>Docker Socket 可用且挂载权限正常。</li>
|
||||
<li>当前面板允许短暂停止和重启。</li>
|
||||
<li>已选择正确的更新通道与目标版本。</li>
|
||||
</ul>
|
||||
</div>
|
||||
</div>
|
||||
</ModalBody>
|
||||
<ModalFooter>
|
||||
<Button
|
||||
isDisabled={systemUpgradeExecuting}
|
||||
variant="light"
|
||||
onPress={onClose}
|
||||
>
|
||||
取消
|
||||
</Button>
|
||||
<Button
|
||||
color="primary"
|
||||
isDisabled={systemUpgradeExecuting}
|
||||
isLoading={systemUpgradeExecuting}
|
||||
onPress={handleConfirmSystemUpgrade}
|
||||
>
|
||||
确认升级
|
||||
</Button>
|
||||
</ModalFooter>
|
||||
</>
|
||||
)}
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
|
||||
{/* Floating Save Button (FAB) */}
|
||||
<AnimatePresence>
|
||||
{hasChanges && (
|
||||
|
||||
@@ -2582,11 +2582,16 @@ export default function NodePage() {
|
||||
/>
|
||||
|
||||
{/* 高级配置 */}
|
||||
<Accordion variant="bordered">
|
||||
<Accordion className="px-0" variant="light">
|
||||
<AccordionItem
|
||||
key="advanced"
|
||||
aria-label="高级配置"
|
||||
title="高级配置"
|
||||
className="border-b-0 [&_[data-slot=accordion-trigger]]:no-underline [&_[data-slot=accordion-trigger]]:hover:no-underline"
|
||||
title={
|
||||
<span className="text-small text-default-500 font-medium">
|
||||
高级配置
|
||||
</span>
|
||||
}
|
||||
>
|
||||
<div className="space-y-4 pb-2">
|
||||
<Input
|
||||
@@ -2683,10 +2688,10 @@ export default function NodePage() {
|
||||
/>
|
||||
)}
|
||||
<div
|
||||
className={`grid grid-cols-1 sm:grid-cols-3 gap-3 bg-default-50 dark:bg-default-100 p-3 rounded-md border border-default-200 dark:border-default-100/30 ${protocolDisabled ? "opacity-70" : ""}`}
|
||||
className={`grid grid-cols-1 sm:grid-cols-3 gap-3 bg-content1/30 dark:bg-content1/20 p-3 rounded-md border border-divider ${protocolDisabled ? "opacity-70" : ""}`}
|
||||
>
|
||||
{/* HTTP tile */}
|
||||
<div className="px-3 py-3 rounded-lg bg-white dark:bg-default-50 border border-default-200 dark:border-default-100/30 hover:border-primary-200 transition-colors">
|
||||
<div className="px-3 py-3 rounded-lg bg-content1/55 dark:bg-content1/35 border border-divider hover:border-primary-200 dark:hover:border-primary-500/30 transition-colors">
|
||||
<div className="flex items-center gap-2 mb-2">
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
@@ -2727,7 +2732,7 @@ export default function NodePage() {
|
||||
</div>
|
||||
|
||||
{/* TLS tile */}
|
||||
<div className="px-3 py-3 rounded-lg bg-white dark:bg-default-50 border border-default-200 dark:border-default-100/30 hover:border-primary-200 transition-colors">
|
||||
<div className="px-3 py-3 rounded-lg bg-content1/55 dark:bg-content1/35 border border-divider hover:border-primary-200 dark:hover:border-primary-500/30 transition-colors">
|
||||
<div className="flex items-center gap-2 mb-2">
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
@@ -2771,7 +2776,7 @@ export default function NodePage() {
|
||||
</div>
|
||||
|
||||
{/* SOCKS tile */}
|
||||
<div className="px-3 py-3 rounded-lg bg-white dark:bg-default-50 border border-default-200 dark:border-default-100/30 hover:border-primary-200 transition-colors">
|
||||
<div className="px-3 py-3 rounded-lg bg-content1/55 dark:bg-content1/35 border border-divider hover:border-primary-200 dark:hover:border-primary-500/30 transition-colors">
|
||||
<div className="flex items-center gap-2 mb-2">
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
|
||||
@@ -63,6 +63,19 @@ const MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY =
|
||||
"monitor_tunnel_quality_enabled";
|
||||
const MONITOR_TUNNEL_QUALITY_ENABLED_EVENT =
|
||||
"monitorTunnelQualityEnabledChanged";
|
||||
const DEFAULT_PROBE_TARGET_LABEL = "www.bing.com:443";
|
||||
|
||||
const probeTargetLabel = (quality?: TunnelQualityApiItem | null) => {
|
||||
if (!quality?.probeTargetHost || !quality.probeTargetPort) {
|
||||
return DEFAULT_PROBE_TARGET_LABEL;
|
||||
}
|
||||
|
||||
const host = quality.probeTargetHost.includes(":")
|
||||
? `[${quality.probeTargetHost}]`
|
||||
: quality.probeTargetHost;
|
||||
|
||||
return `${host}:${quality.probeTargetPort}`;
|
||||
};
|
||||
|
||||
const formatTimestamp = (ts: number, rangeMs?: number): string => {
|
||||
const date = new Date(ts);
|
||||
@@ -177,7 +190,7 @@ function UptimeHistoryBar({
|
||||
}
|
||||
|
||||
const displayLatency = latency >= 0 ? `${latency.toFixed(0)}ms` : "-";
|
||||
const tooltip = `${timeStr} | ${displayLatency} | ${statusText}`;
|
||||
const tooltip = `${timeStr} | ${displayLatency} | ${statusText} | 测试目标: ${probeTargetLabel(q)}`;
|
||||
|
||||
return (
|
||||
<div
|
||||
@@ -305,7 +318,7 @@ const QualityChartCard = React.memo(function QualityChartCard({
|
||||
name === "entryToExit"
|
||||
? "入口→出口"
|
||||
: name === "exitToBing"
|
||||
? "出口→Bing"
|
||||
? "出口→测试目标"
|
||||
: name;
|
||||
|
||||
return [`${n.toFixed(1)}ms`, label];
|
||||
@@ -1001,9 +1014,12 @@ export function TunnelMonitorView({
|
||||
</Card>
|
||||
<Card className="border border-divider/60 shadow-sm hover:shadow-md transition-shadow bg-gradient-to-br from-background to-default-50/50">
|
||||
<CardBody className="py-3 px-4 flex flex-col items-center justify-center min-h-[5rem]">
|
||||
<span className="text-[11px] text-default-500 mb-1.5 flex items-center gap-1">
|
||||
<span
|
||||
className="text-[11px] text-default-500 mb-1.5 flex items-center gap-1"
|
||||
title={probeTargetLabel(quality)}
|
||||
>
|
||||
<Globe className="w-3 h-3" />
|
||||
出口 → Bing 延迟
|
||||
出口 → 测试目标 延迟
|
||||
</span>
|
||||
<LatencyDisplay
|
||||
loading={qualityLoading}
|
||||
@@ -1027,8 +1043,11 @@ export function TunnelMonitorView({
|
||||
</Card>
|
||||
<Card className="border border-divider/60 shadow-sm hover:shadow-md transition-shadow bg-gradient-to-br from-background to-default-50/50">
|
||||
<CardBody className="py-3 px-4 flex flex-col items-center justify-center min-h-[5rem]">
|
||||
<span className="text-[11px] text-default-500 mb-1.5">
|
||||
出口 → Bing 丢包
|
||||
<span
|
||||
className="text-[11px] text-default-500 mb-1.5"
|
||||
title={probeTargetLabel(quality)}
|
||||
>
|
||||
出口 → 测试目标 丢包
|
||||
</span>
|
||||
<span
|
||||
className={`text-sm font-semibold font-mono ${(quality?.exitToBingLoss ?? 0) > 0 ? "text-warning" : ""}`}
|
||||
@@ -1054,6 +1073,9 @@ export function TunnelMonitorView({
|
||||
<span>实时隧道质量检测已关闭</span>
|
||||
</>
|
||||
)}
|
||||
<span className="text-default-400">
|
||||
· 测试目标: {probeTargetLabel(quality)}
|
||||
</span>
|
||||
{quality?.timestamp && (
|
||||
<span className="text-default-400">
|
||||
· 最近更新:{" "}
|
||||
@@ -1148,6 +1170,7 @@ export function TunnelMonitorView({
|
||||
{tunnels.map((tunnel) => {
|
||||
const quality = qualityMap[tunnel.id];
|
||||
const isEnabled = tunnel.status === 1;
|
||||
const targetLabel = probeTargetLabel(quality);
|
||||
|
||||
return (
|
||||
<Card
|
||||
@@ -1203,7 +1226,13 @@ export function TunnelMonitorView({
|
||||
<div className="space-y-1">
|
||||
<div className="text-[10px] text-default-500 flex items-center gap-1">
|
||||
<Globe className="w-3 h-3" />
|
||||
出口→Bing
|
||||
出口→测试目标
|
||||
</div>
|
||||
<div
|
||||
aria-label={`测试目标 ${targetLabel}`}
|
||||
className="text-[10px] text-default-400 truncate"
|
||||
>
|
||||
{targetLabel}
|
||||
</div>
|
||||
<UptimeHistoryBar
|
||||
history={qualityHistoryMap[tunnel.id]}
|
||||
@@ -1259,13 +1288,14 @@ export function TunnelMonitorView({
|
||||
<TableColumn>状态</TableColumn>
|
||||
<TableColumn>名称</TableColumn>
|
||||
<TableColumn>入口→出口</TableColumn>
|
||||
<TableColumn>出口→Bing</TableColumn>
|
||||
<TableColumn>出口→测试目标</TableColumn>
|
||||
<TableColumn>更新时间</TableColumn>
|
||||
</TableHeader>
|
||||
<TableBody emptyContent="暂无隧道">
|
||||
{tunnels.map((tunnel) => {
|
||||
const quality = qualityMap[tunnel.id];
|
||||
const isEnabled = tunnel.status === 1;
|
||||
const targetLabel = probeTargetLabel(quality);
|
||||
|
||||
return (
|
||||
<TableRow
|
||||
@@ -1295,11 +1325,19 @@ export function TunnelMonitorView({
|
||||
/>
|
||||
</TableCell>
|
||||
<TableCell>
|
||||
<UptimeHistoryBar
|
||||
history={qualityHistoryMap[tunnel.id]}
|
||||
latestValue={quality?.exitToBingLatency}
|
||||
type="exitToBing"
|
||||
/>
|
||||
<div
|
||||
aria-label={`出口到测试目标 ${targetLabel}`}
|
||||
className="space-y-1"
|
||||
>
|
||||
<UptimeHistoryBar
|
||||
history={qualityHistoryMap[tunnel.id]}
|
||||
latestValue={quality?.exitToBingLatency}
|
||||
type="exitToBing"
|
||||
/>
|
||||
<span className="block max-w-[160px] truncate text-[10px] text-default-400">
|
||||
{targetLabel}
|
||||
</span>
|
||||
</div>
|
||||
</TableCell>
|
||||
<TableCell>
|
||||
{quality?.timestamp ? (
|
||||
|
||||
@@ -46,6 +46,7 @@ import { Alert } from "@/shadcn-bridge/heroui/alert";
|
||||
import { Checkbox } from "@/shadcn-bridge/heroui/checkbox";
|
||||
import { Progress } from "@/shadcn-bridge/heroui/progress";
|
||||
import { Radio, RadioGroup } from "@/shadcn-bridge/heroui/radio";
|
||||
import { Accordion, AccordionItem } from "@/shadcn-bridge/heroui/accordion";
|
||||
import {
|
||||
Table,
|
||||
TableHeader,
|
||||
@@ -130,11 +131,21 @@ interface Tunnel {
|
||||
flow: number; // 1: 单向, 2: 双向
|
||||
trafficRatio: number;
|
||||
ipPreference?: string;
|
||||
probeTargetHost?: string;
|
||||
probeTargetPort?: number;
|
||||
bestExitState?: BestExitState;
|
||||
status: number;
|
||||
createdTime: string;
|
||||
}
|
||||
|
||||
const DEFAULT_PROBE_TARGET_HOST = "www.bing.com";
|
||||
const DEFAULT_PROBE_TARGET_PORT = 443;
|
||||
|
||||
const getTunnelDiagnosisTarget = (tunnel: Tunnel) => ({
|
||||
targetIp: tunnel.probeTargetHost || DEFAULT_PROBE_TARGET_HOST,
|
||||
targetPort: tunnel.probeTargetPort || DEFAULT_PROBE_TARGET_PORT,
|
||||
});
|
||||
|
||||
interface Node {
|
||||
id: number;
|
||||
name: string;
|
||||
@@ -156,6 +167,8 @@ interface TunnelForm {
|
||||
trafficRatio: number;
|
||||
inIp: string; // 入口IP
|
||||
ipPreference: string;
|
||||
probeTargetHost?: string;
|
||||
probeTargetPort?: number;
|
||||
status: number;
|
||||
}
|
||||
|
||||
@@ -191,6 +204,7 @@ const isObjectRecord = (value: unknown): value is Record<string, unknown> =>
|
||||
const toSafeString = (value: unknown): string => {
|
||||
if (typeof value === "string") return value;
|
||||
if (typeof value === "number" && Number.isFinite(value)) return String(value);
|
||||
|
||||
return "";
|
||||
};
|
||||
|
||||
@@ -252,11 +266,13 @@ const normalizeBestExitState = (value: unknown): BestExitState | undefined => {
|
||||
|
||||
const bestExitOwnerRoleText = (role?: string) => {
|
||||
if (role === "chain") return "中转";
|
||||
|
||||
return "入口";
|
||||
};
|
||||
|
||||
const bestExitDetailTitle = (state?: BestExitState) => {
|
||||
if (!state?.items?.length) return undefined;
|
||||
|
||||
return state.items
|
||||
.map(
|
||||
(item) =>
|
||||
@@ -584,6 +600,8 @@ export default function TunnelPage() {
|
||||
.join("\n")
|
||||
: "",
|
||||
ipPreference: tunnel.ipPreference || "",
|
||||
probeTargetHost: tunnel.probeTargetHost || "",
|
||||
probeTargetPort: tunnel.probeTargetPort || 0,
|
||||
status: tunnel.status,
|
||||
});
|
||||
setErrors({});
|
||||
@@ -860,12 +878,18 @@ export default function TunnelPage() {
|
||||
.map((ip) => ip.trim())
|
||||
.filter((ip) => ip)
|
||||
.join(",");
|
||||
const probeTargetHost = (form.probeTargetHost || "").trim();
|
||||
const probeTargetPort = probeTargetHost
|
||||
? Number(form.probeTargetPort || 0)
|
||||
: 0;
|
||||
|
||||
const data = {
|
||||
...form,
|
||||
inIp: inIpString,
|
||||
outNodeId: cleanedOutNodeId,
|
||||
chainNodes: cleanedChainNodes,
|
||||
probeTargetHost,
|
||||
probeTargetPort,
|
||||
};
|
||||
|
||||
const response = isEdit
|
||||
@@ -890,6 +914,7 @@ export default function TunnelPage() {
|
||||
const handleDiagnose = async (tunnel: Tunnel) => {
|
||||
diagnosisAbortRef.current?.abort();
|
||||
const abortController = new AbortController();
|
||||
const diagnosisTarget = getTunnelDiagnosisTarget(tunnel);
|
||||
|
||||
diagnosisAbortRef.current = abortController;
|
||||
|
||||
@@ -1031,6 +1056,7 @@ export default function TunnelPage() {
|
||||
tunnelType: tunnel.type,
|
||||
description: "诊断失败",
|
||||
message: response.msg || "诊断过程中发生错误",
|
||||
...diagnosisTarget,
|
||||
}),
|
||||
);
|
||||
setDiagnosisProgress({
|
||||
@@ -1062,6 +1088,7 @@ export default function TunnelPage() {
|
||||
tunnelType: tunnel.type,
|
||||
description: "网络错误",
|
||||
message: "无法连接到服务器",
|
||||
...diagnosisTarget,
|
||||
}),
|
||||
);
|
||||
setDiagnosisProgress({
|
||||
@@ -2244,85 +2271,6 @@ export default function TunnelPage() {
|
||||
<SelectItem key="2">隧道转发</SelectItem>
|
||||
</Select>
|
||||
|
||||
<div className="grid grid-cols-1 md:grid-cols-2 gap-4">
|
||||
<Select
|
||||
errorMessage={errors.flow}
|
||||
isInvalid={!!errors.flow}
|
||||
label="流量计算"
|
||||
placeholder="请选择流量计算方式"
|
||||
selectedKeys={[form.flow.toString()]}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
|
||||
if (selectedKey) {
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
flow: parseInt(selectedKey),
|
||||
}));
|
||||
}
|
||||
}}
|
||||
>
|
||||
<SelectItem key="1">单向计算(仅上传)</SelectItem>
|
||||
<SelectItem key="2">双向计算(上传+下载)</SelectItem>
|
||||
</Select>
|
||||
|
||||
<Input
|
||||
errorMessage={errors.trafficRatio}
|
||||
isInvalid={!!errors.trafficRatio}
|
||||
label="流量倍率"
|
||||
max={100}
|
||||
min={0.01}
|
||||
placeholder="例如:0.5 或 1 或 2"
|
||||
step="any"
|
||||
type="number"
|
||||
value={form.trafficRatio.toString()}
|
||||
variant="bordered"
|
||||
onChange={(e) =>
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
trafficRatio: parseFloat(e.target.value) || 0,
|
||||
}))
|
||||
}
|
||||
/>
|
||||
</div>
|
||||
|
||||
<Textarea
|
||||
description="入口IP由系统自动从入口节点采集,无需手动填写。支持多个IP,每行一个地址,留空则使用入口节点IP"
|
||||
errorMessage={errors.inIp}
|
||||
isInvalid={!!errors.inIp}
|
||||
label="入口IP"
|
||||
maxRows={5}
|
||||
minRows={3}
|
||||
placeholder="一行一个IP地址或域名,例如: 192.168.1.100 example.com"
|
||||
value={form.inIp}
|
||||
variant="bordered"
|
||||
onChange={(e) =>
|
||||
setForm((prev) => ({ ...prev, inIp: e.target.value }))
|
||||
}
|
||||
/>
|
||||
|
||||
{form.type === 2 && (
|
||||
<Select
|
||||
description="当节点同时拥有IPv4和IPv6地址时,选择隧道连接使用的地址类型"
|
||||
label="隧道连接地址偏好"
|
||||
placeholder="自动选择"
|
||||
selectedKeys={[form.ipPreference || ""]}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
ipPreference: selectedKey || "",
|
||||
}));
|
||||
}}
|
||||
>
|
||||
<SelectItem key="v4">优先IPv4</SelectItem>
|
||||
<SelectItem key="v6">优先IPv6</SelectItem>
|
||||
</Select>
|
||||
)}
|
||||
|
||||
<Divider />
|
||||
<h3 className="text-lg font-semibold">入口配置</h3>
|
||||
|
||||
@@ -3130,6 +3078,154 @@ export default function TunnelPage() {
|
||||
})()}
|
||||
</>
|
||||
)}
|
||||
|
||||
<Accordion className="px-0" variant="light">
|
||||
<AccordionItem
|
||||
key="advanced"
|
||||
aria-label="高级设置"
|
||||
className="border-b-0 [&_[data-slot=accordion-trigger]]:no-underline [&_[data-slot=accordion-trigger]]:hover:no-underline"
|
||||
title={
|
||||
<span className="text-small text-default-500 font-medium">
|
||||
高级设置
|
||||
</span>
|
||||
}
|
||||
>
|
||||
<div className="space-y-4 pb-2">
|
||||
<div className="grid grid-cols-1 md:grid-cols-2 gap-4">
|
||||
<Select
|
||||
errorMessage={errors.flow}
|
||||
isInvalid={!!errors.flow}
|
||||
label="流量计算"
|
||||
placeholder="请选择流量计算方式"
|
||||
selectedKeys={[form.flow.toString()]}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
|
||||
if (selectedKey) {
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
flow: parseInt(selectedKey),
|
||||
}));
|
||||
}
|
||||
}}
|
||||
>
|
||||
<SelectItem key="1">单向计算(仅上传)</SelectItem>
|
||||
<SelectItem key="2">
|
||||
双向计算(上传+下载)
|
||||
</SelectItem>
|
||||
</Select>
|
||||
|
||||
<Input
|
||||
errorMessage={errors.trafficRatio}
|
||||
isInvalid={!!errors.trafficRatio}
|
||||
label="流量倍率"
|
||||
max={100}
|
||||
min={0.01}
|
||||
placeholder="例如:0.5 或 1 或 2"
|
||||
step="any"
|
||||
type="number"
|
||||
value={form.trafficRatio.toString()}
|
||||
variant="bordered"
|
||||
onChange={(e) =>
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
trafficRatio: parseFloat(e.target.value) || 0,
|
||||
}))
|
||||
}
|
||||
/>
|
||||
</div>
|
||||
|
||||
<Textarea
|
||||
description="入口IP由系统自动从入口节点采集,无需手动填写。支持多个IP,每行一个地址,留空则使用入口节点IP"
|
||||
errorMessage={errors.inIp}
|
||||
isInvalid={!!errors.inIp}
|
||||
label="入口IP"
|
||||
maxRows={5}
|
||||
minRows={3}
|
||||
placeholder="一行一个IP地址或域名,例如: 192.168.1.100 example.com"
|
||||
value={form.inIp}
|
||||
variant="bordered"
|
||||
onChange={(e) =>
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
inIp: e.target.value,
|
||||
}))
|
||||
}
|
||||
/>
|
||||
|
||||
{form.type === 2 && (
|
||||
<Select
|
||||
description="当节点同时拥有IPv4和IPv6地址时,选择隧道连接使用的地址类型"
|
||||
label="隧道连接地址偏好"
|
||||
placeholder="自动选择"
|
||||
selectedKeys={[form.ipPreference || ""]}
|
||||
variant="bordered"
|
||||
onSelectionChange={(keys) => {
|
||||
const selectedKey = Array.from(keys)[0] as string;
|
||||
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
ipPreference: selectedKey || "",
|
||||
}));
|
||||
}}
|
||||
>
|
||||
<SelectItem key="v4">优先IPv4</SelectItem>
|
||||
<SelectItem key="v6">优先IPv6</SelectItem>
|
||||
</Select>
|
||||
)}
|
||||
|
||||
<div>
|
||||
<div className="text-sm font-medium">
|
||||
质量检测目标
|
||||
</div>
|
||||
<p className="text-xs text-default-500 mt-0.5">
|
||||
用于实时隧道质量检测、诊断目标和 best
|
||||
最优出口评分,留空使用 www.bing.com:443
|
||||
</p>
|
||||
</div>
|
||||
<div className="grid grid-cols-1 md:grid-cols-[1fr_140px] gap-3">
|
||||
<Input
|
||||
errorMessage={errors.probeTargetHost}
|
||||
isInvalid={!!errors.probeTargetHost}
|
||||
label="Host"
|
||||
placeholder="www.bing.com"
|
||||
value={form.probeTargetHost || ""}
|
||||
variant="bordered"
|
||||
onChange={(e) =>
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
probeTargetHost: e.target.value,
|
||||
}))
|
||||
}
|
||||
/>
|
||||
<Input
|
||||
errorMessage={errors.probeTargetPort}
|
||||
isInvalid={!!errors.probeTargetPort}
|
||||
label="Port"
|
||||
max={65535}
|
||||
min={1}
|
||||
placeholder="443"
|
||||
type="number"
|
||||
value={
|
||||
form.probeTargetPort
|
||||
? String(form.probeTargetPort)
|
||||
: ""
|
||||
}
|
||||
variant="bordered"
|
||||
onChange={(e) =>
|
||||
setForm((prev) => ({
|
||||
...prev,
|
||||
probeTargetPort: e.target.value
|
||||
? Number(e.target.value)
|
||||
: 0,
|
||||
}))
|
||||
}
|
||||
/>
|
||||
</div>
|
||||
</div>
|
||||
</AccordionItem>
|
||||
</Accordion>
|
||||
</div>
|
||||
</ModalBody>
|
||||
<ModalFooter>
|
||||
|
||||
@@ -27,6 +27,8 @@ export interface DiagnosisFallbackInput {
|
||||
tunnelType: number;
|
||||
description: string;
|
||||
message: string;
|
||||
targetIp?: string;
|
||||
targetPort?: number;
|
||||
}
|
||||
|
||||
export const buildDiagnosisFallbackResult = ({
|
||||
@@ -34,6 +36,8 @@ export const buildDiagnosisFallbackResult = ({
|
||||
tunnelType,
|
||||
description,
|
||||
message,
|
||||
targetIp = "-",
|
||||
targetPort = 443,
|
||||
}: DiagnosisFallbackInput): DiagnosisResult => {
|
||||
return {
|
||||
tunnelName,
|
||||
@@ -45,8 +49,8 @@ export const buildDiagnosisFallbackResult = ({
|
||||
description,
|
||||
nodeName: "-",
|
||||
nodeId: "-",
|
||||
targetIp: "-",
|
||||
targetPort: 443,
|
||||
targetIp,
|
||||
targetPort,
|
||||
message,
|
||||
},
|
||||
],
|
||||
|
||||
@@ -8,6 +8,8 @@ interface TunnelFormInput {
|
||||
inNodeId: TunnelChainNode[];
|
||||
outNodeId?: TunnelChainNode[];
|
||||
trafficRatio: number;
|
||||
probeTargetHost?: string;
|
||||
probeTargetPort?: number;
|
||||
}
|
||||
|
||||
interface TunnelNodeInput {
|
||||
@@ -15,6 +17,96 @@ interface TunnelNodeInput {
|
||||
status: number;
|
||||
}
|
||||
|
||||
const isValidProbeIPv4 = (host: string) => {
|
||||
const parts = host.split(".");
|
||||
|
||||
return (
|
||||
parts.length === 4 &&
|
||||
parts.every((part) => {
|
||||
if (!/^\d+$/.test(part)) {
|
||||
return false;
|
||||
}
|
||||
if (part.length > 1 && part.startsWith("0")) {
|
||||
return false;
|
||||
}
|
||||
|
||||
const value = Number(part);
|
||||
|
||||
return value >= 0 && value <= 255;
|
||||
})
|
||||
);
|
||||
};
|
||||
|
||||
const isIPv4LikeProbeHost = (host: string) =>
|
||||
/^[0-9.]+$/.test(host) && host.includes(".");
|
||||
|
||||
const isValidProbeIPv6 = (host: string) => {
|
||||
let value = host;
|
||||
|
||||
if (host.startsWith("[") || host.endsWith("]")) {
|
||||
if (!host.startsWith("[") || !host.endsWith("]")) {
|
||||
return false;
|
||||
}
|
||||
value = host.slice(1, -1);
|
||||
}
|
||||
|
||||
if (!value.includes(":") || value.includes("[") || value.includes("]")) {
|
||||
return false;
|
||||
}
|
||||
|
||||
try {
|
||||
const url = new URL(`http://[${value}]`);
|
||||
|
||||
return url.hostname.length > 0;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
};
|
||||
|
||||
const isSchemeLikeProbeHost = (host: string) => {
|
||||
if (isValidProbeIPv6(host)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
const colonIndex = host.indexOf(":");
|
||||
|
||||
if (colonIndex <= 0) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return /^[A-Za-z][A-Za-z0-9+.-]*$/.test(host.slice(0, colonIndex));
|
||||
};
|
||||
|
||||
const isValidProbeDomain = (host: string) => {
|
||||
if (!host || host.length > 253) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return host.split(".").every((label) => {
|
||||
if (
|
||||
!label ||
|
||||
label.length > 63 ||
|
||||
label.startsWith("-") ||
|
||||
label.endsWith("-")
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return /^[A-Za-z0-9-]+$/.test(label);
|
||||
});
|
||||
};
|
||||
|
||||
const isValidProbeTargetHost = (host: string) => {
|
||||
if (isValidProbeIPv6(host) || isValidProbeIPv4(host)) {
|
||||
return true;
|
||||
}
|
||||
if (host.includes(":") || isIPv4LikeProbeHost(host)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return isValidProbeDomain(host);
|
||||
};
|
||||
|
||||
export const createTunnelFormDefaults = () => {
|
||||
return {
|
||||
name: "",
|
||||
@@ -26,6 +118,8 @@ export const createTunnelFormDefaults = () => {
|
||||
trafficRatio: 1.0,
|
||||
inIp: "",
|
||||
ipPreference: "",
|
||||
probeTargetHost: "",
|
||||
probeTargetPort: 0,
|
||||
status: 1,
|
||||
};
|
||||
};
|
||||
@@ -63,6 +157,31 @@ export const validateTunnelForm = (
|
||||
errors.trafficRatio = "流量倍率须大于0,支持小数(如 0.5)";
|
||||
}
|
||||
|
||||
const rawProbeHost = form.probeTargetHost || "";
|
||||
const probeHost = rawProbeHost.trim();
|
||||
const probePortInput = form.probeTargetPort;
|
||||
const probePort = Number(probePortInput ?? 0);
|
||||
const hasProbeHostInput = rawProbeHost.length > 0;
|
||||
const hasProbePort = probePortInput != null && probePortInput !== 0;
|
||||
|
||||
if (hasProbeHostInput || hasProbePort) {
|
||||
if (!probeHost) {
|
||||
errors.probeTargetHost = "请输入测试目标 Host";
|
||||
} else if (
|
||||
probeHost.includes("://") ||
|
||||
/[\s/?#]/.test(rawProbeHost) ||
|
||||
isSchemeLikeProbeHost(probeHost)
|
||||
) {
|
||||
errors.probeTargetHost = "Host 不能包含协议、端口、空格或路径";
|
||||
} else if (!isValidProbeTargetHost(probeHost)) {
|
||||
errors.probeTargetHost = "测试目标 Host 格式无效";
|
||||
}
|
||||
|
||||
if (!Number.isInteger(probePort) || probePort < 1 || probePort > 65535) {
|
||||
errors.probeTargetPort = "端口必须是 1-65535";
|
||||
}
|
||||
}
|
||||
|
||||
if (form.type === 2) {
|
||||
if (!form.outNodeId || form.outNodeId.length === 0) {
|
||||
errors.outNodeId = "请至少选择一个出口节点";
|
||||
|
||||
@@ -160,7 +160,6 @@ export function ModalContent({
|
||||
showCloseButton={false}
|
||||
{...props}
|
||||
>
|
||||
<DialogTitle className="sr-only">Modal Dialog</DialogTitle>
|
||||
{renderedChildren}
|
||||
</BaseDialogContent>
|
||||
);
|
||||
@@ -173,15 +172,17 @@ export function ModalHeader({
|
||||
const context = useModalContext();
|
||||
|
||||
return (
|
||||
<div
|
||||
className={cn(
|
||||
"text-lg font-semibold",
|
||||
context?.classNames?.header,
|
||||
className,
|
||||
)}
|
||||
data-slot="modal-header"
|
||||
{...props}
|
||||
/>
|
||||
<DialogTitle asChild>
|
||||
<div
|
||||
className={cn(
|
||||
"text-lg font-semibold",
|
||||
context?.classNames?.header,
|
||||
className,
|
||||
)}
|
||||
data-slot="modal-header"
|
||||
{...props}
|
||||
/>
|
||||
</DialogTitle>
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -23,6 +23,7 @@ interface ClassNameMap {
|
||||
}
|
||||
|
||||
export interface SelectProps<T = unknown> extends FieldMetaProps {
|
||||
"aria-label"?: string;
|
||||
children?: React.ReactNode | ((item: T) => React.ReactNode);
|
||||
className?: string;
|
||||
classNames?: ClassNameMap;
|
||||
@@ -146,6 +147,7 @@ function textSizeClass(size: SelectProps["size"]) {
|
||||
}
|
||||
|
||||
export function Select<T>({
|
||||
"aria-label": ariaLabel,
|
||||
children,
|
||||
className,
|
||||
classNames,
|
||||
@@ -365,6 +367,7 @@ export function Select<T>({
|
||||
aria-controls={`${generatedId}-listbox`}
|
||||
aria-expanded={isExpanded}
|
||||
aria-haspopup="listbox"
|
||||
aria-label={label ? undefined : ariaLabel}
|
||||
className={cn(
|
||||
"flex w-full min-w-0 items-center gap-2 overflow-hidden rounded-md border border-input bg-background px-3 py-2 text-left shadow-sm focus:outline-none focus-visible:ring-2 focus-visible:ring-ring",
|
||||
isDisabled ? "cursor-not-allowed opacity-60" : "",
|
||||
@@ -398,6 +401,7 @@ export function Select<T>({
|
||||
</div>
|
||||
) : (
|
||||
<select
|
||||
aria-label={label ? undefined : ariaLabel}
|
||||
className={cn(
|
||||
"w-full rounded-md border border-input bg-background px-3 py-2 text-foreground shadow-sm focus:outline-none focus-visible:ring-2 focus-visible:ring-ring dark:[color-scheme:dark]",
|
||||
sizeClass(size),
|
||||
|
||||
Reference in New Issue
Block a user