fix: add retry mechanism for tunnel service bind conflicts

When tunnel services encounter 'address already in use' errors during
creation/update, automatically cleanup and retry once instead of failing
immediately. This handles race conditions during rapid tunnel reconfiguration.

Entire-Checkpoint: 39e6fb9de836
This commit is contained in:
sagitchu
2026-03-07 16:21:44 +08:00
parent 9c0e7341c3
commit 87479c2ac1
6 changed files with 549 additions and 2 deletions
@@ -273,6 +273,75 @@ func TestIsCannotAssignRequestedAddressError(t *testing.T) {
}
}
func TestRetryTunnelServiceAddWithCleanupRetriesOnAddressInUse(t *testing.T) {
addCalls := 0
cleanupCalls := 0
err := retryTunnelServiceAddWithCleanup(
func() error {
addCalls++
if addCalls == 1 {
return errors.New("listen tcp 10.0.0.1:32000: bind: address already in use")
}
return nil
},
func() error {
cleanupCalls++
return nil
},
0,
)
if err != nil {
t.Fatalf("expected retry to succeed, got %v", err)
}
if addCalls != 2 {
t.Fatalf("expected 2 add attempts, got %d", addCalls)
}
if cleanupCalls != 1 {
t.Fatalf("expected 1 cleanup attempt, got %d", cleanupCalls)
}
}
func TestRetryTunnelServiceAddWithCleanupSkipsCleanupOnNonBindError(t *testing.T) {
addCalls := 0
cleanupCalls := 0
err := retryTunnelServiceAddWithCleanup(
func() error {
addCalls++
return errors.New("network timeout")
},
func() error {
cleanupCalls++
return nil
},
0,
)
if err == nil {
t.Fatalf("expected hard error")
}
if addCalls != 1 {
t.Fatalf("expected 1 add attempt, got %d", addCalls)
}
if cleanupCalls != 0 {
t.Fatalf("expected 0 cleanup attempts, got %d", cleanupCalls)
}
}
func TestRetryTunnelServiceAddWithCleanupReturnsCleanupError(t *testing.T) {
cleanupErr := errors.New("delete failed")
err := retryTunnelServiceAddWithCleanup(
func() error {
return errors.New("listen tcp 10.0.0.1:32000: bind: address already in use")
},
func() error {
return cleanupErr
},
0,
)
if !errors.Is(err, cleanupErr) {
t.Fatalf("expected cleanup error %v, got %v", cleanupErr, err)
}
}
func TestBuildForwardServiceConfigs_UsesBindIPForListen(t *testing.T) {
forward := &forwardRecord{RemoteAddr: "1.2.3.4:80", Strategy: "fifo", TunnelID: 7}
node := &nodeRecord{TCPListenAddr: "[::]", UDPListenAddr: "[::]"}
@@ -43,6 +43,19 @@ func TestBuildTunnelChainServiceConfig_UsesConnectIPForListen(t *testing.T) {
}
}
func TestBuildTunnelChainServiceConfig_FallsBackToNodeListenAddr(t *testing.T) {
node := &nodeRecord{TCPListenAddr: "10.8.0.5"}
chain := tunnelRuntimeNode{Protocol: "tls", Port: 21002}
services := buildTunnelChainServiceConfig(99, chain, node)
if len(services) != 1 {
t.Fatalf("expected 1 service, got %d", len(services))
}
addr, _ := services[0]["addr"].(string)
if addr != "10.8.0.5:21002" {
t.Fatalf("expected node listen addr 10.8.0.5:21002, got %q", addr)
}
}
func TestBuildTunnelChainServiceConfig_DefaultListenAddrWhenConnectIPEmpty(t *testing.T) {
node := &nodeRecord{TCPListenAddr: "[::]"}
chain := tunnelRuntimeNode{Protocol: "tls", Port: 21001}
+42 -2
View File
@@ -25,6 +25,8 @@ import (
"gorm.io/gorm"
)
const tunnelServiceBindRetryDelay = 150 * time.Millisecond
func (h *Handler) userCreate(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("请求失败"))
@@ -2584,7 +2586,7 @@ func (h *Handler) applyTunnelRuntime(state *tunnelCreateState) ([]int64, []int64
createdChains = append(createdChains, chainNode.NodeID)
serviceData := buildTunnelChainServiceConfig(state.TunnelID, chainNode, state.Nodes[chainNode.NodeID])
if _, err := h.sendNodeCommand(chainNode.NodeID, "AddService", serviceData, true, false); err != nil {
if err := h.addTunnelServiceOnNode(chainNode.NodeID, state.TunnelID, serviceData); err != nil {
return createdChains, createdServices, fmt.Errorf("转发链节点 %s 下发服务失败: %w", nodeDisplayName(state.Nodes[chainNode.NodeID]), err)
}
createdServices = append(createdServices, chainNode.NodeID)
@@ -2596,7 +2598,7 @@ func (h *Handler) applyTunnelRuntime(state *tunnelCreateState) ([]int64, []int64
continue
}
serviceData := buildTunnelChainServiceConfig(state.TunnelID, outNode, state.Nodes[outNode.NodeID])
if _, err := h.sendNodeCommand(outNode.NodeID, "AddService", serviceData, true, false); err != nil {
if err := h.addTunnelServiceOnNode(outNode.NodeID, state.TunnelID, serviceData); err != nil {
return createdChains, createdServices, fmt.Errorf("出口节点 %s 下发服务失败: %w", nodeDisplayName(state.Nodes[outNode.NodeID]), err)
}
createdServices = append(createdServices, outNode.NodeID)
@@ -2605,6 +2607,44 @@ func (h *Handler) applyTunnelRuntime(state *tunnelCreateState) ([]int64, []int64
return createdChains, createdServices, nil
}
func retryTunnelServiceAddWithCleanup(add func() error, cleanup func() error, wait time.Duration) error {
if add == nil {
return errors.New("invalid tunnel service add callback")
}
err := add()
if err == nil || !isAddressAlreadyInUseError(err) {
return err
}
if cleanup == nil {
return err
}
if cleanupErr := cleanup(); cleanupErr != nil {
return cleanupErr
}
if wait > 0 {
time.Sleep(wait)
}
return add()
}
func (h *Handler) addTunnelServiceOnNode(nodeID, tunnelID int64, serviceData []map[string]interface{}) error {
if h == nil {
return errors.New("invalid tunnel service context")
}
serviceName := fmt.Sprintf("%d_tls", tunnelID)
return retryTunnelServiceAddWithCleanup(
func() error {
_, err := h.sendNodeCommand(nodeID, "AddService", serviceData, true, false)
return err
},
func() error {
_, err := h.sendNodeCommand(nodeID, "DeleteService", map[string]interface{}{"services": []string{serviceName}}, false, true)
return err
},
tunnelServiceBindRetryDelay,
)
}
func (h *Handler) rollbackTunnelRuntime(chainNodeIDs, serviceNodeIDs []int64, tunnelID int64) {
if h == nil || tunnelID <= 0 {
return