From c1bc795674732778716642f9e5a0d42f1e9dc3e5 Mon Sep 17 00:00:00 2001 From: sagitchu Date: Tue, 21 Apr 2026 15:41:32 +0800 Subject: [PATCH] fix(kcp): enable congestion control, add FEC, and fix remote node UDP diagnosis - KCP: switch from fast3 to fast2 mode, enable FEC (10/3), enable congestion control (nc=0) for automatic rate adaptation - go-gost/x: support kcp.nc metadata to override mode-initialized NoCongestion value in dialer and listener Init() - Diagnosis: add udpPingViaRemoteNode, fix pingViaRemoteNode to dispatch based on protocol, add Protocol field to federation diagnose request structs --- go-backend/internal/http/client/federation.go | 9 +++-- .../internal/http/handler/control_plane.go | 38 +++++++++++++++++-- .../internal/http/handler/federation.go | 16 +++++--- go-backend/internal/http/handler/mutations.go | 14 ++++--- go-gost/x/dialer/kcp/dialer.go | 4 ++ go-gost/x/listener/kcp/listener.go | 4 ++ 6 files changed, 66 insertions(+), 19 deletions(-) diff --git a/go-backend/internal/http/client/federation.go b/go-backend/internal/http/client/federation.go index 09d6a64..a5a887a 100644 --- a/go-backend/internal/http/client/federation.go +++ b/go-backend/internal/http/client/federation.go @@ -71,10 +71,11 @@ type RuntimeReleaseRoleRequest struct { } type RuntimeDiagnoseRequest struct { - IP string `json:"ip"` - Port int `json:"port"` - Count int `json:"count"` - Timeout int `json:"timeout"` + IP string `json:"ip"` + Port int `json:"port"` + Count int `json:"count"` + Timeout int `json:"timeout"` + Protocol string `json:"protocol"` } type RuntimeNodeCommandRequest struct { diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index ea8e563..ba875c6 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -1402,10 +1402,37 @@ func (h *Handler) tcpPingViaRemoteNode(node *nodeRecord, ip string, port int, op fc := client.NewFederationClientWithTimeout(options.commandTimeout) return fc.Diagnose(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeDiagnoseRequest{ - IP: strings.TrimSpace(ip), - Port: port, - Count: 4, - Timeout: options.pingTimeoutMS, + IP: strings.TrimSpace(ip), + Port: port, + Count: 4, + Timeout: options.pingTimeoutMS, + Protocol: "tcp", + }) +} + +func (h *Handler) udpPingViaRemoteNode(node *nodeRecord, ip string, port int, options diagnosisExecOptions) (map[string]interface{}, error) { + if node == nil { + return nil, errors.New("节点不存在") + } + remoteURL := strings.TrimSpace(node.RemoteURL) + remoteToken := strings.TrimSpace(node.RemoteToken) + if remoteURL == "" || remoteToken == "" { + return nil, errors.New("远程节点缺少共享配置") + } + if options.commandTimeout <= 0 { + options.commandTimeout = diagnosisCommandTimeout + } + if options.pingTimeoutMS <= 0 { + options.pingTimeoutMS = int(diagnosisCommandTimeout / time.Millisecond) + } + + fc := client.NewFederationClientWithTimeout(options.commandTimeout) + return fc.Diagnose(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeDiagnoseRequest{ + IP: strings.TrimSpace(ip), + Port: port, + Count: 4, + Timeout: options.pingTimeoutMS, + Protocol: "udp", }) } @@ -1422,6 +1449,9 @@ func (h *Handler) pingViaNode(nodeID int64, ip string, port int, protocol string } func (h *Handler) pingViaRemoteNode(node *nodeRecord, ip string, port int, protocol string, options diagnosisExecOptions) (map[string]interface{}, error) { + if isUDPBasedProtocol(protocol) { + return h.udpPingViaRemoteNode(node, ip, port, options) + } return h.tcpPingViaRemoteNode(node, ip, port, options) } diff --git a/go-backend/internal/http/handler/federation.go b/go-backend/internal/http/handler/federation.go index 5562604..2e76e97 100644 --- a/go-backend/internal/http/handler/federation.go +++ b/go-backend/internal/http/handler/federation.go @@ -86,10 +86,11 @@ type federationRuntimeReleaseRoleRequest struct { } type federationRuntimeDiagnoseRequest struct { - IP string `json:"ip"` - Port int `json:"port"` - Count int `json:"count"` - Timeout int `json:"timeout"` + IP string `json:"ip"` + Port int `json:"port"` + Count int `json:"count"` + Timeout int `json:"timeout"` + Protocol string `json:"protocol"` } type federationRuntimeCommandRequest struct { @@ -1281,7 +1282,12 @@ func (h *Handler) federationRuntimeDiagnose(w http.ResponseWriter, r *http.Reque commandTimeout = diagnosisCommandTimeout } - res, err := h.sendNodeCommandWithTimeout(share.NodeID, "TcpPing", map[string]interface{}{ + commandType := "TcpPing" + if isUDPBasedProtocol(req.Protocol) { + commandType = "UdpPing" + } + + res, err := h.sendNodeCommandWithTimeout(share.NodeID, commandType, map[string]interface{}{ "ip": req.IP, "port": req.Port, "count": req.Count, diff --git a/go-backend/internal/http/handler/mutations.go b/go-backend/internal/http/handler/mutations.go index fc98f98..0f56988 100644 --- a/go-backend/internal/http/handler/mutations.go +++ b/go-backend/internal/http/handler/mutations.go @@ -3635,13 +3635,14 @@ func buildTunnelDialerConfig(protocol string) map[string]interface{} { dialer["metadata"] = map[string]interface{}{ "kcp.keepalive": 10, "kcp.tcp": false, - "kcp.mode": "fast3", + "kcp.mode": "fast2", "kcp.sndwnd": 2048, "kcp.rcvwnd": 2048, "kcp.mtu": 1400, - "kcp.datashard": 0, - "kcp.parityshard": 0, + "kcp.datashard": 10, + "kcp.parityshard": 3, "kcp.nocomp": true, + "kcp.nc": 0, } } return dialer @@ -3655,13 +3656,14 @@ func buildTunnelListenerConfig(protocol string) map[string]interface{} { listener["metadata"] = map[string]interface{}{ "kcp.keepalive": 10, "kcp.tcp": false, - "kcp.mode": "fast3", + "kcp.mode": "fast2", "kcp.sndwnd": 2048, "kcp.rcvwnd": 2048, "kcp.mtu": 1400, - "kcp.datashard": 0, - "kcp.parityshard": 0, + "kcp.datashard": 10, + "kcp.parityshard": 3, "kcp.nocomp": true, + "kcp.nc": 0, } } return listener diff --git a/go-gost/x/dialer/kcp/dialer.go b/go-gost/x/dialer/kcp/dialer.go index a73df73..3651bea 100644 --- a/go-gost/x/dialer/kcp/dialer.go +++ b/go-gost/x/dialer/kcp/dialer.go @@ -11,6 +11,7 @@ import ( "github.com/go-gost/core/logger" md "github.com/go-gost/core/metadata" kcp_util "github.com/go-gost/x/internal/util/kcp" + mdutil "github.com/go-gost/x/metadata/util" "github.com/go-gost/x/registry" "github.com/xtaci/kcp-go/v5" "github.com/xtaci/smux" @@ -48,6 +49,9 @@ func (d *kcpDialer) Init(md md.Metadata) (err error) { } d.md.config.Init() + if md != nil && md.IsExists("kcp.nc") { + d.md.config.NoCongestion = mdutil.GetInt(md, "kcp.nc") + } return nil } diff --git a/go-gost/x/listener/kcp/listener.go b/go-gost/x/listener/kcp/listener.go index 71dd053..d38175d 100644 --- a/go-gost/x/listener/kcp/listener.go +++ b/go-gost/x/listener/kcp/listener.go @@ -15,6 +15,7 @@ import ( limiter_wrapper "github.com/go-gost/x/limiter/traffic/wrapper" metrics "github.com/go-gost/x/metrics/wrapper" stats "github.com/go-gost/x/observer/stats/wrapper" + mdutil "github.com/go-gost/x/metadata/util" "github.com/go-gost/x/registry" "github.com/xtaci/kcp-go/v5" "github.com/xtaci/smux" @@ -53,6 +54,9 @@ func (l *kcpListener) Init(md md.Metadata) (err error) { config := l.md.config config.Init() + if md != nil && md.IsExists("kcp.nc") { + config.NoCongestion = mdutil.GetInt(md, "kcp.nc") + } var conn net.PacketConn if config.TCP {