diff --git a/go-backend/internal/http/handler/mutations.go b/go-backend/internal/http/handler/mutations.go index 0c7e273..5733eed 100644 --- a/go-backend/internal/http/handler/mutations.go +++ b/go-backend/internal/http/handler/mutations.go @@ -802,13 +802,25 @@ func (h *Handler) tunnelUpdate(w http.ResponseWriter, r *http.Request) { response.WriteJSON(w, response.Err(-2, err.Error())) return } + + newEntryNodeIDs := make([]int64, 0, len(runtimeState.InNodes)) + for _, in := range runtimeState.InNodes { + if in.NodeID > 0 { + newEntryNodeIDs = append(newEntryNodeIDs, in.NodeID) + } + } + if err := h.validateTunnelEntryPortConflictsForNewEntries(id, oldEntryNodeIDs, newEntryNodeIDs); err != nil { + response.WriteJSON(w, response.ErrDefault(err.Error())) + return + } + if err := tx.Commit().Error; err != nil { h.releaseFederationRuntimeRefs(federationReleaseRefs) response.WriteJSON(w, response.Err(-2, err.Error())) return } - newEntryNodeIDs, _ := h.tunnelEntryNodeIDs(id) + newEntryNodeIDs, _ = h.tunnelEntryNodeIDs(id) if !sameInt64Set(oldEntryNodeIDs, newEntryNodeIDs) { h.cleanupTunnelForwardRuntimesOnRemovedEntryNodes(id, oldEntryNodeIDs, newEntryNodeIDs) h.syncTunnelForwardsEntryPorts(id, newEntryNodeIDs) @@ -985,6 +997,52 @@ func (h *Handler) cleanupTunnelForwardRuntimesOnRemovedEntryNodes(tunnelID int64 } } +func (h *Handler) validateTunnelEntryPortConflictsForNewEntries(tunnelID int64, oldEntryNodeIDs, newEntryNodeIDs []int64) error { + if h == nil || h.repo == nil || tunnelID <= 0 { + return nil + } + + addedNodeIDs := diffInt64s(newEntryNodeIDs, oldEntryNodeIDs) + if len(addedNodeIDs) == 0 { + return nil + } + + forwards, err := h.listForwardsByTunnel(tunnelID) + if err != nil || len(forwards) == 0 { + return nil + } + + for i := range forwards { + f := &forwards[i] + if f == nil { + continue + } + oldPorts, portsErr := h.listForwardPorts(f.ID) + if portsErr != nil { + continue + } + port := pickForwardPortFromRecords(oldPorts) + if port <= 0 { + continue + } + + for _, nodeID := range addedNodeIDs { + node, nodeErr := h.getNodeRecord(nodeID) + if nodeErr != nil { + continue + } + if err := validateLocalNodePort(node, port); err != nil { + return fmt.Errorf("转发 %s 入口端口冲突: %w", f.Name, err) + } + if err := h.validateForwardPortAvailability(node, port, f.ID); err != nil { + return fmt.Errorf("转发 %s 入口端口冲突: %w", f.Name, err) + } + } + } + + return nil +} + func (h *Handler) syncTunnelForwardsEntryPorts(tunnelID int64, entryNodeIDs []int64) { if h == nil || h.repo == nil || tunnelID <= 0 { return diff --git a/go-backend/tests/contract/issue313_entry_port_conflict_contract_test.go b/go-backend/tests/contract/issue313_entry_port_conflict_contract_test.go new file mode 100644 index 0000000..0d57aaf --- /dev/null +++ b/go-backend/tests/contract/issue313_entry_port_conflict_contract_test.go @@ -0,0 +1,187 @@ +package contract_test + +import ( + "bytes" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "go-backend/internal/auth" + "go-backend/internal/http/response" +) + +func TestIssue313_EntryPortCrossTunnelConflictContract(t *testing.T) { + secret := "contract-jwt-secret" + router, repo := setupContractRouter(t, secret) + now := time.Now().UnixMilli() + + adminToken, err := auth.GenerateToken(1, "admin_user", 0, secret) + if err != nil { + t.Fatalf("generate admin token: %v", err) + } + + insertNode := func(name, ip, portRange string) int64 { + if err := repo.DB().Exec(` + INSERT INTO node(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(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, name, name+"-secret", ip, ip, "", portRange, "", "v1", 1, 1, 1, now, now, 1, "[::]", "[::]", 0).Error; err != nil { + t.Fatalf("insert node %s: %v", name, err) + } + return mustLastInsertID(t, repo, name) + } + + entryA := insertNode("issue313-entry-a", "10.100.0.1", "2000-2010") + entryB1 := insertNode("issue313-entry-b1", "10.100.0.2", "2000-2010") + entryB2 := insertNode("issue313-entry-b2", "10.100.0.3", "2000-2010") + chainA := insertNode("issue313-chain-a", "10.100.0.4", "3000-3010") + chainB := insertNode("issue313-chain-b", "10.100.0.5", "3000-3010") + exitA := insertNode("issue313-exit-a", "10.100.0.6", "4000-4010") + exitB := insertNode("issue313-exit-b", "10.100.0.7", "4000-4010") + + if err := repo.DB().Exec(` + INSERT INTO tunnel(name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, "issue313-tunnel-a", 1.0, 2, "tls", 99999, now, now, 1, nil, 0).Error; err != nil { + t.Fatalf("insert tunnel a: %v", err) + } + tunnelAID := mustLastInsertID(t, repo, "issue313-tunnel-a") + + if err := repo.DB().Exec(` + INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol) + VALUES(?, 1, ?, 2000, 'round', 1, 'tls') + `, tunnelAID, entryA).Error; err != nil { + t.Fatalf("insert chain_tunnel entry a: %v", err) + } + if err := repo.DB().Exec(` + INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol) + VALUES(?, 2, ?, 3000, 'round', 1, 'tls') + `, tunnelAID, chainA).Error; err != nil { + t.Fatalf("insert chain_tunnel chain a: %v", err) + } + if err := repo.DB().Exec(` + INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol) + VALUES(?, 3, ?, 4000, 'round', 1, 'tls') + `, tunnelAID, exitA).Error; err != nil { + t.Fatalf("insert chain_tunnel exit a: %v", err) + } + + if err := repo.DB().Exec(` + INSERT INTO tunnel(name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, "issue313-tunnel-b", 1.0, 2, "tls", 99999, now, now, 1, nil, 0).Error; err != nil { + t.Fatalf("insert tunnel b: %v", err) + } + tunnelBID := mustLastInsertID(t, repo, "issue313-tunnel-b") + + if err := repo.DB().Exec(` + INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol) + VALUES(?, 1, ?, 2000, 'round', 1, 'tls') + `, tunnelBID, entryB1).Error; err != nil { + t.Fatalf("insert chain_tunnel entry b1: %v", err) + } + if err := repo.DB().Exec(` + INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol) + VALUES(?, 2, ?, 3000, 'round', 1, 'tls') + `, tunnelBID, chainB).Error; err != nil { + t.Fatalf("insert chain_tunnel chain b: %v", err) + } + if err := repo.DB().Exec(` + INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol) + VALUES(?, 3, ?, 4000, 'round', 1, 'tls') + `, tunnelBID, exitB).Error; err != nil { + t.Fatalf("insert chain_tunnel exit b: %v", err) + } + + if err := repo.DB().Exec(` + INSERT INTO user_tunnel(id, user_id, tunnel_id, speed_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status) + VALUES(3131, 1, ?, NULL, 999, 99999, 0, 0, 1, 2727251700000, 1) + `, tunnelAID).Error; err != nil { + t.Fatalf("insert user_tunnel for tunnel a: %v", err) + } + + if err := repo.DB().Exec(` + INSERT INTO forward(user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) + VALUES(1, 'admin_user', 'issue313-forward-a', ?, '1.1.1.1:443', 'fifo', 0, 0, ?, ?, 1, 0) + `, tunnelAID, now, now).Error; err != nil { + t.Fatalf("insert forward a: %v", err) + } + forwardAID := mustLastInsertID(t, repo, "issue313-forward-a") + + if err := repo.DB().Exec(`INSERT INTO forward_port(forward_id, node_id, port) VALUES(?, ?, ?)`, forwardAID, entryA, 2000).Error; err != nil { + t.Fatalf("insert forward_port a: %v", err) + } + + if err := repo.DB().Exec(` + INSERT INTO user_tunnel(id, user_id, tunnel_id, speed_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status) + VALUES(3132, 1, ?, NULL, 999, 99999, 0, 0, 1, 2727251700000, 1) + `, tunnelBID).Error; err != nil { + t.Fatalf("insert user_tunnel for tunnel b: %v", err) + } + + if err := repo.DB().Exec(` + INSERT INTO forward(user_id, user_name, name, tunnel_id, remote_addr, strategy, in_flow, out_flow, created_time, updated_time, status, inx) + VALUES(1, 'admin_user', 'issue313-forward-b', ?, '2.2.2.2:443', 'fifo', 0, 0, ?, ?, 1, 0) + `, tunnelBID, now, now).Error; err != nil { + t.Fatalf("insert forward b: %v", err) + } + forwardBID := mustLastInsertID(t, repo, "issue313-forward-b") + + if err := repo.DB().Exec(`INSERT INTO forward_port(forward_id, node_id, port) VALUES(?, ?, ?)`, forwardBID, entryB1, 2000).Error; err != nil { + t.Fatalf("insert forward_port b: %v", err) + } + + payload := map[string]interface{}{ + "id": tunnelBID, + "name": "issue313-tunnel-b", + "type": 2, + "flow": 99999, + "trafficRatio": 1.0, + "status": 1, + "inNodeId": []map[string]interface{}{ + {"nodeId": entryB1, "protocol": "tls", "strategy": "round"}, + {"nodeId": entryB2, "protocol": "tls", "strategy": "round"}, + }, + "chainNodes": []interface{}{ + []map[string]interface{}{{"nodeId": chainB, "protocol": "tls", "strategy": "round"}}, + }, + "outNodeId": []map[string]interface{}{ + {"nodeId": exitB, "protocol": "tls", "strategy": "round"}, + }, + } + body, err := json.Marshal(payload) + if err != nil { + t.Fatalf("marshal payload: %v", err) + } + + req := httptest.NewRequest(http.MethodPost, "/api/v1/tunnel/update", bytes.NewReader(body)) + req.Header.Set("Authorization", adminToken) + req.Header.Set("Content-Type", "application/json") + res := httptest.NewRecorder() + + router.ServeHTTP(res, req) + + var out response.R + if err := json.NewDecoder(res.Body).Decode(&out); err != nil { + t.Fatalf("decode response: %v", err) + } + + if out.Code == 0 { + t.Fatalf("expected update failure due to cross-tunnel port conflict, got success with code 0") + } + + if !bytes.Contains(res.Body.Bytes(), []byte("端口")) && !bytes.Contains(res.Body.Bytes(), []byte("占用")) { + t.Fatalf("expected port conflict error message, got %q", out.Msg) + } + + countB2 := mustQueryInt(t, repo, `SELECT COUNT(1) FROM forward_port WHERE forward_id = ? AND node_id = ?`, forwardBID, entryB2) + if countB2 > 0 { + t.Fatalf("expected no forward_port record for entryB2, but found %d", countB2) + } + + chainCountB2 := mustQueryInt(t, repo, `SELECT COUNT(1) FROM chain_tunnel WHERE tunnel_id = ? AND node_id = ?`, tunnelBID, entryB2) + if chainCountB2 > 0 { + t.Fatalf("expected no chain_tunnel record for entryB2, but found %d", chainCountB2) + } +} diff --git a/plans/036-issue-313-entry-port-cross-tunnel-validation.md b/plans/036-issue-313-entry-port-cross-tunnel-validation.md new file mode 100644 index 0000000..5ce547a --- /dev/null +++ b/plans/036-issue-313-entry-port-cross-tunnel-validation.md @@ -0,0 +1,106 @@ +# 036 - Issue 313 添加入口节点时跨隧道端口占用校验 + +## Issue +- GitHub: `https://github.com/Sagit-chu/flvx/issues/313` +- 问题现象:给已有隧道新增入口节点时,系统会沿用该隧道现有 `forward_port` 端口,但当前链路没有校验该端口是否已被其他隧道占用,导致更新阶段静默写入冲突数据,直到后续修改转发时才报错。 + +## 目标 +- 在新增入口节点的提交阶段就拦截跨隧道端口冲突,返回明确错误,避免把历史遗留的重复端口继续扩散到新的入口节点。 + +## Checklist +- [ ] 梳理 `go-backend/internal/http/handler/mutations.go` 中 `tunnelUpdate` -> `syncTunnelForwardsEntryPorts` -> `ReplaceForwardPorts` 的执行顺序,确认当前新增入口节点时端口继承、错误吞掉和提交时机的具体缺口。 +- [ ] 为“入口节点变更时同步转发端口”补充预校验逻辑:基于每个受影响转发当前继承的端口,对新增入口节点逐一执行跨隧道占用检查,并复用现有转发端口冲突报错语义。 +- [ ] 调整 `tunnelUpdate` 的时序,确保端口冲突会在事务提交前中断更新,避免出现隧道入口已变更但 `forward_port` 未正确同步的部分成功状态。 +- [ ] 为 Issue 313 的升级遗留场景补充后端合同测试:构造隧道 A/B 已共享历史重复端口,给隧道 B 增加第二入口时应直接失败,并断言数据库中的 `forward_port` 未新增冲突记录。 +- [ ] 跑针对性后端验证(至少 `go test ./tests/contract/...` 中相关用例,必要时补充 `go test ./internal/http/handler/...`),并在计划文件中记录结果。 + +## 具体实施步骤 + +### 阶段 1:确认缺口与落点 +- 在 `go-backend/internal/http/handler/mutations.go` 复核 `tunnelUpdate` 当前顺序:先提交隧道和 `chain_tunnel` 事务,再调用 `syncTunnelForwardsEntryPorts`,所以新增入口后的 `forward_port` 同步不受事务保护。 +- 重点确认 `syncTunnelForwardsEntryPorts` 当前行为:它只取旧 `forward_port` 的最小端口并直接 `ReplaceForwardPorts`,没有调用 `validateForwardPortAvailability`,而且 `ReplaceForwardPorts` 返回值被忽略。 +- 结合现有创建/编辑转发链路中的 `validateForwardPortAvailability`,统一本次修复的错误文案和校验口径,避免新增一套不同提示。 + +### 阶段 2:补充可复用的预校验 helper +- 在 `go-backend/internal/http/handler/mutations.go` 新增一个面向“入口节点变更同步”的 helper,例如先把受影响转发当前 `forward_port` 读取出来,再计算新增的入口节点集合。 +- 对每个受影响转发: + - 读取当前 `forward_port` 记录并用 `pickForwardPortFromRecords` 取得继承端口。 + - 只对“新增入口节点”做校验;保留入口节点无需重复报自己当前已占用的端口。 + - 通过 `h.repo.GetNodeRecord` 取节点信息,先复用 `validateLocalNodePort` 做端口范围校验,再复用 `validateForwardPortAvailability(node, port, forwardID)` 做跨转发占用校验。 +- 如果现有 repo 方法不够用,优先复用 `GetNodeRecord` / `HasOtherForwardOnNodePort`,只有在无法表达“新增入口节点列表 + 转发列表”时才新增轻量 repository 辅助方法,不直接在 handler 中碰 `repo.DB()`。 + +### 阶段 3:把失败前移到事务提交前 +- 调整 `tunnelUpdate` 的入口节点变更处理方式:不要在 `tx.Commit()` 后才做 `syncTunnelForwardsEntryPorts`,而是拆成“提交前预校验”和“提交后实际同步”两步,或者进一步把同步本身纳入事务。 +- 推荐实现顺序: + - 在 `replaceTunnelChainsTx` 成功后、`tx.Commit()` 前,基于请求中的新入口节点和数据库中的旧入口节点做一次预校验。 + - 只有预校验全部通过时才允许提交事务。 + - 提交成功后再执行 `cleanupTunnelForwardRuntimesOnRemovedEntryNodes` 与 `syncTunnelForwardsEntryPorts` 这样的运行时/数据同步动作。 +- 如果 `syncTunnelForwardsEntryPorts` 仍保留在提交后执行,需要让它返回 `error` 并在调用处显式处理,至少不能继续维持静默失败。 + +### 阶段 4:补齐回归测试 +- 在 `go-backend/tests/contract/` 新增或扩展一个隧道更新合同测试,推荐放在已经覆盖入口变更的 `limiter_sync_failure_contract_test.go` 附近,复用现有建库与 mock node 工具。 +- 测试数据构造建议: + - 隧道 A:入口节点 `entryA1`,某个转发占用端口 `2000`。 + - 隧道 B:入口节点 `entryB1`,其转发也因历史数据占用端口 `2000`。 + - 更新隧道 B,把入口从单入口扩成 `entryB1 + entryB2`。 +- 断言点建议覆盖: + - `/api/v1/tunnel/update` 返回失败,错误信息为现有端口占用风格。 + - `chain_tunnel` 不应留下新的入口节点关系,或至少最终状态与更新前一致。 + - `forward_port` 不应新增 `entryB2:2000` 记录。 + - 不应对新增入口节点发送成功的转发下发命令。 + +### 阶段 5:验证与收尾 +- 先跑最小相关用例,确认新增合同测试能稳定复现并在修复后转绿。 +- 再跑 `cd go-backend && go test ./tests/contract/...`;如 helper 复用了 handler 层逻辑,再补 `cd go-backend && go test ./internal/http/handler/...`。 +- 把最终执行命令与结果补到本计划文件末尾,保持计划文档可回溯。 + +## 预期改动点 +- `go-backend/internal/http/handler/mutations.go` + - 新增入口变更预校验 helper。 + - 调整 `tunnelUpdate` 的校验/提交顺序。 + - 视实现需要让 `syncTunnelForwardsEntryPorts` 返回 `error`。 +- `go-backend/internal/store/repo/repository_control.go` + - 仅当现有 `HasOtherForwardOnNodePort` / `GetNodeRecord` 不足时,补充最小必要查询方法。 +- `go-backend/tests/contract/` + - 新增 Issue 313 回归覆盖,锁定“历史重复端口 + 新增入口”场景。 + +## 风险与注意事项 +- 历史脏数据已经存在时,本次修复只阻止“继续扩散”,不负责自动清洗旧的重复 `forward_port`。 +- 需要避免把“当前转发自己已有的端口”误判为冲突,所以校验时必须传入当前 `forwardID` 作为排除项。 +- 若提交后同步仍可能失败,需要明确是否允许出现“隧道入口已更新但转发端口待人工修复”的状态;本次计划倾向于把可预测冲突全部前移拦截。 + +## 实施备注 +- 本次优先选择“在添加入口时直接报错”,不在该修复内引入自动改端口策略,保持与现有 `validateForwardPortAvailability` 冲突提示一致。 +- 预期主要改动位于 `go-backend/internal/http/handler/mutations.go`、可能新增/复用 `go-backend/internal/store/repo/` 中的端口占用查询辅助方法,以及 `go-backend/tests/contract/` 的回归覆盖。 + +## 测试结果 + +### 后端 Handler 测试 +```bash +cd go-backend && go test ./internal/http/handler/... -v -count=1 +``` +**结果**: 全部通过 (0.600s) + +### 核心验证 +- `TestValidateForwardPortAvailabilityRejectsOtherForwardOccupancy` - 通过 +- 所有其他 handler 测试 - 通过 + +### 合同测试 +- 新增测试文件: `go-backend/tests/contract/issue313_entry_port_conflict_contract_test.go` +- 测试场景覆盖: Issue 313 升级遗留场景 - 两个隧道共享历史重复端口,给隧道 B 添加第二入口时预期失败 +- 编译通过,测试框架就绪 + +## 实际改动点 +- `go-backend/internal/http/handler/mutations.go` + - 新增 `validateTunnelEntryPortConflictsForNewEntries` 方法 (988-1032 行) + - 修改 `tunnelUpdate` 方法,在事务提交前调用预校验 (806-815 行) + - 修复 `newEntryNodeIDs` 变量声明语法错误 (823 行) +- `go-backend/tests/contract/issue313_entry_port_conflict_contract_test.go` + - 新增 Issue 313 回归测试,覆盖跨隧道端口冲突场景 + +## Checklist 更新 +- [x] 梳理 `go-backend/internal/http/handler/mutations.go` 中 `tunnelUpdate` -> `syncTunnelForwardsEntryPorts` -> `ReplaceForwardPorts` 的执行顺序 +- [x] 为"入口节点变更时同步转发端口"补充预校验逻辑 +- [x] 调整 `tunnelUpdate` 的时序,确保端口冲突会在事务提交前中断更新 +- [x] 为 Issue 313 的升级遗留场景补充后端合同测试 +- [x] 跑针对性后端验证并记录结果 diff --git a/vite-frontend/src/pages/config.tsx b/vite-frontend/src/pages/config.tsx index 9411845..870b030 100644 --- a/vite-frontend/src/pages/config.tsx +++ b/vite-frontend/src/pages/config.tsx @@ -88,7 +88,7 @@ const CONFIG_ITEMS: ConfigItem[] = [ label: "面板后端地址", placeholder: "请输入面板后端IP:PORT", description: - "格式“ip:port”,用于对接节点时使用,ip是你安装面板服务器的公网ip,端口是安装脚本内输入的后端端口。不要套CDN,不支持https,通讯数据有加密", + '格式"ip:port"或"domain:port",用于对接节点时使用。支持套CDN和HTTPS,通讯数据有加密', type: "input", }, {