diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index 13ae9d6..74378cb 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -290,6 +290,11 @@ func (h *Handler) syncForwardServicesWithWarnings(forward *forwardRecord, method services := buildForwardServiceConfigs(serviceBase, forward, tunnel, node, fp.Port, strings.TrimSpace(fp.InIP), limiterID, tunnelTLSProtocol) _, err = h.sendNodeCommand(node.ID, method, services, true, false) if err != nil && allowFallbackAdd && method == "UpdateService" { + if isNotFoundError(err) { + if delErr := h.deleteForwardServicesOnNode(forward, node.ID); delErr != nil && !isNotFoundError(delErr) { + return warnings, fmt.Errorf("节点 %s 清理旧服务失败: %w", node.Name, delErr) + } + } _, err = h.sendNodeCommand(node.ID, "AddService", services, true, false) } if err != nil && strings.EqualFold(strings.TrimSpace(method), "UpdateService") && isAddressAlreadyInUseError(err) { @@ -437,40 +442,26 @@ func (h *Handler) controlForwardServices(forward *forwardRecord, commandType str candidateTunnelIDs = append(candidateTunnelIDs, allUserTunnelIDs...) bases := buildForwardServiceBaseCandidates(forward.ID, forward.UserID, userTunnelID, candidateTunnelIDs) seen := map[int64]struct{}{} + healed := false for _, fp := range ports { if _, ok := seen[fp.NodeID]; ok { continue } seen[fp.NodeID] = struct{}{} - var lastNotFoundErr error - nodeHandled := false + nodeHandled, lastNotFoundErr, err := h.controlForwardServicesOnNode(fp.NodeID, bases, commandType) + if err != nil { + return err + } - for _, base := range bases { - variants := []string{base + "_tcp", base + "_udp"} - if shouldTryLegacySingleService(commandType) || strings.EqualFold(strings.TrimSpace(commandType), "DeleteService") { - variants = append(variants, base) + if !nodeHandled && lastNotFoundErr != nil && !healed && shouldSelfHealForwardServiceControl(commandType) { + if healErr := h.syncForwardServices(forward, "UpdateService", true); healErr != nil { + return healErr } - - candidateHandled := false - for _, name := range variants { - payload := map[string]interface{}{ - "services": []string{name}, - } - _, err := h.sendNodeCommand(fp.NodeID, commandType, payload, false, false) - if err == nil { - candidateHandled = true - continue - } - if !isNotFoundError(err) { - return err - } - lastNotFoundErr = err - } - - if candidateHandled { - nodeHandled = true - break + healed = true + nodeHandled, lastNotFoundErr, err = h.controlForwardServicesOnNode(fp.NodeID, bases, commandType) + if err != nil { + return err } } @@ -488,6 +479,49 @@ func (h *Handler) controlForwardServices(forward *forwardRecord, commandType str return nil } +func (h *Handler) controlForwardServicesOnNode(nodeID int64, bases []string, commandType string) (bool, error, error) { + return controlForwardServiceCommand(bases, commandType, func(name string) error { + payload := map[string]interface{}{ + "services": []string{name}, + } + _, err := h.sendNodeCommand(nodeID, commandType, payload, false, false) + return err + }) +} + +func controlForwardServiceCommand(bases []string, commandType string, send func(name string) error) (bool, error, error) { + var lastNotFoundErr error + for _, base := range bases { + variants := []string{base + "_tcp", base + "_udp"} + if shouldTryLegacySingleService(commandType) || strings.EqualFold(strings.TrimSpace(commandType), "DeleteService") { + variants = append(variants, base) + } + + candidateHandled := false + for _, name := range variants { + err := send(name) + if err == nil { + candidateHandled = true + continue + } + if !isNotFoundError(err) { + return false, lastNotFoundErr, err + } + lastNotFoundErr = err + } + + if candidateHandled { + return true, nil, nil + } + } + return false, lastNotFoundErr, nil +} + +func shouldSelfHealForwardServiceControl(commandType string) bool { + cmd := strings.ToLower(strings.TrimSpace(commandType)) + return cmd == "pauseservice" || cmd == "resumeservice" +} + func (h *Handler) applyNodeProtocolChange(nodeID int64, httpVal, tlsVal, socksVal int) error { _, err := h.sendNodeCommand(nodeID, "SetProtocol", map[string]interface{}{ "http": httpVal, diff --git a/go-backend/internal/http/handler/control_plane_test.go b/go-backend/internal/http/handler/control_plane_test.go index 47a414a..380f800 100644 --- a/go-backend/internal/http/handler/control_plane_test.go +++ b/go-backend/internal/http/handler/control_plane_test.go @@ -69,6 +69,78 @@ func TestShouldTryLegacySingleService(t *testing.T) { } } +func TestShouldSelfHealForwardServiceControl(t *testing.T) { + if !shouldSelfHealForwardServiceControl("PauseService") { + t.Fatalf("PauseService should trigger self-heal") + } + if !shouldSelfHealForwardServiceControl(" resumeService ") { + t.Fatalf("ResumeService should trigger self-heal") + } + if shouldSelfHealForwardServiceControl("DeleteService") { + t.Fatalf("DeleteService should not trigger self-heal") + } +} + +func TestControlForwardServiceCommandHandledOnKnownVariant(t *testing.T) { + bases := []string{"12_34_56"} + called := make([]string, 0) + handled, lastNotFoundErr, err := controlForwardServiceCommand(bases, "PauseService", func(name string) error { + called = append(called, name) + if name == "12_34_56_udp" { + return nil + } + return errors.New("service " + name + " not found") + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if !handled { + t.Fatalf("expected handled=true") + } + if lastNotFoundErr != nil { + t.Fatalf("expected lastNotFoundErr=nil when handled") + } + wantCalls := []string{"12_34_56_tcp", "12_34_56_udp", "12_34_56"} + if !reflect.DeepEqual(called, wantCalls) { + t.Fatalf("expected calls %v, got %v", wantCalls, called) + } +} + +func TestControlForwardServiceCommandReturnsLastNotFoundWhenAllMissing(t *testing.T) { + bases := []string{"12_34_56"} + handled, lastNotFoundErr, err := controlForwardServiceCommand(bases, "PauseService", func(name string) error { + return errors.New("service " + name + " not found") + }) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if handled { + t.Fatalf("expected handled=false") + } + if lastNotFoundErr == nil { + t.Fatalf("expected lastNotFoundErr when all variants are missing") + } +} + +func TestControlForwardServiceCommandReturnsHardError(t *testing.T) { + bases := []string{"12_34_56"} + handled, lastNotFoundErr, err := controlForwardServiceCommand(bases, "PauseService", func(name string) error { + if name == "12_34_56_tcp" { + return errors.New("network timeout") + } + return nil + }) + if err == nil { + t.Fatalf("expected hard error") + } + if handled { + t.Fatalf("expected handled=false on hard error") + } + if lastNotFoundErr != nil { + t.Fatalf("did not expect not-found error alongside hard error") + } +} + func TestIsAlreadyExistsMessage(t *testing.T) { if !isAlreadyExistsMessage("service demo already exists") { t.Fatalf("expected already exists message to be tolerated") diff --git a/plans/011-forward-service-legacy-compat-and-agent-rollout.md b/plans/011-forward-service-legacy-compat-and-agent-rollout.md new file mode 100644 index 0000000..33aa46e --- /dev/null +++ b/plans/011-forward-service-legacy-compat-and-agent-rollout.md @@ -0,0 +1,28 @@ +# 011 转发服务名升级兼容与节点滚动升级 + +## 目标 +- 修复旧版本升级后编辑转发/隧道出现 `service not found`(service不存在)的问题。 +- 在后端加入兼容自愈逻辑,允许旧命名与新命名共存过渡。 +- 给出低风险节点升级顺序,避免一次性全量切换带来的中断。 + +## Checklist +- [x] 定位回归路径:服务名从 `forward_user_0` 迁移到真实 `user_tunnel_id` 后,与旧运行态不一致导致控制失败。 +- [x] 在 `UpdateService` 的兼容路径加入旧服务清理后重建逻辑。 +- [x] 在 `Pause/Resume` 控制路径加入首次 not found 后自愈重试逻辑。 +- [x] 增加回归测试覆盖兼容行为。 +- [x] 执行 `go-backend` 相关测试并记录结果。 +- [x] 输出运维侧“后端先行 + agent 灰度升级 + 批量重部署”操作步骤。 + +## 变更说明(实施中) +- 后端控制面将在检测到升级期的服务名不一致时进行自动自愈,降低人工干预和手工重建成本。 + +## 测试记录 +- 命令:`cd go-backend && go test ./internal/http/handler/...` +- 结果:通过。 + +## 运维升级顺序(推荐) +1. 先发布本次后端兼容补丁(无需等待所有 agent 同步升级)。 +2. 按 10%-20% 灰度分批升级 agent(低风险节点 -> 非高峰节点 -> 全量)。 +3. 每批升级后执行一次“转发批量重部署”,将运行态统一到新服务命名。 +4. 观察日志中 `service .* not found` 是否清零,再推进下一批。 +5. 全量稳定后保留兼容逻辑至少一个小版本周期,再评估收敛。