mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-28 07:36:38 +08:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ae370382d3 | |||
| e112d81697 | |||
| 11a27d3c67 |
+10
-1
@@ -64,6 +64,14 @@ curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/panel_instal
|
||||
curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh -o install.sh && chmod +x install.sh && ./install.sh
|
||||
```
|
||||
|
||||
Alpine Linux 最小化安装若未包含 `curl`,可使用系统自带的 `wget` 下载:
|
||||
|
||||
```bash
|
||||
wget -O install.sh https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh && chmod +x install.sh && ./install.sh
|
||||
```
|
||||
|
||||
脚本会在 Alpine 上自动安装 Bash,并使用 OpenRC 注册、启动和管理 `flux_agent` 服务;其他受支持的 Linux 发行版继续使用 systemd。
|
||||
|
||||
**安装过程中会提示输入:**
|
||||
- **服务器地址**: 面板端的通信地址(通常是 `http://<面板IP>:<后端端口>`,例如 `http://1.2.3.4:6365`)。
|
||||
- **密钥**: 刚才在面板中获取的节点密钥。
|
||||
@@ -77,7 +85,8 @@ curl -L https://raw.githubusercontent.com/Sagit-chu/flux-panel/main/install.sh -
|
||||
|
||||
### 3. 验证安装
|
||||
安装完成后,服务会自动启动。
|
||||
- 查看状态: `systemctl status flux_agent`
|
||||
- systemd 查看状态: `systemctl status flux_agent`
|
||||
- Alpine/OpenRC 查看状态: `rc-service flux_agent status`
|
||||
- 回到面板 **节点管理** 页面,该节点状态应显示为 **在线**。
|
||||
|
||||
---
|
||||
|
||||
@@ -0,0 +1,650 @@
|
||||
# Forward Flow Reset Implementation Plan
|
||||
|
||||
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
|
||||
|
||||
**Goal:** Add a permission-checked action that resets only one forward rule's displayed upload and download counters.
|
||||
|
||||
**Architecture:** A dedicated repository method updates only the selected `forward` row. A dedicated authenticated handler reuses `resolveForwardAccess`, and the React page calls the endpoint from all three rule views through one confirmation modal.
|
||||
|
||||
**Tech Stack:** Go `net/http`, GORM, SQLite/PostgreSQL-compatible models, React, TypeScript, shadcn bridge components, Tailwind CSS v4.
|
||||
|
||||
## Global Constraints
|
||||
|
||||
- Only `forward.in_flow`, `forward.out_flow`, and `forward.updated_time` may change during reset.
|
||||
- Do not modify `user`, `user_tunnel`, quota, historical statistics, nftables counter state, or running services.
|
||||
- Administrators may reset any rule; non-admin users may reset only their own rules through existing `resolveForwardAccess` behavior.
|
||||
- All API responses must keep the `{code, msg, data, ts}` envelope.
|
||||
- Frontend imports must use `src/shadcn-bridge/heroui/*`; do not add `@heroui/*` or `@nextui-org/*` dependencies.
|
||||
- Do not add frontend test infrastructure.
|
||||
- Do not edit generated protobuf files, `install.sh`, or `panel_install.sh`.
|
||||
|
||||
---
|
||||
|
||||
### Task 1: Add the repository flow-reset primitive
|
||||
|
||||
**Files:**
|
||||
- Create: `go-backend/internal/store/repo/repository_forward_flow_reset_test.go`
|
||||
- Modify: `go-backend/internal/store/repo/repository_mutations.go`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: `model.Forward`, the repository's GORM database handle, and an explicit Unix-millisecond timestamp.
|
||||
- Produces: `func (r *Repository) ResetForwardFlow(forwardID int64, now int64) error`.
|
||||
|
||||
- [ ] **Step 1: Write the failing repository tests**
|
||||
|
||||
Create `go-backend/internal/store/repo/repository_forward_flow_reset_test.go`:
|
||||
|
||||
```go
|
||||
package repo
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestResetForwardFlowOnlyUpdatesSelectedForward(t *testing.T) {
|
||||
r, err := Open(filepath.Join(t.TempDir(), "forward-flow-reset.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open repo: %v", err)
|
||||
}
|
||||
defer r.Close()
|
||||
|
||||
const originalUpdated int64 = 1000
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status)
|
||||
VALUES(2, 'owner', 'pwd', 1, 0, 100, 700, 900, 0, 10, 1000, 1000, 1)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert user: %v", err)
|
||||
}
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx)
|
||||
VALUES(1, 'tunnel', 1, 1, 'tls', 1, 1000, 1000, 1, NULL, 0)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert tunnel: %v", err)
|
||||
}
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO user_tunnel(id, user_id, tunnel_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status)
|
||||
VALUES(10, 2, 1, 10, 100, 500, 600, 0, 0, 1)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert user tunnel: %v", err)
|
||||
}
|
||||
if err := r.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)
|
||||
VALUES
|
||||
(20, 2, 'owner', 'target', 1, '127.0.0.1:80', 'fifo', 111, 222, 1000, ?, 1, 0),
|
||||
(21, 2, 'owner', 'other', 1, '127.0.0.1:81', 'fifo', 333, 444, 1000, ?, 1, 1)
|
||||
`, originalUpdated, originalUpdated).Error; err != nil {
|
||||
t.Fatalf("insert forwards: %v", err)
|
||||
}
|
||||
|
||||
const resetAt int64 = 2000
|
||||
if err := r.ResetForwardFlow(20, resetAt); err != nil {
|
||||
t.Fatalf("ResetForwardFlow: %v", err)
|
||||
}
|
||||
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM forward WHERE id = 20", 0)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM forward WHERE id = 20", 0)
|
||||
assertForwardFlowResetValue(t, r, "SELECT updated_time FROM forward WHERE id = 20", resetAt)
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM forward WHERE id = 21", 333)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM forward WHERE id = 21", 444)
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM user WHERE id = 2", 700)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM user WHERE id = 2", 900)
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM user_tunnel WHERE id = 10", 500)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM user_tunnel WHERE id = 10", 600)
|
||||
}
|
||||
|
||||
func TestResetForwardFlowRejectsUninitializedRepository(t *testing.T) {
|
||||
var r *Repository
|
||||
if err := r.ResetForwardFlow(20, 2000); err == nil {
|
||||
t.Fatal("expected uninitialized repository error")
|
||||
}
|
||||
}
|
||||
|
||||
func assertForwardFlowResetValue(t *testing.T, r *Repository, query string, want int64) {
|
||||
t.Helper()
|
||||
var got int64
|
||||
if err := r.DB().Raw(query).Scan(&got).Error; err != nil {
|
||||
t.Fatalf("query %q: %v", query, err)
|
||||
}
|
||||
if got != want {
|
||||
t.Fatalf("query %q returned %d, want %d", query, got, want)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Run the repository tests and verify the missing method failure**
|
||||
|
||||
Run:
|
||||
|
||||
```bash
|
||||
cd go-backend && go test ./internal/store/repo -run TestResetForwardFlow -count=1
|
||||
```
|
||||
|
||||
Expected: compilation fails because `ResetForwardFlow` is undefined.
|
||||
|
||||
- [ ] **Step 3: Implement the minimal repository method**
|
||||
|
||||
Add to the flow-reset section of `go-backend/internal/store/repo/repository_mutations.go`:
|
||||
|
||||
```go
|
||||
func (r *Repository) ResetForwardFlow(forwardID int64, now int64) error {
|
||||
if r == nil || r.db == nil {
|
||||
return errors.New("repository not initialized")
|
||||
}
|
||||
return r.db.Model(&model.Forward{}).
|
||||
Where("id = ?", forwardID).
|
||||
Updates(map[string]interface{}{
|
||||
"in_flow": 0,
|
||||
"out_flow": 0,
|
||||
"updated_time": now,
|
||||
}).Error
|
||||
}
|
||||
```
|
||||
|
||||
The file already imports `errors` and `model`; do not add a new dependency.
|
||||
|
||||
- [ ] **Step 4: Format and run the focused repository tests**
|
||||
|
||||
Run:
|
||||
|
||||
```bash
|
||||
cd go-backend && gofmt -w internal/store/repo/repository_forward_flow_reset_test.go internal/store/repo/repository_mutations.go
|
||||
go test ./internal/store/repo -run TestResetForwardFlow -count=1
|
||||
```
|
||||
|
||||
Expected: both reset tests pass.
|
||||
|
||||
- [ ] **Step 5: Commit the repository change**
|
||||
|
||||
```bash
|
||||
git add go-backend/internal/store/repo/repository_mutations.go go-backend/internal/store/repo/repository_forward_flow_reset_test.go
|
||||
git commit -m "feat: add forward flow reset repository method"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 2: Add the authenticated reset endpoint
|
||||
|
||||
**Files:**
|
||||
- Create: `go-backend/internal/http/handler/forward_reset_flow_test.go`
|
||||
- Modify: `go-backend/internal/http/handler/handler.go`
|
||||
- Modify: `go-backend/internal/http/handler/mutations.go`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: `POST` JSON `{ "id": number }`, `resolveForwardAccess`, and `Repository.ResetForwardFlow` from Task 1.
|
||||
- Produces: `POST /api/v1/forward/reset-flow` and `func (h *Handler) forwardResetFlow(http.ResponseWriter, *http.Request)`.
|
||||
|
||||
- [ ] **Step 1: Write the failing handler tests**
|
||||
|
||||
Create `go-backend/internal/http/handler/forward_reset_flow_test.go`:
|
||||
|
||||
```go
|
||||
package handler
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"testing"
|
||||
|
||||
"go-backend/internal/auth"
|
||||
"go-backend/internal/http/middleware"
|
||||
"go-backend/internal/store/repo"
|
||||
)
|
||||
|
||||
func TestForwardResetFlowPermissionsAndIsolation(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
actorID int64
|
||||
actorRole int
|
||||
forwardID int64
|
||||
wantCode int
|
||||
wantInFlow int64
|
||||
wantOutFlow int64
|
||||
}{
|
||||
{name: "admin resets another user's rule", actorID: 1, actorRole: 0, forwardID: 20, wantCode: 0, wantInFlow: 0, wantOutFlow: 0},
|
||||
{name: "owner resets own rule", actorID: 2, actorRole: 1, forwardID: 20, wantCode: 0, wantInFlow: 0, wantOutFlow: 0},
|
||||
{name: "user cannot reset another user's rule", actorID: 3, actorRole: 1, forwardID: 20, wantCode: -1, wantInFlow: 111, wantOutFlow: 222},
|
||||
{name: "missing rule is rejected", actorID: 1, actorRole: 0, forwardID: 999, wantCode: -1, wantInFlow: 111, wantOutFlow: 222},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
h, r := setupForwardResetFlowHandler(t)
|
||||
req := newForwardResetFlowRequest(t, http.MethodPost, tt.forwardID, tt.actorID, tt.actorRole)
|
||||
res := httptest.NewRecorder()
|
||||
|
||||
h.forwardResetFlow(res, req)
|
||||
|
||||
if got := decodeForwardResetFlowCode(t, res); got != tt.wantCode {
|
||||
t.Fatalf("code = %d, want %d; body=%s", got, tt.wantCode, res.Body.String())
|
||||
}
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT in_flow FROM forward WHERE id = 20", tt.wantInFlow)
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT out_flow FROM forward WHERE id = 20", tt.wantOutFlow)
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT in_flow FROM user WHERE id = 2", 700)
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT out_flow FROM user_tunnel WHERE id = 10", 600)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestForwardResetFlowRejectsInvalidRequests(t *testing.T) {
|
||||
h, _ := setupForwardResetFlowHandler(t)
|
||||
|
||||
t.Run("non post", func(t *testing.T) {
|
||||
req := newForwardResetFlowRequest(t, http.MethodGet, 20, 1, 0)
|
||||
res := httptest.NewRecorder()
|
||||
h.forwardResetFlow(res, req)
|
||||
if code := decodeForwardResetFlowCode(t, res); code != -1 {
|
||||
t.Fatalf("code = %d, want -1", code)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("invalid id", func(t *testing.T) {
|
||||
req := newForwardResetFlowRequest(t, http.MethodPost, 0, 1, 0)
|
||||
res := httptest.NewRecorder()
|
||||
h.forwardResetFlow(res, req)
|
||||
if code := decodeForwardResetFlowCode(t, res); code != -1 {
|
||||
t.Fatalf("code = %d, want -1", code)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func setupForwardResetFlowHandler(t *testing.T) (*Handler, *repo.Repository) {
|
||||
t.Helper()
|
||||
r, err := repo.Open(filepath.Join(t.TempDir(), "forward-reset-handler.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open repo: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = r.Close() })
|
||||
|
||||
statements := []string{
|
||||
`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(1, 'admin', 'pwd', 0, 0, 100, 0, 0, 0, 10, 1000, 1000, 1)`,
|
||||
`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(2, 'owner', 'pwd', 1, 0, 100, 700, 900, 0, 10, 1000, 1000, 1)`,
|
||||
`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(3, 'other', 'pwd', 1, 0, 100, 0, 0, 0, 10, 1000, 1000, 1)`,
|
||||
`INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) VALUES(1, 'tunnel', 1, 1, 'tls', 1, 1000, 1000, 1, NULL, 0)`,
|
||||
`INSERT INTO user_tunnel(id, user_id, tunnel_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status) VALUES(10, 2, 1, 10, 100, 500, 600, 0, 0, 1)`,
|
||||
`INSERT INTO forward(id, user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) VALUES(20, 2, 'owner', 'target', 1, '127.0.0.1:80', 'fifo', 111, 222, 1000, 1000, 1, 0)`,
|
||||
}
|
||||
for _, statement := range statements {
|
||||
if err := r.DB().Exec(statement).Error; err != nil {
|
||||
t.Fatalf("seed database: %v", err)
|
||||
}
|
||||
}
|
||||
return New(r, "test-secret"), r
|
||||
}
|
||||
|
||||
func newForwardResetFlowRequest(t *testing.T, method string, forwardID, actorID int64, roleID int) *http.Request {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(map[string]int64{"id": forwardID})
|
||||
if err != nil {
|
||||
t.Fatalf("marshal request: %v", err)
|
||||
}
|
||||
req := httptest.NewRequest(method, "/api/v1/forward/reset-flow", bytes.NewReader(body))
|
||||
claims := auth.Claims{Sub: strconv.FormatInt(actorID, 10), RoleID: roleID}
|
||||
return req.WithContext(context.WithValue(req.Context(), middleware.ClaimsContextKey, claims))
|
||||
}
|
||||
|
||||
func decodeForwardResetFlowCode(t *testing.T, res *httptest.ResponseRecorder) int {
|
||||
t.Helper()
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
}
|
||||
if err := json.Unmarshal(res.Body.Bytes(), &payload); err != nil {
|
||||
t.Fatalf("decode response: %v; body=%s", err, res.Body.String())
|
||||
}
|
||||
return payload.Code
|
||||
}
|
||||
|
||||
func assertForwardResetFlowDBValue(t *testing.T, r *repo.Repository, query string, want int64) {
|
||||
t.Helper()
|
||||
var got int64
|
||||
if err := r.DB().Raw(query).Scan(&got).Error; err != nil {
|
||||
t.Fatalf("query %q: %v", query, err)
|
||||
}
|
||||
if got != want {
|
||||
t.Fatalf("query %q returned %d, want %d", query, got, want)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
If the project's default error code differs from `-1`, replace the test expectation with the actual `response.ErrDefault` code after inspecting one existing handler response; do not weaken the success and database assertions.
|
||||
|
||||
- [ ] **Step 2: Run the handler tests and verify the missing handler failure**
|
||||
|
||||
Run:
|
||||
|
||||
```bash
|
||||
cd go-backend && go test ./internal/http/handler -run TestForwardResetFlow -count=1
|
||||
```
|
||||
|
||||
Expected: compilation fails because `forwardResetFlow` is undefined.
|
||||
|
||||
- [ ] **Step 3: Register and implement the endpoint**
|
||||
|
||||
Add this route beside the other forward routes in `go-backend/internal/http/handler/handler.go`:
|
||||
|
||||
```go
|
||||
mux.HandleFunc("/api/v1/forward/reset-flow", h.forwardResetFlow)
|
||||
```
|
||||
|
||||
Add this handler beside `forwardPause` and `forwardResume` in `go-backend/internal/http/handler/mutations.go`:
|
||||
|
||||
```go
|
||||
func (h *Handler) forwardResetFlow(w http.ResponseWriter, r *http.Request) {
|
||||
id := idFromBody(r, w)
|
||||
if id <= 0 {
|
||||
return
|
||||
}
|
||||
if _, _, _, err := h.resolveForwardAccess(r, id); err != nil {
|
||||
if errors.Is(err, errForwardNotFound) {
|
||||
response.WriteJSON(w, response.ErrDefault("转发不存在"))
|
||||
return
|
||||
}
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
if err := h.repo.ResetForwardFlow(id, time.Now().UnixMilli()); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
```
|
||||
|
||||
This deliberately does not call runtime service controls or nftables reconciliation.
|
||||
|
||||
- [ ] **Step 4: Format and run the focused handler tests**
|
||||
|
||||
Run:
|
||||
|
||||
```bash
|
||||
cd go-backend && gofmt -w internal/http/handler/forward_reset_flow_test.go internal/http/handler/handler.go internal/http/handler/mutations.go
|
||||
go test ./internal/http/handler -run TestForwardResetFlow -count=1
|
||||
```
|
||||
|
||||
Expected: all reset endpoint tests pass.
|
||||
|
||||
- [ ] **Step 5: Run all backend tests**
|
||||
|
||||
Run:
|
||||
|
||||
```bash
|
||||
cd go-backend && go test ./...
|
||||
```
|
||||
|
||||
Expected: all backend packages and contract tests pass, excluding environment-gated PostgreSQL tests when `FLVX_POSTGRES_TEST_DSN` is unset.
|
||||
|
||||
- [ ] **Step 6: Commit the endpoint change**
|
||||
|
||||
```bash
|
||||
git add go-backend/internal/http/handler/handler.go go-backend/internal/http/handler/mutations.go go-backend/internal/http/handler/forward_reset_flow_test.go
|
||||
git commit -m "feat: add forward flow reset endpoint"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 3: Add the rule-page reset action and confirmation modal
|
||||
|
||||
**Files:**
|
||||
- Modify: `vite-frontend/src/api/index.ts`
|
||||
- Modify: `vite-frontend/src/pages/forward.tsx`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: `POST /forward/reset-flow`, the page's `Forward` shape, `refreshForwardList`, toast notifications, and existing modal/button bridge components.
|
||||
- Produces: `resetForwardFlow(id: number)`, a shared reset handler, disabled zero-usage actions in all rule views, and one confirmation modal.
|
||||
|
||||
- [ ] **Step 1: Add the frontend API wrapper**
|
||||
|
||||
Add beside the forward control operations in `vite-frontend/src/api/index.ts`:
|
||||
|
||||
```ts
|
||||
export const resetForwardFlow = (forwardId: number) =>
|
||||
Network.post("/forward/reset-flow", { id: forwardId });
|
||||
```
|
||||
|
||||
Import `resetForwardFlow` from `@/api` in `vite-frontend/src/pages/forward.tsx`.
|
||||
|
||||
- [ ] **Step 2: Add page state and shared reset handlers**
|
||||
|
||||
Add state beside the existing delete modal state:
|
||||
|
||||
```ts
|
||||
const [resetFlowModalOpen, setResetFlowModalOpen] = useState(false);
|
||||
const [resetFlowLoading, setResetFlowLoading] = useState(false);
|
||||
const [forwardToResetFlow, setForwardToResetFlow] = useState<Forward | null>(null);
|
||||
```
|
||||
|
||||
Add these handlers beside `handleDelete` and `confirmDelete`:
|
||||
|
||||
```ts
|
||||
const handleResetFlow = (forward: Forward) => {
|
||||
if ((forward.inFlow || 0) + (forward.outFlow || 0) <= 0) return;
|
||||
setForwardToResetFlow(forward);
|
||||
setResetFlowModalOpen(true);
|
||||
};
|
||||
|
||||
const confirmResetFlow = async () => {
|
||||
if (!forwardToResetFlow) return;
|
||||
|
||||
setResetFlowLoading(true);
|
||||
try {
|
||||
const res = await resetForwardFlow(forwardToResetFlow.id);
|
||||
|
||||
if (res.code !== 0) {
|
||||
toast.error(res.msg || "流量清零失败");
|
||||
return;
|
||||
}
|
||||
|
||||
toast.success("规则流量已清零");
|
||||
setResetFlowModalOpen(false);
|
||||
setForwardToResetFlow(null);
|
||||
await refreshForwardList(false);
|
||||
} catch {
|
||||
toast.error("流量清零失败");
|
||||
} finally {
|
||||
setResetFlowLoading(false);
|
||||
}
|
||||
};
|
||||
```
|
||||
|
||||
- [ ] **Step 3: Add one reusable reset icon button to both table row components**
|
||||
|
||||
Pass `handleResetFlow` into `SortableTableRow` and `SortableCompactTableRow` at every render site. Add it to each component's destructured props.
|
||||
|
||||
Insert this button between diagnosis and delete in each table action cell:
|
||||
|
||||
```tsx
|
||||
<Button
|
||||
isIconOnly
|
||||
className="bg-secondary/10 text-secondary hover:bg-secondary/20"
|
||||
isDisabled={(forward.inFlow || 0) + (forward.outFlow || 0) <= 0}
|
||||
size="sm"
|
||||
title="流量清零"
|
||||
onPress={() => handleResetFlow(forward)}
|
||||
>
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
className="h-4 w-4"
|
||||
fill="none"
|
||||
stroke="currentColor"
|
||||
viewBox="0 0 24 24"
|
||||
>
|
||||
<path
|
||||
d="M4 4v6h6M20 20v-6h-6M20 9a8 8 0 00-13.657-3.657L4 8m16 8-2.343 2.657A8 8 0 014 15"
|
||||
strokeLinecap="round"
|
||||
strokeLinejoin="round"
|
||||
strokeWidth={2}
|
||||
/>
|
||||
</svg>
|
||||
</Button>
|
||||
```
|
||||
|
||||
- [ ] **Step 4: Add the reset action to the card view**
|
||||
|
||||
Insert a fourth action button between diagnosis and delete in `renderForwardCard`:
|
||||
|
||||
```tsx
|
||||
<Button
|
||||
className="flex-1 min-h-8"
|
||||
color="secondary"
|
||||
isDisabled={(forward.inFlow || 0) + (forward.outFlow || 0) <= 0}
|
||||
size="sm"
|
||||
startContent={
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
className="w-3 h-3"
|
||||
fill="none"
|
||||
stroke="currentColor"
|
||||
viewBox="0 0 24 24"
|
||||
>
|
||||
<path
|
||||
d="M4 4v6h6M20 20v-6h-6M20 9a8 8 0 00-13.657-3.657L4 8m16 8-2.343 2.657A8 8 0 014 15"
|
||||
strokeLinecap="round"
|
||||
strokeLinejoin="round"
|
||||
strokeWidth={2}
|
||||
/>
|
||||
</svg>
|
||||
}
|
||||
variant="flat"
|
||||
onPress={() => handleResetFlow(forward)}
|
||||
>
|
||||
清零
|
||||
</Button>
|
||||
```
|
||||
|
||||
Change the card action container from `flex gap-1.5 mt-3` to `grid grid-cols-2 gap-1.5 mt-3` so all four actions remain readable at the smallest supported card width.
|
||||
|
||||
- [ ] **Step 5: Add the confirmation modal**
|
||||
|
||||
Add beside the delete confirmation modal:
|
||||
|
||||
```tsx
|
||||
<Modal
|
||||
backdrop="blur"
|
||||
classNames={{
|
||||
base: "!w-[calc(100%-32px)] !mx-auto sm:!w-full rounded-2xl overflow-hidden",
|
||||
}}
|
||||
isOpen={resetFlowModalOpen}
|
||||
placement="center"
|
||||
scrollBehavior="inside"
|
||||
size="lg"
|
||||
onOpenChange={setResetFlowModalOpen}
|
||||
>
|
||||
<ModalContent>
|
||||
{(onClose) => (
|
||||
<>
|
||||
<ModalHeader className="flex flex-col gap-1">
|
||||
<h2 className="text-lg font-bold text-secondary">确认流量清零</h2>
|
||||
</ModalHeader>
|
||||
<ModalBody>
|
||||
<p className="text-default-600">
|
||||
确定要清零规则{" "}
|
||||
<span className="font-semibold text-foreground">
|
||||
"{forwardToResetFlow?.name}"
|
||||
</span>{" "}
|
||||
当前显示的上传和下载流量吗?
|
||||
</p>
|
||||
<p className="text-small text-default-500 mt-2">
|
||||
此操作不可撤销,但不会影响用户总流量、用户隧道配额和历史统计。
|
||||
</p>
|
||||
</ModalBody>
|
||||
<ModalFooter>
|
||||
<Button isDisabled={resetFlowLoading} variant="light" onPress={onClose}>
|
||||
取消
|
||||
</Button>
|
||||
<Button
|
||||
color="secondary"
|
||||
isLoading={resetFlowLoading}
|
||||
onPress={confirmResetFlow}
|
||||
>
|
||||
确认清零
|
||||
</Button>
|
||||
</ModalFooter>
|
||||
</>
|
||||
)}
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
```
|
||||
|
||||
Add this wrapper beside the other reset handlers and pass it to the modal as `onOpenChange={handleResetFlowModalOpenChange}`:
|
||||
|
||||
```ts
|
||||
const handleResetFlowModalOpenChange = (isOpen: boolean) => {
|
||||
if (resetFlowLoading) return;
|
||||
setResetFlowModalOpen(isOpen);
|
||||
if (!isOpen) {
|
||||
setForwardToResetFlow(null);
|
||||
}
|
||||
};
|
||||
```
|
||||
|
||||
- [ ] **Step 6: Format and verify the frontend**
|
||||
|
||||
Run:
|
||||
|
||||
```bash
|
||||
cd vite-frontend && pnpm exec prettier --write src/api/index.ts src/pages/forward.tsx
|
||||
pnpm run build
|
||||
pnpm run lint
|
||||
```
|
||||
|
||||
Expected: TypeScript/Vite build succeeds and ESLint finishes without errors.
|
||||
|
||||
- [ ] **Step 7: Commit the frontend change**
|
||||
|
||||
```bash
|
||||
git add vite-frontend/src/api/index.ts vite-frontend/src/pages/forward.tsx
|
||||
git commit -m "feat: add forward flow reset action"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 4: Perform integrated verification
|
||||
|
||||
**Files:**
|
||||
- Verify only; no planned source changes.
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: the repository method, API endpoint, and rule-page action from Tasks 1-3.
|
||||
- Produces: evidence that the complete feature builds and all affected tests pass.
|
||||
|
||||
- [ ] **Step 1: Run the complete backend suite**
|
||||
|
||||
```bash
|
||||
cd go-backend && go test ./...
|
||||
```
|
||||
|
||||
Expected: all available backend tests pass.
|
||||
|
||||
- [ ] **Step 2: Run the complete frontend checks**
|
||||
|
||||
```bash
|
||||
cd vite-frontend && pnpm run build && pnpm run lint
|
||||
```
|
||||
|
||||
Expected: both commands exit successfully.
|
||||
|
||||
- [ ] **Step 3: Check formatting and working-tree scope**
|
||||
|
||||
```bash
|
||||
git diff --check
|
||||
git status --short
|
||||
git log -4 --oneline
|
||||
```
|
||||
|
||||
Expected: no whitespace errors; the working tree is clean; the three feature commits are visible after the design and implementation-plan commits.
|
||||
|
||||
- [ ] **Step 4: Manually verify the feature when a local panel is available**
|
||||
|
||||
1. Open the Rules page as an administrator and reset a rule with non-zero upload/download traffic.
|
||||
2. Confirm the modal states that user totals, tunnel quota, and history are unaffected.
|
||||
3. Confirm the rule immediately shows zero after success.
|
||||
4. Confirm the user page's total traffic and user-tunnel traffic values did not change.
|
||||
5. Generate new traffic and confirm the rule starts accumulating from zero.
|
||||
6. Log in as a normal user and confirm the user can reset an owned rule but cannot access another user's rule through a direct API request.
|
||||
|
||||
Expected: all six checks match the design specification.
|
||||
@@ -0,0 +1,197 @@
|
||||
# 规则流量清零设计
|
||||
|
||||
## 背景
|
||||
|
||||
Issue #523 希望“规则”页面中每条隧道规则显示的流量使用量支持手动清零。
|
||||
|
||||
当前规则流量保存在 `forward.in_flow` 和 `forward.out_flow`。流量上报时,同一份增量还会累计到用户总流量、用户隧道流量和相关配额统计中。因此,本功能必须将“规则展示计数器清零”与“用户或隧道配额重置”严格区分。
|
||||
|
||||
## 目标
|
||||
|
||||
为单条规则提供手动流量清零能力:
|
||||
|
||||
- 将所选规则的上传流量和下载流量清零。
|
||||
- 管理员可以清零任意规则。
|
||||
- 普通用户只能清零自己的规则。
|
||||
- 清零后,新产生的流量继续从零正常累计。
|
||||
|
||||
## 非目标
|
||||
|
||||
本功能不会:
|
||||
|
||||
- 修改用户总流量 `user.in_flow` 或 `user.out_flow`。
|
||||
- 修改用户隧道流量 `user_tunnel.in_flow` 或 `user_tunnel.out_flow`。
|
||||
- 修改每日或每月配额用量。
|
||||
- 修改历史流量统计。
|
||||
- 重置 nftables 节点计数器或其增量计算基线。
|
||||
- 重启、暂停、恢复或重新部署规则服务。
|
||||
- 增加批量流量清零功能。
|
||||
|
||||
## 后端设计
|
||||
|
||||
### API
|
||||
|
||||
新增接口:
|
||||
|
||||
```text
|
||||
POST /api/v1/forward/reset-flow
|
||||
```
|
||||
|
||||
请求体:
|
||||
|
||||
```json
|
||||
{
|
||||
"id": 123
|
||||
}
|
||||
```
|
||||
|
||||
成功响应沿用统一 envelope:
|
||||
|
||||
```json
|
||||
{
|
||||
"code": 0,
|
||||
"msg": "success",
|
||||
"data": null,
|
||||
"ts": 0
|
||||
}
|
||||
```
|
||||
|
||||
具体 `msg`、`data` 和 `ts` 值继续由现有 response helper 生成。
|
||||
|
||||
### 参数与权限校验
|
||||
|
||||
Handler 执行以下步骤:
|
||||
|
||||
1. 只接受 `POST` 请求。
|
||||
2. 从 JSON 请求体读取正整数规则 ID。
|
||||
3. 调用现有 `resolveForwardAccess`:
|
||||
- 管理员角色可以访问任意存在的规则。
|
||||
- 普通用户仅能访问 `forward.user_id` 等于当前用户 ID 的规则。
|
||||
- 对普通用户访问他人规则的情况,沿用现有逻辑返回“转发不存在”,避免暴露规则存在性。
|
||||
4. 调用 Repository 完成清零。
|
||||
5. 返回统一成功响应。
|
||||
|
||||
### Repository
|
||||
|
||||
新增方法:
|
||||
|
||||
```go
|
||||
func (r *Repository) ResetForwardFlow(forwardID int64, now int64) error
|
||||
```
|
||||
|
||||
该方法只更新指定 `forward` 记录:
|
||||
|
||||
```text
|
||||
in_flow = 0
|
||||
out_flow = 0
|
||||
updated_time = now
|
||||
```
|
||||
|
||||
Repository 不直接操作 Handler 的身份信息,也不更新任何其他表。
|
||||
|
||||
### 并发与后续流量
|
||||
|
||||
清零使用单条 SQL `UPDATE`。agent 流量上报和 nftables 流量采集仍使用原有增量累加逻辑。清零不会重置采集基线,因此下一次采集只会把清零之后新计算出的增量加回规则计数,不会把清零前的累计值整体恢复。
|
||||
|
||||
若清零 SQL 与流量增量 SQL 同时执行,数据库按实际语句执行顺序决定最终值;每条更新本身保持原子性。本功能不引入暂停采集或跨节点同步流程。
|
||||
|
||||
## 前端设计
|
||||
|
||||
### API 封装
|
||||
|
||||
在 `vite-frontend/src/api/index.ts` 新增:
|
||||
|
||||
```ts
|
||||
export const resetForwardFlow = (id: number) =>
|
||||
Network.post("/forward/reset-flow", { id });
|
||||
```
|
||||
|
||||
### 入口
|
||||
|
||||
在规则页面所有单条规则操作入口中增加“流量清零”操作:
|
||||
|
||||
- 分组表格视图。
|
||||
- 精简表格视图。
|
||||
- 卡片视图。
|
||||
|
||||
按钮使用独立的清零/刷新语义图标和提示文本,不复用删除按钮样式。
|
||||
|
||||
当规则的 `inFlow + outFlow` 等于零时,按钮禁用,避免重复请求。
|
||||
|
||||
### 确认交互
|
||||
|
||||
点击按钮后打开确认弹窗,显示规则名称,并明确说明:
|
||||
|
||||
- 仅清零当前规则显示的上传和下载流量。
|
||||
- 不影响用户总流量、用户隧道配额和历史统计。
|
||||
- 操作不可撤销。
|
||||
|
||||
确认期间显示 loading 状态并阻止重复提交。
|
||||
|
||||
### 成功与失败
|
||||
|
||||
- 成功:关闭弹窗,显示成功 toast,并刷新规则列表。
|
||||
- 失败:保留弹窗,显示后端错误信息或通用失败 toast。
|
||||
- 刷新后,该规则上传和下载均显示为零;后续流量继续正常累计。
|
||||
|
||||
## 错误处理
|
||||
|
||||
- 非 POST 请求:返回现有通用请求失败响应。
|
||||
- 请求体无法解析、ID 缺失或 ID 非正数:返回“请求参数错误”。
|
||||
- 规则不存在或普通用户访问他人规则:返回“转发不存在”。
|
||||
- Repository 更新失败:返回包含 Repository 错误信息的统一错误响应。
|
||||
- 前端网络错误:显示“流量清零失败”。
|
||||
|
||||
## 测试策略
|
||||
|
||||
### Repository 测试
|
||||
|
||||
验证:
|
||||
|
||||
- 指定规则的 `in_flow`、`out_flow` 被清零。
|
||||
- 指定规则的 `updated_time` 被更新。
|
||||
- 其他规则的流量不变。
|
||||
- 用户总流量不变。
|
||||
- 用户隧道流量不变。
|
||||
- Repository 未初始化时返回错误。
|
||||
|
||||
### Handler 测试
|
||||
|
||||
验证:
|
||||
|
||||
- 管理员能够清零任意存在的规则。
|
||||
- 普通用户能够清零自己的规则。
|
||||
- 普通用户不能清零他人的规则。
|
||||
- 不存在的规则返回错误。
|
||||
- 无效 ID 返回参数错误。
|
||||
- 非 POST 请求返回请求失败。
|
||||
- 成功请求不修改用户和用户隧道流量。
|
||||
|
||||
### 前端验证
|
||||
|
||||
项目没有配置前端测试框架,因此不新增前端单元测试。使用以下命令验证:
|
||||
|
||||
```bash
|
||||
(cd vite-frontend && pnpm run build)
|
||||
(cd vite-frontend && pnpm run lint)
|
||||
```
|
||||
|
||||
后端使用:
|
||||
|
||||
```bash
|
||||
(cd go-backend && go test ./...)
|
||||
```
|
||||
|
||||
## 文件范围
|
||||
|
||||
预计修改:
|
||||
|
||||
- `go-backend/internal/http/handler/handler.go`
|
||||
- `go-backend/internal/http/handler/mutations.go`
|
||||
- `go-backend/internal/http/handler/*_test.go`
|
||||
- `go-backend/internal/store/repo/repository_mutations.go`
|
||||
- `go-backend/internal/store/repo/*_test.go`
|
||||
- `vite-frontend/src/api/index.ts`
|
||||
- `vite-frontend/src/pages/forward.tsx`
|
||||
|
||||
不需要数据库迁移或新增依赖。
|
||||
@@ -59,6 +59,7 @@ type diagnosisWorkItem struct {
|
||||
type diagnosisExecOptions struct {
|
||||
commandTimeout time.Duration
|
||||
pingTimeoutMS int
|
||||
pingCount int
|
||||
timeoutMessage string
|
||||
}
|
||||
|
||||
@@ -1596,10 +1597,14 @@ func (h *Handler) tcpPingViaNode(nodeID int64, ip string, port int, options diag
|
||||
if options.pingTimeoutMS <= 0 {
|
||||
options.pingTimeoutMS = int(diagnosisCommandTimeout / time.Millisecond)
|
||||
}
|
||||
pingCount := options.pingCount
|
||||
if pingCount <= 0 {
|
||||
pingCount = 4
|
||||
}
|
||||
res, err := h.sendNodeCommandWithTimeout(nodeID, "TcpPing", map[string]interface{}{
|
||||
"ip": ip,
|
||||
"port": port,
|
||||
"count": 4,
|
||||
"count": pingCount,
|
||||
"timeout": options.pingTimeoutMS,
|
||||
}, options.commandTimeout, false, false)
|
||||
if err != nil {
|
||||
@@ -1626,12 +1631,16 @@ func (h *Handler) tcpPingViaRemoteNode(node *nodeRecord, ip string, port int, op
|
||||
if options.pingTimeoutMS <= 0 {
|
||||
options.pingTimeoutMS = int(diagnosisCommandTimeout / time.Millisecond)
|
||||
}
|
||||
pingCount := options.pingCount
|
||||
if pingCount <= 0 {
|
||||
pingCount = 4
|
||||
}
|
||||
|
||||
fc := client.NewFederationClientWithTimeout(options.commandTimeout)
|
||||
return fc.Diagnose(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeDiagnoseRequest{
|
||||
IP: strings.TrimSpace(ip),
|
||||
Port: port,
|
||||
Count: 4,
|
||||
Count: pingCount,
|
||||
Timeout: options.pingTimeoutMS,
|
||||
Protocol: "tcp",
|
||||
})
|
||||
|
||||
@@ -0,0 +1,129 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"testing"
|
||||
|
||||
"go-backend/internal/auth"
|
||||
"go-backend/internal/http/middleware"
|
||||
"go-backend/internal/store/repo"
|
||||
)
|
||||
|
||||
func TestForwardResetFlowPermissionsAndIsolation(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
actorID int64
|
||||
actorRole int
|
||||
forwardID int64
|
||||
wantCode int
|
||||
wantInFlow int64
|
||||
wantOutFlow int64
|
||||
}{
|
||||
{name: "admin resets another user's rule", actorID: 1, actorRole: 0, forwardID: 20, wantCode: 0, wantInFlow: 0, wantOutFlow: 0},
|
||||
{name: "owner resets own rule", actorID: 2, actorRole: 1, forwardID: 20, wantCode: 0, wantInFlow: 0, wantOutFlow: 0},
|
||||
{name: "user cannot reset another user's rule", actorID: 3, actorRole: 1, forwardID: 20, wantCode: -1, wantInFlow: 111, wantOutFlow: 222},
|
||||
{name: "missing rule is rejected", actorID: 1, actorRole: 0, forwardID: 999, wantCode: -1, wantInFlow: 111, wantOutFlow: 222},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
h, r := setupForwardResetFlowHandler(t)
|
||||
req := newForwardResetFlowRequest(t, http.MethodPost, tt.forwardID, tt.actorID, tt.actorRole)
|
||||
res := httptest.NewRecorder()
|
||||
|
||||
h.forwardResetFlow(res, req)
|
||||
|
||||
if got := decodeForwardResetFlowCode(t, res); got != tt.wantCode {
|
||||
t.Fatalf("code = %d, want %d; body=%s", got, tt.wantCode, res.Body.String())
|
||||
}
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT in_flow FROM forward WHERE id = 20", tt.wantInFlow)
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT out_flow FROM forward WHERE id = 20", tt.wantOutFlow)
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT in_flow FROM user WHERE id = 2", 700)
|
||||
assertForwardResetFlowDBValue(t, r, "SELECT out_flow FROM user_tunnel WHERE id = 10", 600)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestForwardResetFlowRejectsInvalidRequests(t *testing.T) {
|
||||
h, _ := setupForwardResetFlowHandler(t)
|
||||
|
||||
t.Run("non post", func(t *testing.T) {
|
||||
req := newForwardResetFlowRequest(t, http.MethodGet, 20, 1, 0)
|
||||
res := httptest.NewRecorder()
|
||||
h.forwardResetFlow(res, req)
|
||||
if code := decodeForwardResetFlowCode(t, res); code != -1 {
|
||||
t.Fatalf("code = %d, want -1", code)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("invalid id", func(t *testing.T) {
|
||||
req := newForwardResetFlowRequest(t, http.MethodPost, 0, 1, 0)
|
||||
res := httptest.NewRecorder()
|
||||
h.forwardResetFlow(res, req)
|
||||
if code := decodeForwardResetFlowCode(t, res); code != -1 {
|
||||
t.Fatalf("code = %d, want -1", code)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func setupForwardResetFlowHandler(t *testing.T) (*Handler, *repo.Repository) {
|
||||
t.Helper()
|
||||
r, err := repo.Open(filepath.Join(t.TempDir(), "forward-reset-handler.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open repo: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = r.Close() })
|
||||
|
||||
statements := []string{
|
||||
`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(2, 'owner', 'pwd', 1, 0, 100, 700, 900, 0, 10, 1000, 1000, 1)`,
|
||||
`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(3, 'other', 'pwd', 1, 0, 100, 0, 0, 0, 10, 1000, 1000, 1)`,
|
||||
`INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) VALUES(1, 'tunnel', 1, 1, 'tls', 1, 1000, 1000, 1, NULL, 0)`,
|
||||
`INSERT INTO user_tunnel(id, user_id, tunnel_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status) VALUES(10, 2, 1, 10, 100, 500, 600, 0, 0, 1)`,
|
||||
`INSERT INTO forward(id, user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) VALUES(20, 2, 'owner', 'target', 1, '127.0.0.1:80', 'fifo', 111, 222, 1000, 1000, 1, 0)`,
|
||||
}
|
||||
for _, statement := range statements {
|
||||
if err := r.DB().Exec(statement).Error; err != nil {
|
||||
t.Fatalf("seed database: %v", err)
|
||||
}
|
||||
}
|
||||
return New(r, "test-secret"), r
|
||||
}
|
||||
|
||||
func newForwardResetFlowRequest(t *testing.T, method string, forwardID, actorID int64, roleID int) *http.Request {
|
||||
t.Helper()
|
||||
body, err := json.Marshal(map[string]int64{"id": forwardID})
|
||||
if err != nil {
|
||||
t.Fatalf("marshal request: %v", err)
|
||||
}
|
||||
req := httptest.NewRequest(method, "/api/v1/forward/reset-flow", bytes.NewReader(body))
|
||||
claims := auth.Claims{Sub: strconv.FormatInt(actorID, 10), RoleID: roleID}
|
||||
return req.WithContext(context.WithValue(req.Context(), middleware.ClaimsContextKey, claims))
|
||||
}
|
||||
|
||||
func decodeForwardResetFlowCode(t *testing.T, res *httptest.ResponseRecorder) int {
|
||||
t.Helper()
|
||||
var payload struct {
|
||||
Code int `json:"code"`
|
||||
}
|
||||
if err := json.Unmarshal(res.Body.Bytes(), &payload); err != nil {
|
||||
t.Fatalf("decode response: %v; body=%s", err, res.Body.String())
|
||||
}
|
||||
return payload.Code
|
||||
}
|
||||
|
||||
func assertForwardResetFlowDBValue(t *testing.T, r *repo.Repository, query string, want int64) {
|
||||
t.Helper()
|
||||
var got int64
|
||||
if err := r.DB().Raw(query).Scan(&got).Error; err != nil {
|
||||
t.Fatalf("query %q: %v", query, err)
|
||||
}
|
||||
if got != want {
|
||||
t.Fatalf("query %q returned %d, want %d", query, got, want)
|
||||
}
|
||||
}
|
||||
@@ -221,6 +221,7 @@ func (h *Handler) Register(mux *http.ServeMux) {
|
||||
mux.HandleFunc("/api/v1/forward/force-delete", h.forwardForceDelete)
|
||||
mux.HandleFunc("/api/v1/forward/pause", h.forwardPause)
|
||||
mux.HandleFunc("/api/v1/forward/resume", h.forwardResume)
|
||||
mux.HandleFunc("/api/v1/forward/reset-flow", h.forwardResetFlow)
|
||||
mux.HandleFunc("/api/v1/forward/diagnose", h.forwardDiagnose)
|
||||
mux.HandleFunc("/api/v1/forward/diagnose/stream", h.forwardDiagnoseStream)
|
||||
mux.HandleFunc("/api/v1/forward/update-order", h.forwardUpdateOrder)
|
||||
@@ -1014,6 +1015,7 @@ func (h *Handler) updateConfigs(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
h.notifyTunnelQualityConfigChanged(key)
|
||||
}
|
||||
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
@@ -1061,6 +1063,7 @@ func (h *Handler) updateSingleConfig(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
h.notifyTunnelQualityConfigChanged(name)
|
||||
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
@@ -1109,11 +1112,23 @@ func normalizeAndValidateConfigValue(key, value string) (string, error) {
|
||||
}
|
||||
case monitoring.ConfigMonitorRetentionDays:
|
||||
return monitoring.NormalizeMonitoringRetentionDays(value)
|
||||
case monitoring.ConfigTunnelQualityProbeIntervalSec:
|
||||
return monitoring.NormalizeTunnelQualityProbeIntervalSeconds(value)
|
||||
default:
|
||||
return value, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) notifyTunnelQualityConfigChanged(key string) {
|
||||
if h == nil || h.qualityProber == nil {
|
||||
return
|
||||
}
|
||||
switch strings.TrimSpace(key) {
|
||||
case monitorTunnelQualityEnabledConfigKey, monitoring.ConfigTunnelQualityProbeIntervalSec:
|
||||
h.qualityProber.NotifyConfigChanged()
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) isTunnelQualityMonitoringEnabled() bool {
|
||||
if h == nil || h.repo == nil {
|
||||
return true
|
||||
|
||||
@@ -2502,6 +2502,26 @@ func (h *Handler) forwardPause(w http.ResponseWriter, r *http.Request) {
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
|
||||
func (h *Handler) forwardResetFlow(w http.ResponseWriter, r *http.Request) {
|
||||
id := idFromBody(r, w)
|
||||
if id <= 0 {
|
||||
return
|
||||
}
|
||||
if _, _, _, err := h.resolveForwardAccess(r, id); err != nil {
|
||||
if errors.Is(err, errForwardNotFound) {
|
||||
response.WriteJSON(w, response.ErrDefault("转发不存在"))
|
||||
return
|
||||
}
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
if err := h.repo.ResetForwardFlow(id, time.Now().UnixMilli()); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
|
||||
func (h *Handler) forwardResume(w http.ResponseWriter, r *http.Request) {
|
||||
id := idFromBody(r, w)
|
||||
if id <= 0 {
|
||||
|
||||
@@ -152,7 +152,11 @@ func evaluateBestExitOwner(owner chainNodeRecord, exits []chainNodeRecord, nodes
|
||||
ownerNode := nodes[owner.NodeID]
|
||||
for _, exit := range exits {
|
||||
exitNode := nodes[exit.NodeID]
|
||||
if exitNode == nil {
|
||||
if !isTunnelProbeNodeOnline(ownerNode) {
|
||||
scores = append(scores, failedBestExitCandidate(owner.NodeID, exit, "owner node offline"))
|
||||
continue
|
||||
}
|
||||
if !isTunnelProbeNodeOnline(exitNode) {
|
||||
scores = append(scores, failedBestExitCandidate(owner.NodeID, exit, "exit node unavailable"))
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -358,9 +358,9 @@ func TestEvaluateBestExitOwnerScoresAllCandidates(t *testing.T) {
|
||||
{NodeID: 31, NodeName: "exit-b", Port: 30031},
|
||||
}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
31: {ID: 31, ServerIP: "10.0.0.31", ServerIPv4: "10.0.0.31", TCPListenAddr: "[::]"},
|
||||
10: {ID: 10, Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Status: 1, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
31: {ID: 31, Status: 1, ServerIP: "10.0.0.31", ServerIPv4: "10.0.0.31", TCPListenAddr: "[::]"},
|
||||
}
|
||||
pinger := func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
switch {
|
||||
@@ -387,12 +387,30 @@ func TestEvaluateBestExitOwnerScoresAllCandidates(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluateBestExitOwnerSkipsOfflineCandidate(t *testing.T) {
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30030}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10"},
|
||||
30: {ID: 30, Status: 0, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30"},
|
||||
}
|
||||
ping := func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
t.Fatalf("offline best-exit candidate should not be probed: node=%d target=%s:%d", nodeID, ip, port)
|
||||
return 0, 100, nil
|
||||
}
|
||||
|
||||
scores := evaluateBestExitOwner(owner, exits, nodes, "", diagnosisExecOptions{}, defaultTunnelProbeTarget(), ping)
|
||||
if len(scores) != 1 || scores[0].Success {
|
||||
t.Fatalf("expected one failed offline candidate, got %+v", scores)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEvaluateBestExitOwnerUsesConfiguredPublicProbeTarget(t *testing.T) {
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry-a"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30001}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, Name: "entry-a", ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10"},
|
||||
30: {ID: 30, Name: "exit-a", ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30"},
|
||||
10: {ID: 10, Name: "entry-a", Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10"},
|
||||
30: {ID: 30, Name: "exit-a", Status: 1, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30"},
|
||||
}
|
||||
target := tunnelProbeTarget{Host: "speed.example.com", Port: 8443}
|
||||
var calls []string
|
||||
@@ -419,8 +437,8 @@ func TestEvaluateBestExitOwnerMarksCandidateFailedWhenOwnerToExitFails(t *testin
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-a", Port: 30030}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
10: {ID: 10, Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Status: 1, ServerIP: "10.0.0.30", ServerIPv4: "10.0.0.30", TCPListenAddr: "[::]"},
|
||||
}
|
||||
pinger := func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
return 0, 100, errBestExitProbeForTest
|
||||
@@ -436,8 +454,8 @@ func TestEvaluateBestExitOwnerMarksCandidateFailedWhenTargetResolutionFails(t *t
|
||||
owner := chainNodeRecord{NodeID: 10, NodeName: "entry"}
|
||||
exits := []chainNodeRecord{{NodeID: 30, NodeName: "exit-v6", Port: 30030}}
|
||||
nodes := map[int64]*nodeRecord{
|
||||
10: {ID: 10, Name: "entry", ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Name: "exit-v6", ServerIP: "2001:db8::30", ServerIPv6: "2001:db8::30", TCPListenAddr: "[::]"},
|
||||
10: {ID: 10, Name: "entry", Status: 1, ServerIP: "10.0.0.10", ServerIPv4: "10.0.0.10", TCPListenAddr: "[::]"},
|
||||
30: {ID: 30, Name: "exit-v6", Status: 1, ServerIP: "2001:db8::30", ServerIPv6: "2001:db8::30", TCPListenAddr: "[::]"},
|
||||
}
|
||||
pinger := func(nodeID int64, ip string, port int, _ diagnosisExecOptions) (float64, float64, error) {
|
||||
t.Fatalf("ping should not be called when target resolution fails: node=%d ip=%s port=%d", nodeID, ip, port)
|
||||
|
||||
@@ -3,6 +3,7 @@ package handler
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"log"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
@@ -13,7 +14,6 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
tunnelQualityProbeInterval = 1 * time.Second
|
||||
tunnelQualityProbeTimeout = 8 * time.Second
|
||||
tunnelQualityPingTimeoutMs = 5000
|
||||
tunnelQualityPruneInterval = 10 * time.Minute
|
||||
@@ -56,7 +56,7 @@ type tunnelQualityProber struct {
|
||||
cache sync.Map // tunnelID (int64) → *tunnelQualitySnapshot
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
interval time.Duration
|
||||
wake chan struct{}
|
||||
lastPrune int64
|
||||
probing int32 // atomic flag: 1 = probeAll running, 0 = idle
|
||||
probeNode bestExitProbeFunc
|
||||
@@ -65,8 +65,8 @@ type tunnelQualityProber struct {
|
||||
// newTunnelQualityProber creates a new prober (not yet running).
|
||||
func newTunnelQualityProber(h *Handler) *tunnelQualityProber {
|
||||
return &tunnelQualityProber{
|
||||
handler: h,
|
||||
interval: tunnelQualityProbeInterval,
|
||||
handler: h,
|
||||
wake: make(chan struct{}, 1),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +86,16 @@ func (p *tunnelQualityProber) Stop() {
|
||||
p.cancel()
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) NotifyConfigChanged() {
|
||||
if p == nil || p.wake == nil {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case p.wake <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
// GetAll returns all cached quality snapshots (latest per tunnel).
|
||||
func (p *tunnelQualityProber) GetAll() []tunnelQualitySnapshot {
|
||||
var items []tunnelQualitySnapshot
|
||||
@@ -109,20 +119,44 @@ func (p *tunnelQualityProber) loop() {
|
||||
// Run once immediately
|
||||
p.probeAll()
|
||||
|
||||
ticker := time.NewTicker(p.interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
timer := time.NewTimer(p.probeInterval())
|
||||
select {
|
||||
case <-p.ctx.Done():
|
||||
stopAndDrainTunnelQualityTimer(timer)
|
||||
return
|
||||
case <-ticker.C:
|
||||
case <-p.wake:
|
||||
stopAndDrainTunnelQualityTimer(timer)
|
||||
continue
|
||||
case <-timer.C:
|
||||
p.probeAll()
|
||||
p.maybePrune()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func stopAndDrainTunnelQualityTimer(timer *time.Timer) {
|
||||
if timer == nil || timer.Stop() {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-timer.C:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probeInterval() time.Duration {
|
||||
if p == nil || p.handler == nil || p.handler.repo == nil {
|
||||
return time.Duration(monitoring.DefaultTunnelQualityProbeIntervalSec) * time.Second
|
||||
}
|
||||
cfg, err := p.handler.repo.GetConfigsByNames([]string{monitoring.ConfigTunnelQualityProbeIntervalSec})
|
||||
if err != nil {
|
||||
return time.Duration(monitoring.DefaultTunnelQualityProbeIntervalSec) * time.Second
|
||||
}
|
||||
seconds := monitoring.TunnelQualityProbeIntervalSecondsFromConfigMap(cfg)
|
||||
return time.Duration(seconds) * time.Second
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) isEnabled() bool {
|
||||
if p == nil || p.handler == nil {
|
||||
return true
|
||||
@@ -246,15 +280,19 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
options := diagnosisExecOptions{
|
||||
commandTimeout: tunnelQualityProbeTimeout,
|
||||
pingTimeoutMS: tunnelQualityPingTimeoutMs,
|
||||
pingCount: 1,
|
||||
timeoutMessage: "探测超时",
|
||||
}
|
||||
p.probeBestExitOwners(tunnelID, inNodes, midNodesGrouped, outNodes, ipPreference, options, probeTarget)
|
||||
|
||||
entry, _, entryOnline := p.firstOnlineChainNode(inNodes)
|
||||
exit, _, exitOnline := p.firstOnlineChainNode(outNodes)
|
||||
|
||||
switch tunnel.Type {
|
||||
case 1:
|
||||
// Port forwarding: entry → public probe target only.
|
||||
if len(inNodes) > 0 {
|
||||
lat, loss, err := p.pingNode(inNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if entryOnline {
|
||||
lat, loss, err := p.pingNode(entry.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -262,24 +300,42 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
} else {
|
||||
snap.ErrorMessage = err.Error()
|
||||
}
|
||||
} else {
|
||||
snap.ErrorMessage = "入口节点均不在线"
|
||||
}
|
||||
case 2:
|
||||
// Tunnel forwarding: entry → exit + exit → Bing
|
||||
probeOK := true
|
||||
|
||||
if len(inNodes) > 0 && len(outNodes) > 0 {
|
||||
if !entryOnline {
|
||||
probeOK = false
|
||||
snap.ErrorMessage = "入口节点均不在线"
|
||||
snap.EntryToExitLatency = -1
|
||||
snap.EntryToExitLoss = 100
|
||||
} else if !exitOnline {
|
||||
probeOK = false
|
||||
snap.ErrorMessage = "出口节点均不在线"
|
||||
snap.EntryToExitLatency = -1
|
||||
snap.EntryToExitLoss = 100
|
||||
} else {
|
||||
var hops []TunnelQualityHop
|
||||
var totalLat float64
|
||||
remainingSuccessProb := 1.0
|
||||
|
||||
nodesInPath := make([]chainNodeRecord, 0, 2+len(midNodesGrouped))
|
||||
nodesInPath = append(nodesInPath, inNodes[0])
|
||||
nodesInPath = append(nodesInPath, entry)
|
||||
for _, midGroup := range midNodesGrouped {
|
||||
if len(midGroup) > 0 {
|
||||
nodesInPath = append(nodesInPath, midGroup[0])
|
||||
mid, _, online := p.firstOnlineChainNode(midGroup)
|
||||
if !online {
|
||||
probeOK = false
|
||||
snap.ErrorMessage = "中间节点组均不在线"
|
||||
break
|
||||
}
|
||||
nodesInPath = append(nodesInPath, mid)
|
||||
}
|
||||
if probeOK {
|
||||
nodesInPath = append(nodesInPath, exit)
|
||||
}
|
||||
nodesInPath = append(nodesInPath, outNodes[0])
|
||||
|
||||
for i := 0; i < len(nodesInPath)-1; i++ {
|
||||
source := nodesInPath[i]
|
||||
@@ -293,7 +349,7 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
}
|
||||
|
||||
targetNode, nodeErr := h.getNodeRecord(target.NodeID)
|
||||
if nodeErr != nil || targetNode == nil {
|
||||
if nodeErr != nil || !isTunnelProbeNodeOnline(targetNode) {
|
||||
snap.ErrorMessage = "节点 " + target.NodeName + " 不可用"
|
||||
probeOK = false
|
||||
hop.Latency = -1
|
||||
@@ -351,8 +407,8 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
}
|
||||
|
||||
// Exit → Bing
|
||||
if len(outNodes) > 0 {
|
||||
lat, loss, err := p.pingNode(outNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if exitOnline {
|
||||
lat, loss, err := p.pingNode(exit.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -367,8 +423,8 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
snap.Success = probeOK
|
||||
default:
|
||||
// Unknown type: entry → public probe target.
|
||||
if len(inNodes) > 0 {
|
||||
lat, loss, err := p.pingNode(inNodes[0].NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if entryOnline {
|
||||
lat, loss, err := p.pingNode(entry.NodeID, probeTarget.Host, probeTarget.Port, options)
|
||||
if err == nil {
|
||||
snap.ExitToBingLatency = lat
|
||||
snap.ExitToBingLoss = loss
|
||||
@@ -376,12 +432,31 @@ func (p *tunnelQualityProber) probeTunnel(tunnelID int64) {
|
||||
} else {
|
||||
snap.ErrorMessage = err.Error()
|
||||
}
|
||||
} else {
|
||||
snap.ErrorMessage = "入口节点均不在线"
|
||||
}
|
||||
}
|
||||
|
||||
p.storeResult(snap)
|
||||
}
|
||||
|
||||
func isTunnelProbeNodeOnline(node *nodeRecord) bool {
|
||||
return node != nil && (node.IsRemote == 1 || node.Status == 1)
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) firstOnlineChainNode(nodes []chainNodeRecord) (chainNodeRecord, *nodeRecord, bool) {
|
||||
if p == nil || p.handler == nil {
|
||||
return chainNodeRecord{}, nil, false
|
||||
}
|
||||
for _, candidate := range nodes {
|
||||
node, err := p.handler.getNodeRecord(candidate.NodeID)
|
||||
if err == nil && isTunnelProbeNodeOnline(node) {
|
||||
return candidate, node, true
|
||||
}
|
||||
}
|
||||
return chainNodeRecord{}, nil, false
|
||||
}
|
||||
|
||||
func (p *tunnelQualityProber) probeBestExitOwners(tunnelID int64, inNodes []chainNodeRecord, chainHops [][]chainNodeRecord, outNodes []chainNodeRecord, ipPreference string, options diagnosisExecOptions, probeTarget tunnelProbeTarget) {
|
||||
if p == nil || p.handler == nil || p.handler.bestExit == nil || len(outNodes) <= 1 {
|
||||
return
|
||||
@@ -444,6 +519,9 @@ func (p *tunnelQualityProber) tcpPingNode(nodeID int64, ip string, port int, opt
|
||||
if nodeErr != nil {
|
||||
return 0, 100, nodeErr
|
||||
}
|
||||
if !isTunnelProbeNodeOnline(node) {
|
||||
return 0, 100, errors.New("节点不在线")
|
||||
}
|
||||
|
||||
var pingData map[string]interface{}
|
||||
var pingErr error
|
||||
|
||||
@@ -26,6 +26,9 @@ func TestTunnelQualityProberUsesConfiguredProbeTarget(t *testing.T) {
|
||||
p := newTunnelQualityProber(h)
|
||||
var calls []string
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
if options.pingCount != 1 {
|
||||
t.Fatalf("expected real-time quality probe count 1, got %d", options.pingCount)
|
||||
}
|
||||
calls = append(calls, fmt.Sprintf("%d|%s|%d", nodeID, ip, port))
|
||||
return 10, 0, nil
|
||||
}
|
||||
@@ -46,6 +49,63 @@ func TestTunnelQualityProberUsesConfiguredProbeTarget(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberSkipsAllOfflineExits(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedQualityForwardTunnel(t, h, 81, []int{0, 0, 0})
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
probeCalls := 0
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
probeCalls++
|
||||
return 0, 100, fmt.Errorf("unexpected probe node=%d target=%s:%d", nodeID, ip, port)
|
||||
}
|
||||
p.probeTunnel(81)
|
||||
|
||||
if probeCalls != 0 {
|
||||
t.Fatalf("expected no TCP probes when all exits are offline, got %d", probeCalls)
|
||||
}
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 {
|
||||
t.Fatalf("expected one quality snapshot, got %+v", snaps)
|
||||
}
|
||||
if snaps[0].Success || snaps[0].ErrorMessage != "出口节点均不在线" {
|
||||
t.Fatalf("expected offline exit snapshot, got %+v", snaps[0])
|
||||
}
|
||||
if snaps[0].EntryToExitLoss != 100 {
|
||||
t.Fatalf("expected 100%% entry-to-exit loss, got %+v", snaps[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberUsesOnlineBackupExit(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedQualityForwardTunnel(t, h, 82, []int{0, 1})
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
var calls []string
|
||||
p.probeNode = func(nodeID int64, ip string, port int, options diagnosisExecOptions) (float64, float64, error) {
|
||||
if options.pingCount != 1 {
|
||||
t.Fatalf("expected real-time quality probe count 1, got %d", options.pingCount)
|
||||
}
|
||||
calls = append(calls, fmt.Sprintf("%d|%s|%d", nodeID, ip, port))
|
||||
return 10, 0, nil
|
||||
}
|
||||
p.probeTunnel(82)
|
||||
|
||||
if slices.Contains(calls, "10|10.0.0.30|30030") {
|
||||
t.Fatalf("did not expect probe to offline primary exit, calls=%+v", calls)
|
||||
}
|
||||
if !slices.Contains(calls, "10|10.0.0.31|30031") {
|
||||
t.Fatalf("expected entry probe to online backup exit, calls=%+v", calls)
|
||||
}
|
||||
if !slices.Contains(calls, "31|www.bing.com|443") {
|
||||
t.Fatalf("expected public probe from online backup exit, calls=%+v", calls)
|
||||
}
|
||||
snaps := p.GetAll()
|
||||
if len(snaps) != 1 || !snaps[0].Success {
|
||||
t.Fatalf("expected successful backup exit snapshot, got %+v", snaps)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberStoresProbeTargetWhenChainIncomplete(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
seedProbeTargetTunnel(t, h, 78, "quality-target-incomplete", "speed.example.com", 8443)
|
||||
@@ -67,3 +127,69 @@ func TestTunnelQualityProberStoresProbeTargetWhenChainIncomplete(t *testing.T) {
|
||||
t.Fatalf("unexpected snapshot target metadata: %+v", snaps[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberUsesConfiguredInterval(t *testing.T) {
|
||||
h := setupProbeTargetTunnelHandler(t)
|
||||
if err := h.repo.UpsertConfig("monitor_tunnel_quality_interval_sec", "15", time.Now().UnixMilli()); err != nil {
|
||||
t.Fatalf("upsert interval config: %v", err)
|
||||
}
|
||||
|
||||
p := newTunnelQualityProber(h)
|
||||
if got := p.probeInterval(); got != 15*time.Second {
|
||||
t.Fatalf("probe interval = %s, want 15s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTunnelQualityProberConfigNotificationIsCoalesced(t *testing.T) {
|
||||
p := newTunnelQualityProber(nil)
|
||||
p.NotifyConfigChanged()
|
||||
p.NotifyConfigChanged()
|
||||
|
||||
if got := len(p.wake); got != 1 {
|
||||
t.Fatalf("wake notifications = %d, want 1", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelQualityProbeIntervalConfigValue(t *testing.T) {
|
||||
got, err := normalizeAndValidateConfigValue("monitor_tunnel_quality_interval_sec", " 15 ")
|
||||
if err != nil || got != "15" {
|
||||
t.Fatalf("normalize interval = %q, %v", got, err)
|
||||
}
|
||||
if _, err := normalizeAndValidateConfigValue("monitor_tunnel_quality_interval_sec", "0"); err == nil {
|
||||
t.Fatalf("expected invalid interval to be rejected")
|
||||
}
|
||||
}
|
||||
|
||||
func seedQualityForwardTunnel(t *testing.T, h *Handler, tunnelID int64, exitStatuses []int) {
|
||||
t.Helper()
|
||||
now := time.Now().UnixMilli()
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, inx, ip_preference, probe_target_host, probe_target_port)
|
||||
VALUES(?, ?, 1, 2, 'tls', 1, ?, ?, 1, ?, '', '', 0)
|
||||
`, tunnelID, fmt.Sprintf("quality-forward-%d", tunnelID), now, now, tunnelID).Error; err != nil {
|
||||
t.Fatalf("insert forwarding tunnel: %v", err)
|
||||
}
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol)
|
||||
VALUES(?, '1', 10, 30001, 'fifo', 1, 'tls')
|
||||
`, tunnelID).Error; err != nil {
|
||||
t.Fatalf("insert entry chain: %v", err)
|
||||
}
|
||||
for i, status := range exitStatuses {
|
||||
nodeID := int64(30 + i)
|
||||
port := 30030 + i
|
||||
ip := fmt.Sprintf("10.0.0.%d", nodeID)
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO node(id, name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx)
|
||||
VALUES(?, ?, ?, ?, ?, '', '30000-30100', '', 'v1', 1, 1, 1, ?, ?, ?, '[::]', '[::]', 0)
|
||||
`, nodeID, fmt.Sprintf("exit-%d", i+1), fmt.Sprintf("exit-secret-%d", i+1), ip, ip, now, now, status).Error; err != nil {
|
||||
t.Fatalf("insert exit node %d: %v", nodeID, err)
|
||||
}
|
||||
if err := h.repo.DB().Exec(`
|
||||
INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol)
|
||||
VALUES(?, '3', ?, ?, 'fifo', ?, 'tls')
|
||||
`, tunnelID, nodeID, port, i+1).Error; err != nil {
|
||||
t.Fatalf("insert exit chain %d: %v", nodeID, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,52 @@
|
||||
package monitoring
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
const (
|
||||
ConfigTunnelQualityProbeIntervalSec = "monitor_tunnel_quality_interval_sec"
|
||||
DefaultTunnelQualityProbeIntervalSec = 1
|
||||
MinTunnelQualityProbeIntervalSec = 1
|
||||
MaxTunnelQualityProbeIntervalSec = 3600
|
||||
)
|
||||
|
||||
func TunnelQualityProbeIntervalSecondsFromConfigMap(cfg map[string]string) int {
|
||||
if cfg == nil {
|
||||
return DefaultTunnelQualityProbeIntervalSec
|
||||
}
|
||||
seconds, err := parseTunnelQualityProbeIntervalSeconds(cfg[ConfigTunnelQualityProbeIntervalSec])
|
||||
if err != nil {
|
||||
return DefaultTunnelQualityProbeIntervalSec
|
||||
}
|
||||
return seconds
|
||||
}
|
||||
|
||||
func NormalizeTunnelQualityProbeIntervalSeconds(value string) (string, error) {
|
||||
seconds, err := parseTunnelQualityProbeIntervalSeconds(value)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return strconv.Itoa(seconds), nil
|
||||
}
|
||||
|
||||
func parseTunnelQualityProbeIntervalSeconds(value string) (int, error) {
|
||||
trimmed := strings.TrimSpace(value)
|
||||
if trimmed == "" {
|
||||
return 0, fmt.Errorf("隧道质量探测间隔不能为空")
|
||||
}
|
||||
seconds, err := strconv.Atoi(trimmed)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("隧道质量探测间隔必须是整数")
|
||||
}
|
||||
if seconds < MinTunnelQualityProbeIntervalSec || seconds > MaxTunnelQualityProbeIntervalSec {
|
||||
return 0, fmt.Errorf(
|
||||
"隧道质量探测间隔必须在 %d 到 %d 秒之间",
|
||||
MinTunnelQualityProbeIntervalSec,
|
||||
MaxTunnelQualityProbeIntervalSec,
|
||||
)
|
||||
}
|
||||
return seconds, nil
|
||||
}
|
||||
@@ -0,0 +1,37 @@
|
||||
package monitoring
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestTunnelQualityProbeIntervalSecondsFromConfigMap(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
cfg map[string]string
|
||||
want int
|
||||
}{
|
||||
{name: "missing config", cfg: nil, want: DefaultTunnelQualityProbeIntervalSec},
|
||||
{name: "configured", cfg: map[string]string{ConfigTunnelQualityProbeIntervalSec: "15"}, want: 15},
|
||||
{name: "invalid", cfg: map[string]string{ConfigTunnelQualityProbeIntervalSec: "0"}, want: DefaultTunnelQualityProbeIntervalSec},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if got := TunnelQualityProbeIntervalSecondsFromConfigMap(tt.cfg); got != tt.want {
|
||||
t.Fatalf("interval = %d, want %d", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNormalizeTunnelQualityProbeIntervalSeconds(t *testing.T) {
|
||||
for _, value := range []string{"1", "15", "3600"} {
|
||||
if got, err := NormalizeTunnelQualityProbeIntervalSeconds(value); err != nil || got != value {
|
||||
t.Fatalf("normalize %q = %q, %v", value, got, err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, value := range []string{"", "0", "3601", "1.5", "abc"} {
|
||||
if got, err := NormalizeTunnelQualityProbeIntervalSeconds(value); err == nil {
|
||||
t.Fatalf("normalize %q unexpectedly succeeded with %q", value, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
package repo
|
||||
|
||||
import (
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestResetForwardFlowOnlyUpdatesSelectedForward(t *testing.T) {
|
||||
r, err := Open(filepath.Join(t.TempDir(), "forward-flow-reset.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("open repo: %v", err)
|
||||
}
|
||||
defer r.Close()
|
||||
|
||||
const originalUpdated int64 = 1000
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status)
|
||||
VALUES(2, 'owner', 'pwd', 1, 0, 100, 700, 900, 0, 10, 1000, 1000, 1)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert user: %v", err)
|
||||
}
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx)
|
||||
VALUES(1, 'tunnel', 1, 1, 'tls', 1, 1000, 1000, 1, NULL, 0)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert tunnel: %v", err)
|
||||
}
|
||||
if err := r.DB().Exec(`
|
||||
INSERT INTO user_tunnel(id, user_id, tunnel_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status)
|
||||
VALUES(10, 2, 1, 10, 100, 500, 600, 0, 0, 1)
|
||||
`).Error; err != nil {
|
||||
t.Fatalf("insert user tunnel: %v", err)
|
||||
}
|
||||
if err := r.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)
|
||||
VALUES
|
||||
(20, 2, 'owner', 'target', 1, '127.0.0.1:80', 'fifo', 111, 222, 1000, ?, 1, 0),
|
||||
(21, 2, 'owner', 'other', 1, '127.0.0.1:81', 'fifo', 333, 444, 1000, ?, 1, 1)
|
||||
`, originalUpdated, originalUpdated).Error; err != nil {
|
||||
t.Fatalf("insert forwards: %v", err)
|
||||
}
|
||||
|
||||
const resetAt int64 = 2000
|
||||
if err := r.ResetForwardFlow(20, resetAt); err != nil {
|
||||
t.Fatalf("ResetForwardFlow: %v", err)
|
||||
}
|
||||
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM forward WHERE id = 20", 0)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM forward WHERE id = 20", 0)
|
||||
assertForwardFlowResetValue(t, r, "SELECT updated_time FROM forward WHERE id = 20", resetAt)
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM forward WHERE id = 21", 333)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM forward WHERE id = 21", 444)
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM user WHERE id = 2", 700)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM user WHERE id = 2", 900)
|
||||
assertForwardFlowResetValue(t, r, "SELECT in_flow FROM user_tunnel WHERE id = 10", 500)
|
||||
assertForwardFlowResetValue(t, r, "SELECT out_flow FROM user_tunnel WHERE id = 10", 600)
|
||||
}
|
||||
|
||||
func TestResetForwardFlowRejectsUninitializedRepository(t *testing.T) {
|
||||
var r *Repository
|
||||
if err := r.ResetForwardFlow(20, 2000); err == nil {
|
||||
t.Fatal("expected uninitialized repository error")
|
||||
}
|
||||
}
|
||||
|
||||
func assertForwardFlowResetValue(t *testing.T, r *Repository, query string, want int64) {
|
||||
t.Helper()
|
||||
var got int64
|
||||
if err := r.DB().Raw(query).Scan(&got).Error; err != nil {
|
||||
t.Fatalf("query %q: %v", query, err)
|
||||
}
|
||||
if got != want {
|
||||
t.Fatalf("query %q returned %d, want %d", query, got, want)
|
||||
}
|
||||
}
|
||||
@@ -197,6 +197,19 @@ func (r *Repository) ResetUserFlowByUserTunnel(userTunnelID int64) {
|
||||
Updates(map[string]interface{}{"in_flow": 0, "out_flow": 0}).Error
|
||||
}
|
||||
|
||||
func (r *Repository) ResetForwardFlow(forwardID int64, now int64) error {
|
||||
if r == nil || r.db == nil {
|
||||
return errors.New("repository not initialized")
|
||||
}
|
||||
return r.db.Model(&model.Forward{}).
|
||||
Where("id = ?", forwardID).
|
||||
Updates(map[string]interface{}{
|
||||
"in_flow": 0,
|
||||
"out_flow": 0,
|
||||
"updated_time": now,
|
||||
}).Error
|
||||
}
|
||||
|
||||
func (r *Repository) GetUsernameByID(userID int64) string {
|
||||
if r == nil || r.db == nil {
|
||||
return ""
|
||||
|
||||
@@ -151,6 +151,7 @@ const (
|
||||
initialBackoff = 2 * time.Second // 重连初始退避
|
||||
maxBackoff = 2 * time.Minute // 重连最大退避
|
||||
defaultMetricReportInterval = 5 * time.Second
|
||||
maxConcurrentTCPPings = 8
|
||||
)
|
||||
|
||||
type WebSocketReporter struct {
|
||||
@@ -172,6 +173,7 @@ type WebSocketReporter struct {
|
||||
connecting bool // 正在连接状态
|
||||
connMutex sync.Mutex // 连接状态锁
|
||||
aesCrypto *crypto.AESCrypto // AES加密器
|
||||
tcpPingSem chan struct{} // 限制诊断探测并发,避免离线目标耗尽连接
|
||||
}
|
||||
|
||||
var wsDial = func(dialer *websocket.Dialer, rawURL string) (*websocket.Conn, *http.Response, error) {
|
||||
@@ -201,6 +203,29 @@ func NewWebSocketReporter(serverURL string, secret string) *WebSocketReporter {
|
||||
connected: false,
|
||||
connecting: false,
|
||||
aesCrypto: aesCrypto,
|
||||
tcpPingSem: make(chan struct{}, maxConcurrentTCPPings),
|
||||
}
|
||||
}
|
||||
|
||||
func (w *WebSocketReporter) tryAcquireTCPPingSlot() bool {
|
||||
if w == nil || w.tcpPingSem == nil {
|
||||
return false
|
||||
}
|
||||
select {
|
||||
case w.tcpPingSem <- struct{}{}:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
func (w *WebSocketReporter) releaseTCPPingSlot() {
|
||||
if w == nil || w.tcpPingSem == nil {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-w.tcpPingSem:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
@@ -840,9 +865,14 @@ func (w *WebSocketReporter) routeCommand(cmd CommandMessage) {
|
||||
|
||||
// TCP Ping 诊断命令(只读,不需要保存配置)
|
||||
case "TcpPing":
|
||||
response.Type = "TcpPingResponse"
|
||||
if !w.tryAcquireTCPPingSlot() {
|
||||
err = fmt.Errorf("TCP探测任务过多,请稍后重试")
|
||||
break
|
||||
}
|
||||
defer w.releaseTCPPingSlot()
|
||||
var tcpPingResult TcpPingResponse
|
||||
tcpPingResult, err = w.handleTcpPing(cmd.Data)
|
||||
response.Type = "TcpPingResponse"
|
||||
response.Data = tcpPingResult
|
||||
// needSaveConfig = false (默认值)
|
||||
|
||||
|
||||
@@ -148,6 +148,25 @@ func TestNewWebSocketReporterUsesReducedMetricInterval(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestWebSocketReporterLimitsConcurrentTCPPings(t *testing.T) {
|
||||
reporter := &WebSocketReporter{tcpPingSem: make(chan struct{}, maxConcurrentTCPPings)}
|
||||
for i := 0; i < maxConcurrentTCPPings; i++ {
|
||||
if !reporter.tryAcquireTCPPingSlot() {
|
||||
t.Fatalf("expected TCP ping slot %d to be available", i)
|
||||
}
|
||||
}
|
||||
if reporter.tryAcquireTCPPingSlot() {
|
||||
t.Fatalf("expected TCP ping concurrency limit at %d", maxConcurrentTCPPings)
|
||||
}
|
||||
for i := 0; i < maxConcurrentTCPPings; i++ {
|
||||
reporter.releaseTCPPingSlot()
|
||||
}
|
||||
if !reporter.tryAcquireTCPPingSlot() {
|
||||
t.Fatalf("expected released TCP ping slot to be reusable")
|
||||
}
|
||||
reporter.releaseTCPPingSlot()
|
||||
}
|
||||
|
||||
func TestFormatWebSocketDialErrorIncludesHTTPStatus(t *testing.T) {
|
||||
err := errors.New("websocket: bad handshake")
|
||||
resp := &http.Response{
|
||||
|
||||
+247
-40
@@ -1,4 +1,31 @@
|
||||
#!/bin/bash
|
||||
#!/bin/sh
|
||||
# shellcheck shell=bash
|
||||
|
||||
# Alpine 默认不带 Bash。先用系统自带的 /bin/sh 安装/切换到 Bash,
|
||||
# 后续主体继续使用 Bash 语法,避免要求用户手动准备运行环境。
|
||||
if [ -z "${BASH_VERSION:-}" ]; then
|
||||
if command -v bash >/dev/null 2>&1; then
|
||||
exec bash "$0" "$@"
|
||||
fi
|
||||
|
||||
if [ -f /etc/alpine-release ] && command -v apk >/dev/null 2>&1; then
|
||||
if [ "$(id -u)" -eq 0 ]; then
|
||||
apk add --no-cache bash
|
||||
elif command -v sudo >/dev/null 2>&1; then
|
||||
sudo apk add --no-cache bash
|
||||
elif command -v doas >/dev/null 2>&1; then
|
||||
doas apk add --no-cache bash
|
||||
else
|
||||
echo "❌ Alpine 安装需要 root 权限,或已配置 sudo/doas。" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
exec bash "$0" "$@"
|
||||
fi
|
||||
|
||||
echo "❌ 此安装脚本需要 Bash。" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# GitHub repo used for release downloads
|
||||
REPO="Sagit-chu/flux-panel"
|
||||
@@ -24,11 +51,14 @@ get_architecture() {
|
||||
|
||||
# 安装目录
|
||||
INSTALL_DIR="/etc/flux_agent"
|
||||
FLUX_AGENT_SYSTEMD_SERVICE_FILE="/etc/systemd/system/flux_agent.service"
|
||||
FLUX_AGENT_OPENRC_SERVICE_FILE="/etc/init.d/flux_agent"
|
||||
LEGACY_GOST_BINARY="/usr/local/bin/gost"
|
||||
LEGACY_GOST_CONFIG_DIR="/etc/gost"
|
||||
LEGACY_GOST_SERVICE_FILE_ETC="/etc/systemd/system/gost.service"
|
||||
LEGACY_GOST_SERVICE_FILE_LIB="/lib/systemd/system/gost.service"
|
||||
LEGACY_GOST_SERVICE_FILE_USR_LIB="/usr/lib/systemd/system/gost.service"
|
||||
SERVICE_MANAGER="${SERVICE_MANAGER:-}"
|
||||
|
||||
# 镜像加速配置(可由面板传入或交互式询问)
|
||||
PROXY_ENABLED="${PROXY_ENABLED:-}"
|
||||
@@ -256,6 +286,199 @@ write_flux_agent_config() {
|
||||
"$(json_escape "$SECRET")" > "$path"
|
||||
}
|
||||
|
||||
ensure_service_manager() {
|
||||
if [[ -n "$SERVICE_MANAGER" ]]; then
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd|openrc)
|
||||
return 0
|
||||
;;
|
||||
*)
|
||||
echo "❌ 不支持的服务管理器: $SERVICE_MANAGER" >&2
|
||||
return 1
|
||||
;;
|
||||
esac
|
||||
fi
|
||||
|
||||
if command -v systemctl >/dev/null 2>&1 && [[ -d /run/systemd/system ]]; then
|
||||
SERVICE_MANAGER="systemd"
|
||||
return 0
|
||||
fi
|
||||
if command -v rc-service >/dev/null 2>&1 && command -v rc-update >/dev/null 2>&1; then
|
||||
SERVICE_MANAGER="openrc"
|
||||
return 0
|
||||
fi
|
||||
|
||||
echo "❌ 未检测到受支持的服务管理器(systemd 或 OpenRC)。" >&2
|
||||
return 1
|
||||
}
|
||||
|
||||
flux_agent_service_exists() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
[[ -f "$FLUX_AGENT_SYSTEMD_SERVICE_FILE" ]] || \
|
||||
systemctl list-units --full -all 2>/dev/null | grep -Fq "flux_agent.service"
|
||||
;;
|
||||
openrc)
|
||||
[[ -f "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]]
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
stop_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl stop flux_agent 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent stop 2>/dev/null || true
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
disable_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl disable flux_agent 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rc-update del flux_agent default 2>/dev/null || true
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
write_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
mkdir -p "$(dirname "$FLUX_AGENT_SYSTEMD_SERVICE_FILE")"
|
||||
cat > "$FLUX_AGENT_SYSTEMD_SERVICE_FILE" <<EOF
|
||||
[Unit]
|
||||
Description=Flux_agent Proxy Service
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
WorkingDirectory=$INSTALL_DIR
|
||||
ExecStart=$INSTALL_DIR/flux_agent
|
||||
Restart=on-failure
|
||||
StandardOutput=null
|
||||
StandardError=null
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
EOF
|
||||
;;
|
||||
openrc)
|
||||
mkdir -p "$(dirname "$FLUX_AGENT_OPENRC_SERVICE_FILE")"
|
||||
cat > "$FLUX_AGENT_OPENRC_SERVICE_FILE" <<EOF
|
||||
#!/sbin/openrc-run
|
||||
|
||||
name="flux_agent"
|
||||
description="Flux_agent Proxy Service"
|
||||
command="$INSTALL_DIR/flux_agent"
|
||||
directory="$INSTALL_DIR"
|
||||
command_background="yes"
|
||||
pidfile="/run/\${RC_SVCNAME}.pid"
|
||||
output_log="/dev/null"
|
||||
error_log="/dev/null"
|
||||
|
||||
depend() {
|
||||
need net
|
||||
}
|
||||
EOF
|
||||
chmod +x "$FLUX_AGENT_OPENRC_SERVICE_FILE"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
enable_and_start_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl daemon-reload
|
||||
systemctl enable flux_agent
|
||||
systemctl start flux_agent
|
||||
;;
|
||||
openrc)
|
||||
rc-update add flux_agent default
|
||||
rc-service flux_agent start
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
start_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl start flux_agent
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent start
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
flux_agent_service_is_active() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl is-active --quiet flux_agent
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent status >/dev/null 2>&1
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
flux_agent_service_status() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
systemctl is-active flux_agent 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rc-service flux_agent status 2>/dev/null || true
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
remove_flux_agent_service() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
rm -f "$FLUX_AGENT_SYSTEMD_SERVICE_FILE"
|
||||
systemctl daemon-reload 2>/dev/null || true
|
||||
;;
|
||||
openrc)
|
||||
rm -f "$FLUX_AGENT_OPENRC_SERVICE_FILE"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
flux_agent_service_status_hint() {
|
||||
ensure_service_manager || return 1
|
||||
|
||||
case "$SERVICE_MANAGER" in
|
||||
systemd)
|
||||
echo "systemctl status flux_agent --no-pager"
|
||||
;;
|
||||
openrc)
|
||||
echo "rc-service flux_agent status"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
cleanup_legacy_gost_installation() {
|
||||
local matched_service_files=()
|
||||
local service_file=""
|
||||
@@ -341,19 +564,21 @@ install_flux_agent() {
|
||||
|
||||
get_config_params
|
||||
|
||||
# 检查并安装 tcpkill
|
||||
# 检查并安装 tcpkill
|
||||
check_and_install_tcpkill
|
||||
|
||||
ensure_service_manager || exit 1
|
||||
|
||||
mkdir -p "$INSTALL_DIR"
|
||||
|
||||
local tmp_binary="$INSTALL_DIR/flux_agent.new"
|
||||
|
||||
# 停止并禁用已有服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
if flux_agent_service_exists; then
|
||||
echo "🔍 检测到已存在的flux_agent服务"
|
||||
systemctl stop flux_agent 2>/dev/null && echo "🛑 停止服务"
|
||||
systemctl disable flux_agent 2>/dev/null && echo "🚫 禁用自启"
|
||||
stop_flux_agent_service
|
||||
echo "🛑 停止服务"
|
||||
disable_flux_agent_service
|
||||
echo "🚫 禁用自启"
|
||||
fi
|
||||
|
||||
# 下载 flux_agent
|
||||
@@ -392,38 +617,21 @@ EOF
|
||||
# 加强权限
|
||||
chmod 600 "$INSTALL_DIR"/*.json
|
||||
|
||||
# 创建 systemd 服务
|
||||
SERVICE_FILE="/etc/systemd/system/flux_agent.service"
|
||||
cat > "$SERVICE_FILE" <<EOF
|
||||
[Unit]
|
||||
Description=Flux_agent Proxy Service
|
||||
After=network.target
|
||||
|
||||
[Service]
|
||||
WorkingDirectory=$INSTALL_DIR
|
||||
ExecStart=$INSTALL_DIR/flux_agent
|
||||
Restart=on-failure
|
||||
StandardOutput=null
|
||||
StandardError=null
|
||||
|
||||
[Install]
|
||||
WantedBy=multi-user.target
|
||||
EOF
|
||||
# 创建 systemd 或 OpenRC 服务
|
||||
write_flux_agent_service
|
||||
|
||||
# 启动服务
|
||||
systemctl daemon-reload
|
||||
systemctl enable flux_agent
|
||||
systemctl start flux_agent
|
||||
enable_and_start_flux_agent_service
|
||||
|
||||
# 检查状态
|
||||
echo "🔄 检查服务状态..."
|
||||
if systemctl is-active --quiet flux_agent; then
|
||||
if flux_agent_service_is_active; then
|
||||
echo "✅ 安装完成,flux_agent服务已启动并设置为开机启动。"
|
||||
echo "📁 配置目录: $INSTALL_DIR"
|
||||
echo "🔧 服务状态: $(systemctl is-active flux_agent)"
|
||||
echo "🔧 服务状态: $(flux_agent_service_status)"
|
||||
else
|
||||
echo "❌ flux_agent服务启动失败,请执行以下命令查看状态:"
|
||||
echo "systemctl status flux_agent --no-pager"
|
||||
flux_agent_service_status_hint
|
||||
fi
|
||||
}
|
||||
|
||||
@@ -443,6 +651,7 @@ update_flux_agent() {
|
||||
|
||||
# 检查并安装 tcpkill
|
||||
check_and_install_tcpkill
|
||||
ensure_service_manager || return 1
|
||||
|
||||
# 先下载新版本
|
||||
echo "⬇️ 下载最新版本..."
|
||||
@@ -455,9 +664,9 @@ update_flux_agent() {
|
||||
cleanup_legacy_gost_installation
|
||||
|
||||
# 停止服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
if flux_agent_service_exists; then
|
||||
echo "🛑 停止 flux_agent 服务..."
|
||||
systemctl stop flux_agent
|
||||
stop_flux_agent_service
|
||||
fi
|
||||
|
||||
# 替换文件
|
||||
@@ -469,7 +678,7 @@ update_flux_agent() {
|
||||
|
||||
# 重启服务
|
||||
echo "🔄 重启服务..."
|
||||
systemctl start flux_agent
|
||||
start_flux_agent_service
|
||||
|
||||
echo "✅ 更新完成,服务已重新启动。"
|
||||
}
|
||||
@@ -477,6 +686,7 @@ update_flux_agent() {
|
||||
# 卸载功能
|
||||
uninstall_flux_agent() {
|
||||
echo "🗑️ 开始卸载 flux_agent..."
|
||||
ensure_service_manager || return 1
|
||||
|
||||
read -p "确认卸载 flux_agent 吗?此操作将删除所有相关文件 (y/N): " confirm
|
||||
if [[ "$confirm" != "y" && "$confirm" != "Y" ]]; then
|
||||
@@ -485,15 +695,15 @@ uninstall_flux_agent() {
|
||||
fi
|
||||
|
||||
# 停止并禁用服务
|
||||
if systemctl list-units --full -all | grep -Fq "flux_agent.service"; then
|
||||
if flux_agent_service_exists; then
|
||||
echo "🛑 停止并禁用服务..."
|
||||
systemctl stop flux_agent 2>/dev/null
|
||||
systemctl disable flux_agent 2>/dev/null
|
||||
stop_flux_agent_service
|
||||
disable_flux_agent_service
|
||||
fi
|
||||
|
||||
# 删除服务文件
|
||||
if [[ -f "/etc/systemd/system/flux_agent.service" ]]; then
|
||||
rm -f "/etc/systemd/system/flux_agent.service"
|
||||
if [[ -f "$FLUX_AGENT_SYSTEMD_SERVICE_FILE" || -f "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]]; then
|
||||
remove_flux_agent_service
|
||||
echo "🧹 删除服务文件"
|
||||
fi
|
||||
|
||||
@@ -503,9 +713,6 @@ uninstall_flux_agent() {
|
||||
echo "🧹 删除安装目录: $INSTALL_DIR"
|
||||
fi
|
||||
|
||||
# 重载 systemd
|
||||
systemctl daemon-reload
|
||||
|
||||
echo "✅ 卸载完成"
|
||||
}
|
||||
|
||||
|
||||
@@ -90,6 +90,7 @@ test_update_flux_agent_asks_for_proxy_config() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
@@ -165,6 +166,7 @@ test_install_flux_agent_preserves_legacy_gost_when_download_fails() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
@@ -199,6 +201,7 @@ test_update_flux_agent_preserves_legacy_gost_when_download_fails() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
cat > "$INSTALL_DIR/flux_agent" <<'EOF'
|
||||
#!/bin/bash
|
||||
@@ -236,7 +239,9 @@ test_install_flux_agent_writes_json_safe_config() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
SERVICE_MANAGER="systemd"
|
||||
INSTALL_DIR=$(mktemp -d)
|
||||
FLUX_AGENT_SYSTEMD_SERVICE_FILE="$INSTALL_DIR/flux_agent.service"
|
||||
SERVER_ADDR='panel"addr'
|
||||
SECRET='sec\ret"1'
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
@@ -274,6 +279,117 @@ EOF
|
||||
assert_equals "$expected" "$actual" "install_flux_agent should JSON-escape config values"
|
||||
)
|
||||
|
||||
test_install_script_bootstraps_bash_for_alpine() (
|
||||
set -euo pipefail
|
||||
|
||||
local shebang
|
||||
shebang=$(head -n 1 "$ROOT_DIR/install.sh")
|
||||
assert_equals "#!/bin/sh" "$shebang" "install.sh should start with Alpine's default shell"
|
||||
|
||||
grep -Fq 'apk add --no-cache bash' "$ROOT_DIR/install.sh" || \
|
||||
fail "install.sh should bootstrap Bash through apk on Alpine"
|
||||
)
|
||||
|
||||
test_install_flux_agent_uses_openrc() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
local temp_root
|
||||
temp_root=$(mktemp -d)
|
||||
INSTALL_DIR="$temp_root/flux_agent"
|
||||
FLUX_AGENT_OPENRC_SERVICE_FILE="$temp_root/init.d/flux_agent"
|
||||
SERVICE_MANAGER="openrc"
|
||||
SERVER_ADDR="panel.example.com:443"
|
||||
SECRET="secret"
|
||||
DOWNLOAD_URL="https://example.com/gost"
|
||||
|
||||
local rc_service_calls=""
|
||||
local rc_update_calls=""
|
||||
|
||||
ask_proxy_config() { :; }
|
||||
ensure_download_url_initialized() { :; }
|
||||
get_config_params() { :; }
|
||||
check_and_install_tcpkill() { :; }
|
||||
cleanup_legacy_gost_installation() { :; }
|
||||
curl() {
|
||||
local output=""
|
||||
while [[ $# -gt 0 ]]; do
|
||||
if [[ "$1" == "-o" ]]; then
|
||||
output="$2"
|
||||
shift 2
|
||||
continue
|
||||
fi
|
||||
shift
|
||||
done
|
||||
|
||||
cat > "$output" <<'EOF'
|
||||
#!/bin/sh
|
||||
echo "new version"
|
||||
EOF
|
||||
chmod +x "$output"
|
||||
}
|
||||
rc-service() {
|
||||
rc_service_calls+=$'\n'"$*"
|
||||
if [[ "$2" == "status" ]]; then
|
||||
echo "status: started"
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
rc-update() {
|
||||
rc_update_calls+=$'\n'"$*"
|
||||
return 0
|
||||
}
|
||||
|
||||
install_flux_agent >/dev/null
|
||||
|
||||
[[ -x "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]] || fail "OpenRC service file should be executable"
|
||||
grep -Fq '#!/sbin/openrc-run' "$FLUX_AGENT_OPENRC_SERVICE_FILE" || \
|
||||
fail "OpenRC service should use openrc-run"
|
||||
grep -Fq "command=\"$INSTALL_DIR/flux_agent\"" "$FLUX_AGENT_OPENRC_SERVICE_FILE" || \
|
||||
fail "OpenRC service should launch the installed flux_agent binary"
|
||||
grep -Fq 'command_background="yes"' "$FLUX_AGENT_OPENRC_SERVICE_FILE" || \
|
||||
fail "OpenRC service should run flux_agent in the background"
|
||||
if command -v openrc-run >/dev/null 2>&1; then
|
||||
"$FLUX_AGENT_OPENRC_SERVICE_FILE" describe >/dev/null 2>&1
|
||||
fi
|
||||
[[ "$rc_update_calls" == *"add flux_agent default"* ]] || \
|
||||
fail "OpenRC install should enable flux_agent in the default runlevel"
|
||||
[[ "$rc_service_calls" == *"start"* ]] || fail "OpenRC install should start flux_agent"
|
||||
[[ "$rc_service_calls" == *"status"* ]] || fail "OpenRC install should verify flux_agent status"
|
||||
)
|
||||
|
||||
test_remove_flux_agent_service_uses_openrc() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
|
||||
local temp_root
|
||||
temp_root=$(mktemp -d)
|
||||
SERVICE_MANAGER="openrc"
|
||||
FLUX_AGENT_OPENRC_SERVICE_FILE="$temp_root/init.d/flux_agent"
|
||||
mkdir -p "$(dirname "$FLUX_AGENT_OPENRC_SERVICE_FILE")"
|
||||
: > "$FLUX_AGENT_OPENRC_SERVICE_FILE"
|
||||
|
||||
local rc_service_calls=""
|
||||
local rc_update_calls=""
|
||||
rc-service() {
|
||||
rc_service_calls+=$'\n'"$*"
|
||||
return 0
|
||||
}
|
||||
rc-update() {
|
||||
rc_update_calls+=$'\n'"$*"
|
||||
return 0
|
||||
}
|
||||
|
||||
stop_flux_agent_service
|
||||
disable_flux_agent_service
|
||||
remove_flux_agent_service
|
||||
|
||||
[[ "$rc_service_calls" == *"stop"* ]] || fail "OpenRC uninstall should stop flux_agent"
|
||||
[[ "$rc_update_calls" == *"del flux_agent default"* ]] || \
|
||||
fail "OpenRC uninstall should remove flux_agent from the default runlevel"
|
||||
[[ ! -e "$FLUX_AGENT_OPENRC_SERVICE_FILE" ]] || fail "OpenRC uninstall should remove its service file"
|
||||
)
|
||||
|
||||
test_cleanup_legacy_gost_installation_removes_service_and_binary() (
|
||||
set -euo pipefail
|
||||
load_script_without_main "$ROOT_DIR/install.sh"
|
||||
@@ -504,6 +620,9 @@ test_update_flux_agent_skips_proxy_prompt_when_not_installed
|
||||
test_install_flux_agent_preserves_legacy_gost_when_download_fails
|
||||
test_update_flux_agent_preserves_legacy_gost_when_download_fails
|
||||
test_install_flux_agent_writes_json_safe_config
|
||||
test_install_script_bootstraps_bash_for_alpine
|
||||
test_install_flux_agent_uses_openrc
|
||||
test_remove_flux_agent_service_uses_openrc
|
||||
test_cleanup_legacy_gost_installation_removes_service_and_binary
|
||||
test_cleanup_legacy_gost_installation_preserves_unrelated_gost
|
||||
test_install_script_accepts_proxy_url_env_without_prompt
|
||||
@@ -514,4 +633,4 @@ test_panel_install_script_uses_default_proxy
|
||||
test_panel_install_script_accepts_proxy_url_env_without_prompt
|
||||
test_panel_install_script_defaults_proxy_on_eof
|
||||
|
||||
echo "install script proxy tests passed"
|
||||
echo "install script tests passed"
|
||||
|
||||
@@ -217,6 +217,8 @@ export const pauseForwardService = (forwardId: number) =>
|
||||
Network.post("/forward/pause", { id: forwardId });
|
||||
export const resumeForwardService = (forwardId: number) =>
|
||||
Network.post("/forward/resume", { id: forwardId });
|
||||
export const resetForwardFlow = (forwardId: number) =>
|
||||
Network.post("/forward/reset-flow", { id: forwardId });
|
||||
|
||||
// 转发诊断操作
|
||||
export const diagnoseForward = (forwardId: number) =>
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
export const TUNNEL_QUALITY_INTERVAL_CONFIG_KEY =
|
||||
"monitor_tunnel_quality_interval_sec";
|
||||
export const DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC = 1;
|
||||
export const MIN_TUNNEL_QUALITY_INTERVAL_SEC = 1;
|
||||
export const MAX_TUNNEL_QUALITY_INTERVAL_SEC = 3600;
|
||||
|
||||
export const parseTunnelQualityIntervalSeconds = (value: unknown): number => {
|
||||
const seconds = Number(value);
|
||||
|
||||
return Number.isInteger(seconds) &&
|
||||
seconds >= MIN_TUNNEL_QUALITY_INTERVAL_SEC &&
|
||||
seconds <= MAX_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
? seconds
|
||||
: DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC;
|
||||
};
|
||||
|
||||
export const validateTunnelQualityInterval = (value: string): string | null => {
|
||||
const normalized = value.trim();
|
||||
|
||||
if (!normalized) {
|
||||
return "请输入探测间隔";
|
||||
}
|
||||
const seconds = Number(normalized);
|
||||
|
||||
if (!Number.isInteger(seconds)) {
|
||||
return "探测间隔必须是整数";
|
||||
}
|
||||
if (
|
||||
seconds < MIN_TUNNEL_QUALITY_INTERVAL_SEC ||
|
||||
seconds > MAX_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
) {
|
||||
return `探测间隔必须在 ${MIN_TUNNEL_QUALITY_INTERVAL_SEC} 到 ${MAX_TUNNEL_QUALITY_INTERVAL_SEC} 秒之间`;
|
||||
}
|
||||
|
||||
return null;
|
||||
};
|
||||
|
||||
export const tunnelQualityIntervalLabel = (seconds: number): string =>
|
||||
seconds === 1 ? "每秒" : `每 ${seconds} 秒`;
|
||||
@@ -43,6 +43,14 @@ import { BackIcon, SettingsIcon } from "@/components/icons";
|
||||
import { ThemeSettings } from "@/components/theme-settings";
|
||||
import { isAdmin } from "@/utils/auth";
|
||||
import { getCachedConfigs, configCache, updateSiteConfig } from "@/config/site";
|
||||
import {
|
||||
DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
MAX_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
MIN_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
parseTunnelQualityIntervalSeconds,
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
validateTunnelQualityInterval,
|
||||
} from "@/config/tunnel-quality";
|
||||
import {
|
||||
type UpdateReleaseChannel,
|
||||
getUpdateReleaseChannel,
|
||||
@@ -157,6 +165,16 @@ const CONFIG_ITEMS: ConfigItem[] = [
|
||||
"关闭后,前端停止自动刷新,后端停止实时隧道质量探测(全局配置)",
|
||||
type: "switch",
|
||||
},
|
||||
{
|
||||
key: TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
label: "隧道质量探测间隔",
|
||||
placeholder: String(DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC),
|
||||
description:
|
||||
"设置实时隧道质量检测的执行频率,单位为秒;允许 1–3600 秒,默认 1 秒。",
|
||||
type: "input",
|
||||
dependsOn: "monitor_tunnel_quality_enabled",
|
||||
dependsValue: "true",
|
||||
},
|
||||
{
|
||||
key: "monitor_retention_days",
|
||||
label: "监控数据保留天数",
|
||||
@@ -239,6 +257,7 @@ const getInitialConfigs = (): Record<string, string> => {
|
||||
"cloudflare_secret_key",
|
||||
"forward_compact_mode",
|
||||
"monitor_tunnel_quality_enabled",
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
"monitor_retention_days",
|
||||
"ip",
|
||||
"panel_domain",
|
||||
@@ -622,6 +641,19 @@ export default function ConfigPage() {
|
||||
|
||||
// 保存配置
|
||||
const handleSave = async () => {
|
||||
const intervalValue = configs[TUNNEL_QUALITY_INTERVAL_CONFIG_KEY];
|
||||
const intervalChanged =
|
||||
intervalValue !== originalConfigs[TUNNEL_QUALITY_INTERVAL_CONFIG_KEY];
|
||||
const intervalError = intervalChanged
|
||||
? validateTunnelQualityInterval(intervalValue || "")
|
||||
: null;
|
||||
|
||||
if (intervalError) {
|
||||
toast.error(intervalError);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
setSaving(true);
|
||||
try {
|
||||
const changedKeys = Object.keys(configs).filter(
|
||||
@@ -667,12 +699,22 @@ export default function ConfigPage() {
|
||||
}),
|
||||
);
|
||||
|
||||
// 如果隧道质量检测开关变更,通知 tunnel-monitor-view
|
||||
if (changedKeys.includes("monitor_tunnel_quality_enabled")) {
|
||||
// 如果隧道质量检测配置变更,通知 tunnel-monitor-view
|
||||
if (
|
||||
changedKeys.some((key) =>
|
||||
[
|
||||
"monitor_tunnel_quality_enabled",
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
].includes(key),
|
||||
)
|
||||
) {
|
||||
window.dispatchEvent(
|
||||
new CustomEvent("monitorTunnelQualityEnabledChanged", {
|
||||
detail: {
|
||||
enabled: configs["monitor_tunnel_quality_enabled"] === "true",
|
||||
intervalSec: parseTunnelQualityIntervalSeconds(
|
||||
configs[TUNNEL_QUALITY_INTERVAL_CONFIG_KEY],
|
||||
),
|
||||
},
|
||||
}),
|
||||
);
|
||||
@@ -1075,11 +1117,21 @@ export default function ConfigPage() {
|
||||
case "bg_image":
|
||||
return renderBgImageUploader();
|
||||
|
||||
case "input":
|
||||
case "input": {
|
||||
if (isBrandPreviewKey(item.key)) {
|
||||
return renderBrandAssetUploader(item.key, isChanged);
|
||||
}
|
||||
|
||||
const isTunnelQualityInterval =
|
||||
item.key === TUNNEL_QUALITY_INTERVAL_CONFIG_KEY;
|
||||
const intervalValue = isTunnelQualityInterval
|
||||
? (configs[item.key] ?? String(DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC))
|
||||
: (configs[item.key] ?? "");
|
||||
const intervalError =
|
||||
isTunnelQualityInterval && configs[item.key] !== undefined
|
||||
? validateTunnelQualityInterval(intervalValue)
|
||||
: null;
|
||||
|
||||
return (
|
||||
<Input
|
||||
classNames={{
|
||||
@@ -1091,14 +1143,30 @@ export default function ConfigPage() {
|
||||
description={
|
||||
isCommercialDisabled ? "需商业版授权才能修改此项" : undefined
|
||||
}
|
||||
endContent={isTunnelQualityInterval ? "秒" : undefined}
|
||||
errorMessage={intervalError || undefined}
|
||||
isDisabled={isCommercialDisabled}
|
||||
isInvalid={Boolean(intervalError)}
|
||||
max={
|
||||
isTunnelQualityInterval
|
||||
? MAX_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
: undefined
|
||||
}
|
||||
min={
|
||||
isTunnelQualityInterval
|
||||
? MIN_TUNNEL_QUALITY_INTERVAL_SEC
|
||||
: undefined
|
||||
}
|
||||
placeholder={item.placeholder}
|
||||
size="md"
|
||||
value={configs[item.key] || ""}
|
||||
step={isTunnelQualityInterval ? 1 : undefined}
|
||||
type={isTunnelQualityInterval ? "number" : "text"}
|
||||
value={intervalValue}
|
||||
variant="bordered"
|
||||
onChange={(e) => handleConfigChange(item.key, e.target.value)}
|
||||
/>
|
||||
);
|
||||
}
|
||||
|
||||
case "switch":
|
||||
return (
|
||||
|
||||
@@ -69,6 +69,7 @@ import {
|
||||
getNodeList,
|
||||
pauseForwardService,
|
||||
resumeForwardService,
|
||||
resetForwardFlow,
|
||||
diagnoseForward,
|
||||
updateForwardOrder,
|
||||
getConfigByName,
|
||||
@@ -230,7 +231,7 @@ const FORWARD_GROUPED_TABLE_COLUMN_CLASS = {
|
||||
strategy: "w-[100px]",
|
||||
totalFlow: "w-[120px]",
|
||||
status: "w-[100px]",
|
||||
actions: "w-[144px] text-right",
|
||||
actions: "w-[176px] text-right",
|
||||
} as const;
|
||||
|
||||
const normalizeForwardUserName = (userName?: string): string => {
|
||||
@@ -764,6 +765,7 @@ const SortableTableRow = ({
|
||||
handleEdit,
|
||||
handleDelete,
|
||||
handleDiagnose,
|
||||
handleResetFlow,
|
||||
showAddressModal,
|
||||
formatFlow,
|
||||
}: any) => {
|
||||
@@ -919,6 +921,29 @@ const SortableTableRow = ({
|
||||
/>
|
||||
</svg>
|
||||
</Button>
|
||||
<Button
|
||||
isIconOnly
|
||||
className="bg-secondary/10 text-secondary hover:bg-secondary/20"
|
||||
isDisabled={(forward.inFlow || 0) + (forward.outFlow || 0) <= 0}
|
||||
size="sm"
|
||||
title="流量清零"
|
||||
onPress={() => handleResetFlow(forward)}
|
||||
>
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
className="h-4 w-4"
|
||||
fill="none"
|
||||
stroke="currentColor"
|
||||
viewBox="0 0 24 24"
|
||||
>
|
||||
<path
|
||||
d="M4 4v6h6M20 20v-6h-6M20 9a8 8 0 00-13.657-3.657L4 8m16 8-2.343 2.657A8 8 0 014 15"
|
||||
strokeLinecap="round"
|
||||
strokeLinejoin="round"
|
||||
strokeWidth={2}
|
||||
/>
|
||||
</svg>
|
||||
</Button>
|
||||
<Button
|
||||
isIconOnly
|
||||
className="bg-danger/10 text-danger hover:bg-danger/20"
|
||||
@@ -958,6 +983,7 @@ const SortableCompactTableRow = ({
|
||||
handleEdit,
|
||||
handleDelete,
|
||||
handleDiagnose,
|
||||
handleResetFlow,
|
||||
showAddressModal,
|
||||
hasMultipleAddresses,
|
||||
formatFlow,
|
||||
@@ -1144,6 +1170,29 @@ const SortableCompactTableRow = ({
|
||||
/>
|
||||
</svg>
|
||||
</Button>
|
||||
<Button
|
||||
isIconOnly
|
||||
className="bg-secondary/10 text-secondary hover:bg-secondary/20"
|
||||
isDisabled={(forward.inFlow || 0) + (forward.outFlow || 0) <= 0}
|
||||
size="sm"
|
||||
title="流量清零"
|
||||
onPress={() => handleResetFlow(forward)}
|
||||
>
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
className="h-4 w-4"
|
||||
fill="none"
|
||||
stroke="currentColor"
|
||||
viewBox="0 0 24 24"
|
||||
>
|
||||
<path
|
||||
d="M4 4v6h6M20 20v-6h-6M20 9a8 8 0 00-13.657-3.657L4 8m16 8-2.343 2.657A8 8 0 014 15"
|
||||
strokeLinecap="round"
|
||||
strokeLinejoin="round"
|
||||
strokeWidth={2}
|
||||
/>
|
||||
</svg>
|
||||
</Button>
|
||||
<Button
|
||||
isIconOnly
|
||||
className="bg-danger/10 text-danger hover:bg-danger/20"
|
||||
@@ -1289,13 +1338,18 @@ export default function ForwardPage() {
|
||||
const [modalOpen, setModalOpen] = useState(false);
|
||||
// isFilterModalOpen removed
|
||||
const [deleteModalOpen, setDeleteModalOpen] = useState(false);
|
||||
const [resetFlowModalOpen, setResetFlowModalOpen] = useState(false);
|
||||
const [addressModalOpen, setAddressModalOpen] = useState(false);
|
||||
const [diagnosisModalOpen, setDiagnosisModalOpen] = useState(false);
|
||||
const [isEdit, setIsEdit] = useState(false);
|
||||
const [submitLoading, setSubmitLoading] = useState(false);
|
||||
const [deleteLoading, setDeleteLoading] = useState(false);
|
||||
const [resetFlowLoading, setResetFlowLoading] = useState(false);
|
||||
const [diagnosisLoading, setDiagnosisLoading] = useState(false);
|
||||
const [forwardToDelete, setForwardToDelete] = useState<Forward | null>(null);
|
||||
const [forwardToResetFlow, setForwardToResetFlow] = useState<Forward | null>(
|
||||
null,
|
||||
);
|
||||
const [currentDiagnosisForward, setCurrentDiagnosisForward] =
|
||||
useState<Forward | null>(null);
|
||||
const [diagnosisResult, setDiagnosisResult] =
|
||||
@@ -2171,7 +2225,8 @@ export default function ForwardPage() {
|
||||
maxConn: forward.maxConn ?? 0,
|
||||
proxyProtocol: forward.proxyProtocol ?? 0,
|
||||
proxyProtocolReceive: forward.proxyProtocolReceive ?? 0,
|
||||
proxyProtocolSend: forward.proxyProtocolSend ?? forward.proxyProtocol ?? 0,
|
||||
proxyProtocolSend:
|
||||
forward.proxyProtocolSend ?? forward.proxyProtocol ?? 0,
|
||||
});
|
||||
setErrors({});
|
||||
setModalOpen(true);
|
||||
@@ -2183,6 +2238,46 @@ export default function ForwardPage() {
|
||||
setDeleteModalOpen(true);
|
||||
};
|
||||
|
||||
const handleResetFlow = (forward: Forward) => {
|
||||
if ((forward.inFlow || 0) + (forward.outFlow || 0) <= 0) return;
|
||||
|
||||
setForwardToResetFlow(forward);
|
||||
setResetFlowModalOpen(true);
|
||||
};
|
||||
|
||||
const handleResetFlowModalOpenChange = (isOpen: boolean) => {
|
||||
if (resetFlowLoading) return;
|
||||
|
||||
setResetFlowModalOpen(isOpen);
|
||||
if (!isOpen) {
|
||||
setForwardToResetFlow(null);
|
||||
}
|
||||
};
|
||||
|
||||
const confirmResetFlow = async () => {
|
||||
if (!forwardToResetFlow) return;
|
||||
|
||||
setResetFlowLoading(true);
|
||||
try {
|
||||
const res = await resetForwardFlow(forwardToResetFlow.id);
|
||||
|
||||
if (res.code !== 0) {
|
||||
toast.error(res.msg || "流量清零失败");
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
toast.success("规则流量已清零");
|
||||
setResetFlowModalOpen(false);
|
||||
setForwardToResetFlow(null);
|
||||
await refreshForwardList(false);
|
||||
} catch {
|
||||
toast.error("流量清零失败");
|
||||
} finally {
|
||||
setResetFlowLoading(false);
|
||||
}
|
||||
};
|
||||
|
||||
// 确认删除规则
|
||||
const confirmDelete = async () => {
|
||||
if (!forwardToDelete) return;
|
||||
@@ -3957,7 +4052,7 @@ export default function ForwardPage() {
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div className="flex gap-1.5 mt-3">
|
||||
<div className="grid grid-cols-2 gap-1.5 mt-3">
|
||||
<Button
|
||||
className="flex-1 min-h-8"
|
||||
color="primary"
|
||||
@@ -4000,6 +4095,32 @@ export default function ForwardPage() {
|
||||
>
|
||||
诊断
|
||||
</Button>
|
||||
<Button
|
||||
className="flex-1 min-h-8"
|
||||
color="secondary"
|
||||
isDisabled={(forward.inFlow || 0) + (forward.outFlow || 0) <= 0}
|
||||
size="sm"
|
||||
startContent={
|
||||
<svg
|
||||
aria-hidden="true"
|
||||
className="w-3 h-3"
|
||||
fill="none"
|
||||
stroke="currentColor"
|
||||
viewBox="0 0 24 24"
|
||||
>
|
||||
<path
|
||||
d="M4 4v6h6M20 20v-6h-6M20 9a8 8 0 00-13.657-3.657L4 8m16 8-2.343 2.657A8 8 0 014 15"
|
||||
strokeLinecap="round"
|
||||
strokeLinejoin="round"
|
||||
strokeWidth={2}
|
||||
/>
|
||||
</svg>
|
||||
}
|
||||
variant="flat"
|
||||
onPress={() => handleResetFlow(forward)}
|
||||
>
|
||||
清零
|
||||
</Button>
|
||||
<Button
|
||||
className="flex-1 min-h-8"
|
||||
color="danger"
|
||||
@@ -4302,7 +4423,7 @@ export default function ForwardPage() {
|
||||
<TableColumn className="w-[80px]">策略</TableColumn>
|
||||
<TableColumn className="w-[100px]">用量</TableColumn>
|
||||
<TableColumn className="w-[80px]">状态</TableColumn>
|
||||
<TableColumn align="left" className="w-[120px] pl-4">
|
||||
<TableColumn align="left" className="w-[160px] pl-4">
|
||||
操作
|
||||
</TableColumn>
|
||||
</TableHeader>
|
||||
@@ -4321,6 +4442,7 @@ export default function ForwardPage() {
|
||||
handleDelete={handleDelete}
|
||||
handleDiagnose={handleDiagnose}
|
||||
handleEdit={handleEdit}
|
||||
handleResetFlow={handleResetFlow}
|
||||
handleServiceToggle={handleServiceToggle}
|
||||
hasMultipleAddresses={hasMultipleAddresses}
|
||||
selectMode={selectMode}
|
||||
@@ -4616,6 +4738,7 @@ export default function ForwardPage() {
|
||||
handleDelete={handleDelete}
|
||||
handleDiagnose={handleDiagnose}
|
||||
handleEdit={handleEdit}
|
||||
handleResetFlow={handleResetFlow}
|
||||
handleServiceToggle={
|
||||
handleServiceToggle
|
||||
}
|
||||
@@ -5172,6 +5295,59 @@ export default function ForwardPage() {
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
|
||||
{/* 规则流量清零确认模态框 */}
|
||||
<Modal
|
||||
backdrop="blur"
|
||||
classNames={{
|
||||
base: "!w-[calc(100%-32px)] !mx-auto sm:!w-full rounded-2xl overflow-hidden",
|
||||
}}
|
||||
isOpen={resetFlowModalOpen}
|
||||
placement="center"
|
||||
scrollBehavior="inside"
|
||||
size="lg"
|
||||
onOpenChange={handleResetFlowModalOpenChange}
|
||||
>
|
||||
<ModalContent>
|
||||
{(onClose) => (
|
||||
<>
|
||||
<ModalHeader className="flex flex-col gap-1">
|
||||
<h2 className="text-lg font-bold text-secondary">
|
||||
确认流量清零
|
||||
</h2>
|
||||
</ModalHeader>
|
||||
<ModalBody>
|
||||
<p className="text-default-600">
|
||||
确定要清零规则{" "}
|
||||
<span className="font-semibold text-foreground">
|
||||
"{forwardToResetFlow?.name}"
|
||||
</span>{" "}
|
||||
当前显示的上传和下载流量吗?
|
||||
</p>
|
||||
<p className="text-small text-default-500 mt-2">
|
||||
此操作不可撤销,但不会影响用户总流量、用户隧道配额和历史统计。
|
||||
</p>
|
||||
</ModalBody>
|
||||
<ModalFooter>
|
||||
<Button
|
||||
isDisabled={resetFlowLoading}
|
||||
variant="light"
|
||||
onPress={onClose}
|
||||
>
|
||||
取消
|
||||
</Button>
|
||||
<Button
|
||||
color="secondary"
|
||||
isLoading={resetFlowLoading}
|
||||
onPress={confirmResetFlow}
|
||||
>
|
||||
确认清零
|
||||
</Button>
|
||||
</ModalFooter>
|
||||
</>
|
||||
)}
|
||||
</ModalContent>
|
||||
</Modal>
|
||||
|
||||
{/* 地址列表弹窗 */}
|
||||
<Modal
|
||||
classNames={{
|
||||
|
||||
@@ -53,12 +53,17 @@ import {
|
||||
TableRow,
|
||||
TableCell,
|
||||
} from "@/shadcn-bridge/heroui/table";
|
||||
import {
|
||||
DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
parseTunnelQualityIntervalSeconds,
|
||||
TUNNEL_QUALITY_INTERVAL_CONFIG_KEY,
|
||||
tunnelQualityIntervalLabel,
|
||||
} from "@/config/tunnel-quality";
|
||||
|
||||
interface TunnelMonitorViewProps {
|
||||
viewMode?: "list" | "grid";
|
||||
}
|
||||
|
||||
const QUALITY_POLL_INTERVAL = 1_000; // 1 second
|
||||
const MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY =
|
||||
"monitor_tunnel_quality_enabled";
|
||||
const MONITOR_TUNNEL_QUALITY_ENABLED_EVENT =
|
||||
@@ -562,6 +567,9 @@ export function TunnelMonitorView({
|
||||
const qualityTimerRef = useRef<number | null>(null);
|
||||
const [monitorTunnelQualityEnabled, setMonitorTunnelQualityEnabled] =
|
||||
useState(true);
|
||||
const [tunnelQualityIntervalSec, setTunnelQualityIntervalSec] = useState(
|
||||
DEFAULT_TUNNEL_QUALITY_INTERVAL_SEC,
|
||||
);
|
||||
|
||||
// Detail view state
|
||||
const [detailTunnelId, setDetailTunnelId] = useState<number | null>(null);
|
||||
@@ -618,26 +626,28 @@ export function TunnelMonitorView({
|
||||
}
|
||||
}, []);
|
||||
|
||||
const loadMonitorTunnelQualityEnabled = useCallback(async () => {
|
||||
try {
|
||||
const response = await getConfigByName(
|
||||
MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY,
|
||||
);
|
||||
const loadTunnelQualityConfig = useCallback(async () => {
|
||||
const [enabledResponse, intervalResponse] = await Promise.all([
|
||||
getConfigByName(MONITOR_TUNNEL_QUALITY_ENABLED_CONFIG_KEY).catch(
|
||||
() => null,
|
||||
),
|
||||
getConfigByName(TUNNEL_QUALITY_INTERVAL_CONFIG_KEY).catch(() => null),
|
||||
]);
|
||||
|
||||
setMonitorTunnelQualityEnabled(
|
||||
typeof response.data?.value === "string"
|
||||
? response.data.value === "true"
|
||||
: true,
|
||||
);
|
||||
} catch {
|
||||
setMonitorTunnelQualityEnabled(true);
|
||||
}
|
||||
setMonitorTunnelQualityEnabled(
|
||||
typeof enabledResponse?.data?.value === "string"
|
||||
? enabledResponse.data.value === "true"
|
||||
: true,
|
||||
);
|
||||
setTunnelQualityIntervalSec(
|
||||
parseTunnelQualityIntervalSeconds(intervalResponse?.data?.value),
|
||||
);
|
||||
}, []);
|
||||
|
||||
useEffect(() => {
|
||||
void loadTunnels();
|
||||
void loadMonitorTunnelQualityEnabled();
|
||||
}, [loadMonitorTunnelQualityEnabled, loadTunnels]);
|
||||
void loadTunnelQualityConfig();
|
||||
}, [loadTunnelQualityConfig, loadTunnels]);
|
||||
|
||||
useEffect(() => {
|
||||
const timer = window.setInterval(() => {
|
||||
@@ -649,8 +659,11 @@ export function TunnelMonitorView({
|
||||
|
||||
useEffect(() => {
|
||||
const handleMonitorTunnelQualityEnabledChanged = (event: Event) => {
|
||||
const enabled = (event as CustomEvent<{ enabled?: boolean }>).detail
|
||||
?.enabled;
|
||||
const detail = (
|
||||
event as CustomEvent<{ enabled?: boolean; intervalSec?: number }>
|
||||
).detail;
|
||||
const enabled = detail?.enabled;
|
||||
const intervalSec = detail?.intervalSec;
|
||||
|
||||
if (typeof enabled === "boolean") {
|
||||
setMonitorTunnelQualityEnabled(enabled);
|
||||
@@ -658,7 +671,12 @@ export function TunnelMonitorView({
|
||||
setQualityLoading(false);
|
||||
}
|
||||
} else {
|
||||
void loadMonitorTunnelQualityEnabled();
|
||||
void loadTunnelQualityConfig();
|
||||
}
|
||||
if (typeof intervalSec === "number") {
|
||||
setTunnelQualityIntervalSec(
|
||||
parseTunnelQualityIntervalSeconds(String(intervalSec)),
|
||||
);
|
||||
}
|
||||
};
|
||||
|
||||
@@ -673,7 +691,7 @@ export function TunnelMonitorView({
|
||||
handleMonitorTunnelQualityEnabledChanged as EventListener,
|
||||
);
|
||||
};
|
||||
}, [loadMonitorTunnelQualityEnabled]);
|
||||
}, [loadTunnelQualityConfig]);
|
||||
|
||||
useEffect(() => {
|
||||
if (tunnels.length > 0 && !initialHistoryFetched.current) {
|
||||
@@ -726,7 +744,7 @@ export function TunnelMonitorView({
|
||||
}
|
||||
}, [tunnels]);
|
||||
|
||||
// --- Load quality snapshots (auto-polling every 10s) ---
|
||||
// --- Load quality snapshots using the configured probe interval ---
|
||||
const loadQuality = useCallback(async (options?: { silent?: boolean }) => {
|
||||
const silent = options?.silent ?? false;
|
||||
|
||||
@@ -788,7 +806,7 @@ export function TunnelMonitorView({
|
||||
|
||||
qualityTimerRef.current = window.setInterval(() => {
|
||||
void loadQuality({ silent: true });
|
||||
}, QUALITY_POLL_INTERVAL);
|
||||
}, tunnelQualityIntervalSec * 1000);
|
||||
|
||||
return () => {
|
||||
if (qualityTimerRef.current) {
|
||||
@@ -796,7 +814,7 @@ export function TunnelMonitorView({
|
||||
qualityTimerRef.current = null;
|
||||
}
|
||||
};
|
||||
}, [loadQuality, monitorTunnelQualityEnabled]);
|
||||
}, [loadQuality, monitorTunnelQualityEnabled, tunnelQualityIntervalSec]);
|
||||
|
||||
// --- Load quality history for detail chart ---
|
||||
const loadQualityHistory = useCallback(
|
||||
@@ -1065,7 +1083,7 @@ export function TunnelMonitorView({
|
||||
{monitorTunnelQualityEnabled ? (
|
||||
<>
|
||||
<LiveDot />
|
||||
<span>自动探测中(每秒测试,30秒上报)</span>
|
||||
<span>{`自动探测中(${tunnelQualityIntervalLabel(tunnelQualityIntervalSec)}测试)`}</span>
|
||||
</>
|
||||
) : (
|
||||
<>
|
||||
@@ -1129,7 +1147,7 @@ export function TunnelMonitorView({
|
||||
{monitorTunnelQualityEnabled ? (
|
||||
<>
|
||||
<LiveDot />
|
||||
<span>每秒探测 · 更新于 {lastQualityUpdate}</span>
|
||||
<span>{`${tunnelQualityIntervalLabel(tunnelQualityIntervalSec)}探测 · 更新于 ${lastQualityUpdate}`}</span>
|
||||
</>
|
||||
) : (
|
||||
<>
|
||||
|
||||
Reference in New Issue
Block a user