mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-29 16:06:36 +08:00
Merge pull request #103 from Sagit-chu/opencode/tidy-panda
fix: tls udp 转发
This commit is contained in:
@@ -1000,14 +1000,19 @@ func (h *Handler) federationRuntimeApplyRole(w http.ResponseWriter, r *http.Requ
|
|||||||
response.WriteJSON(w, response.ErrDefault("Invalid target"))
|
response.WriteJSON(w, response.ErrDefault("Invalid target"))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
targetProtocol := defaultString(target.Protocol, protocol)
|
||||||
|
connector := map[string]interface{}{
|
||||||
|
"type": "relay",
|
||||||
|
}
|
||||||
|
if isTLSTunnelProtocol(targetProtocol) {
|
||||||
|
connector["metadata"] = map[string]interface{}{"nodelay": true}
|
||||||
|
}
|
||||||
nodeItems = append(nodeItems, map[string]interface{}{
|
nodeItems = append(nodeItems, map[string]interface{}{
|
||||||
"name": fmt.Sprintf("node_%d", i+1),
|
"name": fmt.Sprintf("node_%d", i+1),
|
||||||
"addr": processServerAddress(fmt.Sprintf("%s:%d", host, target.Port)),
|
"addr": processServerAddress(fmt.Sprintf("%s:%d", host, target.Port)),
|
||||||
"connector": map[string]interface{}{
|
"connector": connector,
|
||||||
"type": "relay",
|
|
||||||
},
|
|
||||||
"dialer": map[string]interface{}{
|
"dialer": map[string]interface{}{
|
||||||
"type": defaultString(target.Protocol, protocol),
|
"type": targetProtocol,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -1046,6 +1051,9 @@ func (h *Handler) federationRuntimeApplyRole(w http.ResponseWriter, r *http.Requ
|
|||||||
"type": protocol,
|
"type": protocol,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
if isTLSTunnelProtocol(protocol) {
|
||||||
|
service["handler"].(map[string]interface{})["metadata"] = map[string]interface{}{"nodelay": true}
|
||||||
|
}
|
||||||
if req.Role == "middle" {
|
if req.Role == "middle" {
|
||||||
service["handler"].(map[string]interface{})["chain"] = chainName
|
service["handler"].(map[string]interface{})["chain"] = chainName
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2564,14 +2564,19 @@ func buildTunnelChainConfig(tunnelID int64, fromNodeID int64, targets []tunnelRu
|
|||||||
if port <= 0 {
|
if port <= 0 {
|
||||||
return nil, errors.New("节点端口不能为空")
|
return nil, errors.New("节点端口不能为空")
|
||||||
}
|
}
|
||||||
|
protocol := defaultString(target.Protocol, "tls")
|
||||||
|
connector := map[string]interface{}{
|
||||||
|
"type": "relay",
|
||||||
|
}
|
||||||
|
if isTLSTunnelProtocol(protocol) {
|
||||||
|
connector["metadata"] = map[string]interface{}{"nodelay": true}
|
||||||
|
}
|
||||||
nodeItems = append(nodeItems, map[string]interface{}{
|
nodeItems = append(nodeItems, map[string]interface{}{
|
||||||
"name": fmt.Sprintf("node_%d", idx+1),
|
"name": fmt.Sprintf("node_%d", idx+1),
|
||||||
"addr": processServerAddress(fmt.Sprintf("%s:%d", host, port)),
|
"addr": processServerAddress(fmt.Sprintf("%s:%d", host, port)),
|
||||||
"connector": map[string]interface{}{
|
"connector": connector,
|
||||||
"type": "relay",
|
|
||||||
},
|
|
||||||
"dialer": map[string]interface{}{
|
"dialer": map[string]interface{}{
|
||||||
"type": defaultString(target.Protocol, "tls"),
|
"type": protocol,
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -2600,14 +2605,19 @@ func buildTunnelChainServiceConfig(tunnelID int64, chainNode tunnelRuntimeNode,
|
|||||||
if node == nil {
|
if node == nil {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
protocol := defaultString(chainNode.Protocol, "tls")
|
||||||
|
handlerCfg := map[string]interface{}{
|
||||||
|
"type": "relay",
|
||||||
|
}
|
||||||
|
if isTLSTunnelProtocol(protocol) {
|
||||||
|
handlerCfg["metadata"] = map[string]interface{}{"nodelay": true}
|
||||||
|
}
|
||||||
service := map[string]interface{}{
|
service := map[string]interface{}{
|
||||||
"name": fmt.Sprintf("%d_tls", tunnelID),
|
"name": fmt.Sprintf("%d_tls", tunnelID),
|
||||||
"addr": fmt.Sprintf("%s:%d", node.TCPListenAddr, chainNode.Port),
|
"addr": fmt.Sprintf("%s:%d", node.TCPListenAddr, chainNode.Port),
|
||||||
"handler": map[string]interface{}{
|
"handler": handlerCfg,
|
||||||
"type": "relay",
|
|
||||||
},
|
|
||||||
"listener": map[string]interface{}{
|
"listener": map[string]interface{}{
|
||||||
"type": defaultString(chainNode.Protocol, "tls"),
|
"type": protocol,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
if chainNode.ChainType == 2 {
|
if chainNode.ChainType == 2 {
|
||||||
@@ -2653,6 +2663,10 @@ func nodeDisplayName(node *nodeRecord) string {
|
|||||||
return fmt.Sprintf("node_%d", node.ID)
|
return fmt.Sprintf("node_%d", node.ID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func isTLSTunnelProtocol(protocol string) bool {
|
||||||
|
return strings.EqualFold(strings.TrimSpace(defaultString(protocol, "tls")), "tls")
|
||||||
|
}
|
||||||
|
|
||||||
func nodeSupportsV4(node *nodeRecord) bool {
|
func nodeSupportsV4(node *nodeRecord) bool {
|
||||||
if node == nil {
|
if node == nil {
|
||||||
return false
|
return false
|
||||||
|
|||||||
@@ -54,10 +54,10 @@ var needWrap = false
|
|||||||
|
|
||||||
// SetProtocolBlock sets protocol blocking switches and recomputes wrapper need
|
// SetProtocolBlock sets protocol blocking switches and recomputes wrapper need
|
||||||
func SetProtocolBlock(httpOn int, tlsOn int, socksOn int) {
|
func SetProtocolBlock(httpOn int, tlsOn int, socksOn int) {
|
||||||
isHttp = httpOn
|
isHttp = httpOn
|
||||||
isTls = tlsOn
|
isTls = tlsOn
|
||||||
isSocks = socksOn
|
isSocks = socksOn
|
||||||
needWrap = isTls+isSocks+isHttp > 0
|
needWrap = isTls+isSocks+isHttp > 0
|
||||||
}
|
}
|
||||||
|
|
||||||
type Option func(opts *options)
|
type Option func(opts *options)
|
||||||
@@ -292,7 +292,9 @@ func (s *defaultService) Serve() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if err := s.handler.Handle(ctx, conn); err != nil {
|
if err := s.handler.Handle(ctx, conn); err != nil {
|
||||||
log.Error(err)
|
if !errors.Is(err, net.ErrClosed) {
|
||||||
|
log.Error(err)
|
||||||
|
}
|
||||||
if v := xmetrics.GetCounter(xmetrics.MetricServiceHandlerErrorsCounter,
|
if v := xmetrics.GetCounter(xmetrics.MetricServiceHandlerErrorsCounter,
|
||||||
metrics.Labels{"service": s.name, "client": clientIP}); v != nil {
|
metrics.Labels{"service": s.name, "client": clientIP}); v != nil {
|
||||||
v.Inc()
|
v.Inc()
|
||||||
@@ -403,12 +405,12 @@ func (s *defaultService) observeStats(ctx context.Context) {
|
|||||||
TotalErrs: st.Get(stats.KindTotalErrs),
|
TotalErrs: st.Get(stats.KindTotalErrs),
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
// 将流量累积到全局管理器,而不是立即上报
|
// 将流量累积到全局管理器,而不是立即上报
|
||||||
if outputBytes > 0 || inputBytes > 0 {
|
if outputBytes > 0 || inputBytes > 0 {
|
||||||
globalManager := GetGlobalTrafficManager()
|
globalManager := GetGlobalTrafficManager()
|
||||||
globalManager.AddTraffic(s.name, int64(outputBytes), int64(inputBytes))
|
globalManager.AddTraffic(s.name, int64(outputBytes), int64(inputBytes))
|
||||||
|
|
||||||
// 立即重置流量计数(因为已经记录到全局管理器中)
|
// 立即重置流量计数(因为已经记录到全局管理器中)
|
||||||
if xstats, ok := st.(*xstats.Stats); ok {
|
if xstats, ok := st.(*xstats.Stats); ok {
|
||||||
xstats.ResetTraffic(st.Get(stats.KindInputBytes)-inputBytes, st.Get(stats.KindOutputBytes)-outputBytes)
|
xstats.ResetTraffic(st.Get(stats.KindInputBytes)-inputBytes, st.Get(stats.KindOutputBytes)-outputBytes)
|
||||||
|
|||||||
Reference in New Issue
Block a user