mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-11 11:46:37 +08:00
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
This commit is contained in:
@@ -71,10 +71,11 @@ type RuntimeReleaseRoleRequest struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type RuntimeDiagnoseRequest struct {
|
type RuntimeDiagnoseRequest struct {
|
||||||
IP string `json:"ip"`
|
IP string `json:"ip"`
|
||||||
Port int `json:"port"`
|
Port int `json:"port"`
|
||||||
Count int `json:"count"`
|
Count int `json:"count"`
|
||||||
Timeout int `json:"timeout"`
|
Timeout int `json:"timeout"`
|
||||||
|
Protocol string `json:"protocol"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type RuntimeNodeCommandRequest struct {
|
type RuntimeNodeCommandRequest struct {
|
||||||
|
|||||||
@@ -1402,10 +1402,37 @@ func (h *Handler) tcpPingViaRemoteNode(node *nodeRecord, ip string, port int, op
|
|||||||
|
|
||||||
fc := client.NewFederationClientWithTimeout(options.commandTimeout)
|
fc := client.NewFederationClientWithTimeout(options.commandTimeout)
|
||||||
return fc.Diagnose(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeDiagnoseRequest{
|
return fc.Diagnose(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeDiagnoseRequest{
|
||||||
IP: strings.TrimSpace(ip),
|
IP: strings.TrimSpace(ip),
|
||||||
Port: port,
|
Port: port,
|
||||||
Count: 4,
|
Count: 4,
|
||||||
Timeout: options.pingTimeoutMS,
|
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) {
|
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)
|
return h.tcpPingViaRemoteNode(node, ip, port, options)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -86,10 +86,11 @@ type federationRuntimeReleaseRoleRequest struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type federationRuntimeDiagnoseRequest struct {
|
type federationRuntimeDiagnoseRequest struct {
|
||||||
IP string `json:"ip"`
|
IP string `json:"ip"`
|
||||||
Port int `json:"port"`
|
Port int `json:"port"`
|
||||||
Count int `json:"count"`
|
Count int `json:"count"`
|
||||||
Timeout int `json:"timeout"`
|
Timeout int `json:"timeout"`
|
||||||
|
Protocol string `json:"protocol"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type federationRuntimeCommandRequest struct {
|
type federationRuntimeCommandRequest struct {
|
||||||
@@ -1281,7 +1282,12 @@ func (h *Handler) federationRuntimeDiagnose(w http.ResponseWriter, r *http.Reque
|
|||||||
commandTimeout = diagnosisCommandTimeout
|
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,
|
"ip": req.IP,
|
||||||
"port": req.Port,
|
"port": req.Port,
|
||||||
"count": req.Count,
|
"count": req.Count,
|
||||||
|
|||||||
@@ -3635,13 +3635,14 @@ func buildTunnelDialerConfig(protocol string) map[string]interface{} {
|
|||||||
dialer["metadata"] = map[string]interface{}{
|
dialer["metadata"] = map[string]interface{}{
|
||||||
"kcp.keepalive": 10,
|
"kcp.keepalive": 10,
|
||||||
"kcp.tcp": false,
|
"kcp.tcp": false,
|
||||||
"kcp.mode": "fast3",
|
"kcp.mode": "fast2",
|
||||||
"kcp.sndwnd": 2048,
|
"kcp.sndwnd": 2048,
|
||||||
"kcp.rcvwnd": 2048,
|
"kcp.rcvwnd": 2048,
|
||||||
"kcp.mtu": 1400,
|
"kcp.mtu": 1400,
|
||||||
"kcp.datashard": 0,
|
"kcp.datashard": 10,
|
||||||
"kcp.parityshard": 0,
|
"kcp.parityshard": 3,
|
||||||
"kcp.nocomp": true,
|
"kcp.nocomp": true,
|
||||||
|
"kcp.nc": 0,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return dialer
|
return dialer
|
||||||
@@ -3655,13 +3656,14 @@ func buildTunnelListenerConfig(protocol string) map[string]interface{} {
|
|||||||
listener["metadata"] = map[string]interface{}{
|
listener["metadata"] = map[string]interface{}{
|
||||||
"kcp.keepalive": 10,
|
"kcp.keepalive": 10,
|
||||||
"kcp.tcp": false,
|
"kcp.tcp": false,
|
||||||
"kcp.mode": "fast3",
|
"kcp.mode": "fast2",
|
||||||
"kcp.sndwnd": 2048,
|
"kcp.sndwnd": 2048,
|
||||||
"kcp.rcvwnd": 2048,
|
"kcp.rcvwnd": 2048,
|
||||||
"kcp.mtu": 1400,
|
"kcp.mtu": 1400,
|
||||||
"kcp.datashard": 0,
|
"kcp.datashard": 10,
|
||||||
"kcp.parityshard": 0,
|
"kcp.parityshard": 3,
|
||||||
"kcp.nocomp": true,
|
"kcp.nocomp": true,
|
||||||
|
"kcp.nc": 0,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return listener
|
return listener
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"github.com/go-gost/core/logger"
|
"github.com/go-gost/core/logger"
|
||||||
md "github.com/go-gost/core/metadata"
|
md "github.com/go-gost/core/metadata"
|
||||||
kcp_util "github.com/go-gost/x/internal/util/kcp"
|
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/go-gost/x/registry"
|
||||||
"github.com/xtaci/kcp-go/v5"
|
"github.com/xtaci/kcp-go/v5"
|
||||||
"github.com/xtaci/smux"
|
"github.com/xtaci/smux"
|
||||||
@@ -48,6 +49,9 @@ func (d *kcpDialer) Init(md md.Metadata) (err error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
d.md.config.Init()
|
d.md.config.Init()
|
||||||
|
if md != nil && md.IsExists("kcp.nc") {
|
||||||
|
d.md.config.NoCongestion = mdutil.GetInt(md, "kcp.nc")
|
||||||
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ import (
|
|||||||
limiter_wrapper "github.com/go-gost/x/limiter/traffic/wrapper"
|
limiter_wrapper "github.com/go-gost/x/limiter/traffic/wrapper"
|
||||||
metrics "github.com/go-gost/x/metrics/wrapper"
|
metrics "github.com/go-gost/x/metrics/wrapper"
|
||||||
stats "github.com/go-gost/x/observer/stats/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/go-gost/x/registry"
|
||||||
"github.com/xtaci/kcp-go/v5"
|
"github.com/xtaci/kcp-go/v5"
|
||||||
"github.com/xtaci/smux"
|
"github.com/xtaci/smux"
|
||||||
@@ -53,6 +54,9 @@ func (l *kcpListener) Init(md md.Metadata) (err error) {
|
|||||||
|
|
||||||
config := l.md.config
|
config := l.md.config
|
||||||
config.Init()
|
config.Init()
|
||||||
|
if md != nil && md.IsExists("kcp.nc") {
|
||||||
|
config.NoCongestion = mdutil.GetInt(md, "kcp.nc")
|
||||||
|
}
|
||||||
|
|
||||||
var conn net.PacketConn
|
var conn net.PacketConn
|
||||||
if config.TCP {
|
if config.TCP {
|
||||||
|
|||||||
Reference in New Issue
Block a user