forked from Mxmilu666/frp
Merge remote-tracking branch 'upstream/dev' into dev
# Conflicts: # .github/workflows/build-and-push-image.yml # cmd/frpc/sub/verify.go # go.mod # go.sum # pkg/util/version/version.go
This commit is contained in:
@@ -2,12 +2,16 @@ package client
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/fatedier/frp/client/configmgmt"
|
||||
"github.com/fatedier/frp/pkg/config/source"
|
||||
v1 "github.com/fatedier/frp/pkg/config/v1"
|
||||
"github.com/fatedier/frp/pkg/policy/security"
|
||||
"github.com/fatedier/frp/pkg/vnet"
|
||||
)
|
||||
|
||||
func newTestRawTCPProxyConfig(name string) *v1.TCPProxyConfig {
|
||||
@@ -22,6 +26,256 @@ func newTestRawTCPProxyConfig(name string) *v1.TCPProxyConfig {
|
||||
}
|
||||
}
|
||||
|
||||
func newTestVirtualNetProxyConfig(name string) *v1.STCPProxyConfig {
|
||||
return &v1.STCPProxyConfig{
|
||||
ProxyBaseConfig: v1.ProxyBaseConfig{
|
||||
Name: name,
|
||||
Type: "stcp",
|
||||
ProxyBackend: v1.ProxyBackend{
|
||||
Plugin: v1.TypedClientPluginOptions{
|
||||
Type: v1.PluginVirtualNet,
|
||||
ClientPluginOptions: &v1.VirtualNetPluginOptions{Type: v1.PluginVirtualNet},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func newTestVirtualNetVisitorConfig(name string) *v1.STCPVisitorConfig {
|
||||
return &v1.STCPVisitorConfig{
|
||||
VisitorBaseConfig: v1.VisitorBaseConfig{
|
||||
Name: name,
|
||||
Type: "stcp",
|
||||
ServerName: "vnet-server",
|
||||
SecretKey: "secret",
|
||||
BindPort: -1,
|
||||
Plugin: v1.TypedVisitorPluginOptions{
|
||||
Type: v1.VisitorPluginVirtualNet,
|
||||
VisitorPluginOptions: &v1.VirtualNetVisitorPluginOptions{
|
||||
Type: v1.VisitorPluginVirtualNet,
|
||||
DestinationIP: "100.86.0.1",
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func TestServiceConfigManagerReloadVirtualNetRuntimeDependency(t *testing.T) {
|
||||
const runtimeErr = "VirtualNet-dependent configuration requires a VirtualNet runtime enabled at startup"
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
startupVirtualNetAddr string
|
||||
nextConfig string
|
||||
wantRuntimeDependency bool
|
||||
}{
|
||||
{
|
||||
name: "unrelated common config",
|
||||
nextConfig: `serverAddr = "0.0.0.0"`,
|
||||
},
|
||||
{
|
||||
name: "VirtualNet address without startup runtime",
|
||||
nextConfig: `featureGates = { VirtualNet = true }
|
||||
virtualNet.address = "100.86.0.4/24"
|
||||
`,
|
||||
wantRuntimeDependency: true,
|
||||
},
|
||||
{
|
||||
name: "VirtualNet proxy without startup runtime",
|
||||
nextConfig: `[[proxies]]
|
||||
name = "vnet-proxy"
|
||||
type = "stcp"
|
||||
secretKey = "secret"
|
||||
[proxies.plugin]
|
||||
type = "virtual_net"
|
||||
`,
|
||||
wantRuntimeDependency: true,
|
||||
},
|
||||
{
|
||||
name: "VirtualNet visitor without startup runtime",
|
||||
nextConfig: `[[visitors]]
|
||||
name = "vnet-visitor"
|
||||
type = "stcp"
|
||||
serverName = "vnet-server"
|
||||
secretKey = "secret"
|
||||
bindPort = -1
|
||||
[visitors.plugin]
|
||||
type = "virtual_net"
|
||||
destinationIP = "100.86.0.1"
|
||||
`,
|
||||
wantRuntimeDependency: true,
|
||||
},
|
||||
{
|
||||
name: "existing VirtualNet startup runtime",
|
||||
startupVirtualNetAddr: "100.86.0.4/24",
|
||||
nextConfig: `featureGates = { VirtualNet = true }
|
||||
virtualNet.address = "100.86.0.5/24"
|
||||
`,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
current := &v1.ClientCommonConfig{}
|
||||
if tc.startupVirtualNetAddr != "" {
|
||||
current.FeatureGates = map[string]bool{"VirtualNet": true}
|
||||
current.VirtualNet.Address = tc.startupVirtualNetAddr
|
||||
}
|
||||
if err := current.Complete(); err != nil {
|
||||
t.Fatalf("complete current config: %v", err)
|
||||
}
|
||||
|
||||
configFile := filepath.Join(t.TempDir(), "frpc.toml")
|
||||
if err := os.WriteFile(configFile, []byte(tc.nextConfig), 0o600); err != nil {
|
||||
t.Fatalf("write config: %v", err)
|
||||
}
|
||||
|
||||
configSource := source.NewConfigSource()
|
||||
aggregator := source.NewAggregator(configSource)
|
||||
svr := &Service{
|
||||
common: current,
|
||||
reloadCommon: current,
|
||||
configFilePath: configFile,
|
||||
unsafeFeatures: security.NewUnsafeFeatures(nil),
|
||||
aggregator: aggregator,
|
||||
configSource: configSource,
|
||||
}
|
||||
if tc.startupVirtualNetAddr != "" {
|
||||
svr.vnetController = vnet.NewController(current.VirtualNet)
|
||||
}
|
||||
|
||||
err := (&serviceConfigManager{svr: svr}).ReloadFromFile(true)
|
||||
if tc.wantRuntimeDependency {
|
||||
if !errors.Is(err, configmgmt.ErrApplyConfig) || !strings.Contains(err.Error(), runtimeErr) {
|
||||
t.Fatalf("expected VirtualNet runtime dependency error, got: %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
t.Fatalf("reload config: %v", err)
|
||||
}
|
||||
if svr.common != current {
|
||||
t.Fatal("reload should not replace startup common config")
|
||||
}
|
||||
if tc.startupVirtualNetAddr == "" && svr.vnetController != nil {
|
||||
t.Fatal("reload should not enable startup-only VirtualNet runtime state")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestServiceConfigManagerReloadVirtualNetRuntimeDependencyUsesMergedSources(t *testing.T) {
|
||||
const runtimeErr = "VirtualNet-dependent configuration requires a VirtualNet runtime enabled at startup"
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
nextConfig string
|
||||
storeProxy v1.ProxyConfigurer
|
||||
storeVisitor v1.VisitorConfigurer
|
||||
wantRuntimeDependency bool
|
||||
wantProxyPlugin string
|
||||
}{
|
||||
{
|
||||
name: "Store VirtualNet proxy is rejected",
|
||||
nextConfig: `serverAddr = "0.0.0.0"`,
|
||||
storeProxy: newTestVirtualNetProxyConfig("store-vnet"),
|
||||
wantRuntimeDependency: true,
|
||||
},
|
||||
{
|
||||
name: "Store VirtualNet visitor is rejected",
|
||||
nextConfig: `serverAddr = "0.0.0.0"`,
|
||||
storeVisitor: newTestVirtualNetVisitorConfig("store-vnet"),
|
||||
wantRuntimeDependency: true,
|
||||
},
|
||||
{
|
||||
name: "Store VirtualNet proxy overrides file proxy",
|
||||
nextConfig: `[[proxies]]
|
||||
name = "shared"
|
||||
type = "tcp"
|
||||
localPort = 10080
|
||||
remotePort = 10081
|
||||
`,
|
||||
storeProxy: newTestVirtualNetProxyConfig("shared"),
|
||||
wantRuntimeDependency: true,
|
||||
},
|
||||
{
|
||||
name: "Store non-VirtualNet proxy overrides file VirtualNet proxy",
|
||||
nextConfig: `[[proxies]]
|
||||
name = "shared"
|
||||
type = "stcp"
|
||||
secretKey = "secret"
|
||||
[proxies.plugin]
|
||||
type = "virtual_net"
|
||||
`,
|
||||
storeProxy: newTestRawTCPProxyConfig("shared"),
|
||||
wantProxyPlugin: "",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
current := &v1.ClientCommonConfig{}
|
||||
if err := current.Complete(); err != nil {
|
||||
t.Fatalf("complete current config: %v", err)
|
||||
}
|
||||
|
||||
configFile := filepath.Join(t.TempDir(), "frpc.toml")
|
||||
if err := os.WriteFile(configFile, []byte(tc.nextConfig), 0o600); err != nil {
|
||||
t.Fatalf("write config: %v", err)
|
||||
}
|
||||
|
||||
storeSource, err := source.NewStoreSource(source.StoreSourceConfig{
|
||||
Path: filepath.Join(t.TempDir(), "store.json"),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("new store source: %v", err)
|
||||
}
|
||||
if tc.storeProxy != nil {
|
||||
if err := storeSource.AddProxy(tc.storeProxy); err != nil {
|
||||
t.Fatalf("add store proxy: %v", err)
|
||||
}
|
||||
}
|
||||
if tc.storeVisitor != nil {
|
||||
if err := storeSource.AddVisitor(tc.storeVisitor); err != nil {
|
||||
t.Fatalf("add store visitor: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
configSource := source.NewConfigSource()
|
||||
aggregator := source.NewAggregator(configSource)
|
||||
aggregator.SetStoreSource(storeSource)
|
||||
svr := &Service{
|
||||
common: current,
|
||||
reloadCommon: current,
|
||||
configFilePath: configFile,
|
||||
unsafeFeatures: security.NewUnsafeFeatures(nil),
|
||||
aggregator: aggregator,
|
||||
configSource: configSource,
|
||||
storeSource: storeSource,
|
||||
}
|
||||
|
||||
err = (&serviceConfigManager{svr: svr}).ReloadFromFile(true)
|
||||
if tc.wantRuntimeDependency {
|
||||
if !errors.Is(err, configmgmt.ErrApplyConfig) || !strings.Contains(err.Error(), runtimeErr) {
|
||||
t.Fatalf("expected VirtualNet runtime dependency error, got: %v", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("reload config: %v", err)
|
||||
}
|
||||
|
||||
if len(svr.proxyCfgs) != 1 {
|
||||
t.Fatalf("expected one applied proxy, got %d", len(svr.proxyCfgs))
|
||||
}
|
||||
if got := svr.proxyCfgs[0].GetBaseConfig().Plugin.Type; got != tc.wantProxyPlugin {
|
||||
t.Fatalf("unexpected applied proxy plugin: %q", got)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestServiceConfigManagerCreateStoreProxyConflict(t *testing.T) {
|
||||
storeSource, err := source.NewStoreSource(source.StoreSourceConfig{
|
||||
Path: filepath.Join(t.TempDir(), "store.json"),
|
||||
|
||||
+11
-2
@@ -49,6 +49,8 @@ type SessionContext struct {
|
||||
Connector MessageConnector
|
||||
// Virtual net controller
|
||||
VnetController *vnet.Controller
|
||||
// UDPPacketCodec is immutable for the lifetime of this negotiated session.
|
||||
UDPPacketCodec string
|
||||
}
|
||||
|
||||
type Control struct {
|
||||
@@ -94,9 +96,16 @@ func NewControl(ctx context.Context, sessionCtx *SessionContext) (*Control, erro
|
||||
ctl.registerMsgHandlers()
|
||||
ctl.msgTransporter = transport.NewMessageTransporter(ctl.msgDispatcher)
|
||||
|
||||
ctl.pm = proxy.NewManager(ctl.ctx, sessionCtx.Common, sessionCtx.Auth.EncryptionKey(), ctl.msgTransporter, sessionCtx.VnetController)
|
||||
ctl.pm = proxy.NewManager(
|
||||
ctl.ctx,
|
||||
sessionCtx.Common,
|
||||
sessionCtx.Auth.EncryptionKey(),
|
||||
ctl.msgTransporter,
|
||||
sessionCtx.VnetController,
|
||||
sessionCtx.UDPPacketCodec,
|
||||
)
|
||||
ctl.vm = visitor.NewManager(ctl.ctx, sessionCtx.RunID, sessionCtx.Common,
|
||||
ctl.connectServer, ctl.msgTransporter, sessionCtx.VnetController)
|
||||
ctl.connectServer, ctl.msgTransporter, sessionCtx.VnetController, sessionCtx.UDPPacketCodec)
|
||||
return ctl, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -99,6 +99,7 @@ func (d *controlSessionDialer) Dial(previousRunID string) (*SessionContext, erro
|
||||
Auth: d.auth,
|
||||
Connector: newMessageConnector(connector, d.common.Transport.WireProtocol),
|
||||
VnetController: d.vnetController,
|
||||
UDPPacketCodec: loginResult.udpPacketCodec,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -127,8 +128,9 @@ func (d *controlSessionDialer) buildLoginMsg(previousRunID string) (*msg.Login,
|
||||
}
|
||||
|
||||
type loginExchangeResult struct {
|
||||
resp *msg.LoginResp
|
||||
crypto *wire.CryptoContext
|
||||
resp *msg.LoginResp
|
||||
crypto *wire.CryptoContext
|
||||
udpPacketCodec string
|
||||
}
|
||||
|
||||
func (d *controlSessionDialer) exchangeLogin(conn net.Conn, loginMsg *msg.Login) (*loginExchangeResult, error) {
|
||||
@@ -172,6 +174,7 @@ func (d *controlSessionDialer) exchangeLogin(conn net.Conn, loginMsg *msg.Login)
|
||||
}()
|
||||
|
||||
var cryptoContext *wire.CryptoContext
|
||||
var udpPacketCodec string
|
||||
if wireConn != nil {
|
||||
serverHelloFrame, err := wireConn.ReadFrame()
|
||||
if err != nil {
|
||||
@@ -191,6 +194,7 @@ func (d *controlSessionDialer) exchangeLogin(conn net.Conn, loginMsg *msg.Login)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
udpPacketCodec = serverHello.Selected.Message.UDPPacketCodec
|
||||
}
|
||||
|
||||
var loginRespMsg msg.LoginResp
|
||||
@@ -198,8 +202,9 @@ func (d *controlSessionDialer) exchangeLogin(conn net.Conn, loginMsg *msg.Login)
|
||||
return nil, err
|
||||
}
|
||||
return &loginExchangeResult{
|
||||
resp: &loginRespMsg,
|
||||
crypto: cryptoContext,
|
||||
resp: &loginRespMsg,
|
||||
crypto: cryptoContext,
|
||||
udpPacketCodec: udpPacketCodec,
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -117,6 +117,7 @@ func TestControlSessionDialerDialV1(t *testing.T) {
|
||||
defer sessionCtx.Connector.Close()
|
||||
|
||||
require.Equal(t, "run-v1", sessionCtx.RunID)
|
||||
require.Empty(t, sessionCtx.UDPPacketCodec)
|
||||
require.NotNil(t, sessionCtx.Conn)
|
||||
require.NotNil(t, sessionCtx.Connector)
|
||||
require.False(t, connector.closed.Load())
|
||||
@@ -225,6 +226,7 @@ func TestControlSessionDialerDialV2(t *testing.T) {
|
||||
defer sessionCtx.Connector.Close()
|
||||
|
||||
require.Equal(t, "run-v2", sessionCtx.RunID)
|
||||
require.Equal(t, wire.UDPPacketCodecBinary, sessionCtx.UDPPacketCodec)
|
||||
require.NotNil(t, sessionCtx.Conn)
|
||||
require.NotNil(t, sessionCtx.Connector)
|
||||
require.False(t, connector.closed.Load())
|
||||
|
||||
@@ -0,0 +1,125 @@
|
||||
// Copyright 2026 The frp Authors
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//go:build !frps
|
||||
|
||||
package client
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"net"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
clientproxy "github.com/fatedier/frp/client/proxy"
|
||||
"github.com/fatedier/frp/pkg/auth"
|
||||
v1 "github.com/fatedier/frp/pkg/config/v1"
|
||||
"github.com/fatedier/frp/pkg/msg"
|
||||
"github.com/fatedier/frp/pkg/proto/wire"
|
||||
)
|
||||
|
||||
func TestControlPropagatesBinaryUDPPacketCodecToWorkConn(t *testing.T) {
|
||||
echoConn, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.ParseIP("127.0.0.1")})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = echoConn.Close() })
|
||||
|
||||
echoDone := make(chan error, 1)
|
||||
go func() {
|
||||
buf := make([]byte, 64)
|
||||
n, addr, err := echoConn.ReadFromUDP(buf)
|
||||
if err == nil {
|
||||
_, err = echoConn.WriteToUDP(buf[:n], addr)
|
||||
}
|
||||
echoDone <- err
|
||||
}()
|
||||
|
||||
authRuntime, err := auth.BuildClientAuth(&v1.AuthClientConfig{
|
||||
Method: v1.AuthMethodToken,
|
||||
Token: "token",
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
controlConn, controlPeer := net.Pipe()
|
||||
t.Cleanup(func() {
|
||||
_ = controlConn.Close()
|
||||
_ = controlPeer.Close()
|
||||
})
|
||||
common := &v1.ClientCommonConfig{
|
||||
Transport: v1.ClientTransportConfig{WireProtocol: wire.ProtocolV2},
|
||||
UDPPacketSize: 1500,
|
||||
}
|
||||
ctl, err := NewControl(context.Background(), &SessionContext{
|
||||
Common: common,
|
||||
RunID: "binary-udp-test",
|
||||
Conn: msg.NewConn(controlConn, msg.NewV2ReadWriter(controlConn)),
|
||||
Auth: authRuntime,
|
||||
UDPPacketCodec: wire.UDPPacketCodecBinary,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(ctl.pm.Close)
|
||||
|
||||
echoAddr := echoConn.LocalAddr().(*net.UDPAddr)
|
||||
proxyCfg := &v1.UDPProxyConfig{
|
||||
ProxyBaseConfig: v1.ProxyBaseConfig{
|
||||
Name: "udp",
|
||||
Type: string(v1.ProxyTypeUDP),
|
||||
ProxyBackend: v1.ProxyBackend{
|
||||
LocalIP: "127.0.0.1",
|
||||
LocalPort: echoAddr.Port,
|
||||
},
|
||||
},
|
||||
}
|
||||
ctl.pm.UpdateAll([]v1.ProxyConfigurer{proxyCfg})
|
||||
require.Eventually(t, func() bool {
|
||||
status, ok := ctl.pm.GetProxyStatus("udp")
|
||||
return ok && status.Phase == clientproxy.ProxyPhaseWaitStart
|
||||
}, time.Second, 10*time.Millisecond)
|
||||
require.NoError(t, ctl.pm.StartProxy("udp", "", ""))
|
||||
|
||||
workClient, workServer := net.Pipe()
|
||||
t.Cleanup(func() {
|
||||
_ = workClient.Close()
|
||||
_ = workServer.Close()
|
||||
})
|
||||
deadline := time.Now().Add(3 * time.Second)
|
||||
require.NoError(t, workClient.SetDeadline(deadline))
|
||||
require.NoError(t, workServer.SetDeadline(deadline))
|
||||
ctl.pm.HandleWorkConn("udp", workClient, &msg.StartWorkConn{ProxyName: "udp"})
|
||||
|
||||
serverRW, err := msg.NewUDPPacketReadWriter(workServer, wire.ProtocolV2, wire.UDPPacketCodecBinary)
|
||||
require.NoError(t, err)
|
||||
writeDone := make(chan error, 1)
|
||||
in := &msg.UDPPacket{
|
||||
Content: []byte("binary udp"),
|
||||
RemoteAddr: &net.UDPAddr{IP: net.ParseIP("192.0.2.1"), Port: 12345},
|
||||
}
|
||||
go func() {
|
||||
writeDone <- serverRW.WriteMsg(in)
|
||||
}()
|
||||
|
||||
frame, err := wire.NewConn(workServer).ReadFrame()
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, wire.FrameTypeMessage, frame.Type)
|
||||
require.GreaterOrEqual(t, len(frame.Payload), 2)
|
||||
require.Equal(t, msg.V2TypeUDPPacketBinary, binary.BigEndian.Uint16(frame.Payload[:2]))
|
||||
out, err := msg.DecodeUDPPacketBinary(frame.Payload[2:])
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, in.Content, out.Content)
|
||||
require.Equal(t, in.RemoteAddr.String(), out.RemoteAddr.String())
|
||||
require.NoError(t, <-writeDone)
|
||||
require.NoError(t, <-echoDone)
|
||||
}
|
||||
+27
-9
@@ -63,11 +63,12 @@ func NewProxy(
|
||||
encryptionKey []byte,
|
||||
msgTransporter transport.MessageTransporter,
|
||||
vnetController *vnet.Controller,
|
||||
udpPacketCodec string,
|
||||
) (pxy Proxy) {
|
||||
var limiter *rate.Limiter
|
||||
limitBytes := pxyConf.GetBaseConfig().Transport.BandwidthLimit.Bytes()
|
||||
if limitBytes > 0 && pxyConf.GetBaseConfig().Transport.BandwidthLimitMode == types.BandwidthLimitModeClient {
|
||||
limiter = rate.NewLimiter(rate.Limit(float64(limitBytes)), int(limitBytes))
|
||||
limiter = limit.NewBandwidthLimiter(limitBytes)
|
||||
}
|
||||
|
||||
baseProxy := BaseProxy{
|
||||
@@ -80,6 +81,7 @@ func NewProxy(
|
||||
vnetController: vnetController,
|
||||
xl: xlog.FromContextSafe(ctx),
|
||||
ctx: ctx,
|
||||
udpPacketCodec: udpPacketCodec,
|
||||
}
|
||||
|
||||
factory := proxyFactoryRegistry[reflect.TypeOf(pxyConf)]
|
||||
@@ -102,9 +104,10 @@ type BaseProxy struct {
|
||||
proxyPlugin plugin.Plugin
|
||||
inWorkConnCallback func(*v1.ProxyBaseConfig, net.Conn, *msg.StartWorkConn) /* continue */ bool
|
||||
|
||||
mu sync.RWMutex
|
||||
xl *xlog.Logger
|
||||
ctx context.Context
|
||||
mu sync.RWMutex
|
||||
xl *xlog.Logger
|
||||
ctx context.Context
|
||||
udpPacketCodec string
|
||||
}
|
||||
|
||||
func (pxy *BaseProxy) Run() error {
|
||||
@@ -209,6 +212,26 @@ func (pxy *BaseProxy) HandleTCPWorkConnection(workConn net.Conn, m *msg.StartWor
|
||||
xl.Tracef("handle tcp work connection, useEncryption: %t, useCompression: %t",
|
||||
baseCfg.Transport.UseEncryption, baseCfg.Transport.UseCompression)
|
||||
|
||||
var srcAddr, dstAddr *net.TCPAddr
|
||||
if m.SrcAddr != "" && m.SrcPort != 0 {
|
||||
if m.DstAddr == "" {
|
||||
m.DstAddr = "127.0.0.1"
|
||||
}
|
||||
var err error
|
||||
srcAddr, err = net.ResolveTCPAddr("tcp", net.JoinHostPort(m.SrcAddr, strconv.Itoa(int(m.SrcPort))))
|
||||
if err != nil {
|
||||
xl.Warnf("resolve source address [%s] error: %v", m.SrcAddr, err)
|
||||
_ = workConn.Close()
|
||||
return
|
||||
}
|
||||
dstAddr, err = net.ResolveTCPAddr("tcp", net.JoinHostPort(m.DstAddr, strconv.Itoa(int(m.DstPort))))
|
||||
if err != nil {
|
||||
xl.Warnf("resolve destination address [%s] error: %v", m.DstAddr, err)
|
||||
_ = workConn.Close()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
remote, recycleFn, err := pxy.wrapWorkConn(workConn, encKey)
|
||||
if err != nil {
|
||||
xl.Errorf("wrap work connection: %v", err)
|
||||
@@ -218,11 +241,6 @@ func (pxy *BaseProxy) HandleTCPWorkConnection(workConn net.Conn, m *msg.StartWor
|
||||
// check if we need to send proxy protocol info
|
||||
var connInfo plugin.ConnectionInfo
|
||||
if m.SrcAddr != "" && m.SrcPort != 0 {
|
||||
if m.DstAddr == "" {
|
||||
m.DstAddr = "127.0.0.1"
|
||||
}
|
||||
srcAddr, _ := net.ResolveTCPAddr("tcp", net.JoinHostPort(m.SrcAddr, strconv.Itoa(int(m.SrcPort))))
|
||||
dstAddr, _ := net.ResolveTCPAddr("tcp", net.JoinHostPort(m.DstAddr, strconv.Itoa(int(m.DstPort))))
|
||||
connInfo.SrcAddr = srcAddr
|
||||
connInfo.DstAddr = dstAddr
|
||||
}
|
||||
|
||||
@@ -43,7 +43,8 @@ type Manager struct {
|
||||
encryptionKey []byte
|
||||
clientCfg *v1.ClientCommonConfig
|
||||
|
||||
ctx context.Context
|
||||
ctx context.Context
|
||||
udpPacketCodec string
|
||||
}
|
||||
|
||||
func NewManager(
|
||||
@@ -52,6 +53,7 @@ func NewManager(
|
||||
encryptionKey []byte,
|
||||
msgTransporter transport.MessageTransporter,
|
||||
vnetController *vnet.Controller,
|
||||
udpPacketCodec string,
|
||||
) *Manager {
|
||||
return &Manager{
|
||||
proxies: make(map[string]*Wrapper),
|
||||
@@ -61,6 +63,7 @@ func NewManager(
|
||||
encryptionKey: encryptionKey,
|
||||
clientCfg: clientCfg,
|
||||
ctx: ctx,
|
||||
udpPacketCodec: udpPacketCodec,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -166,7 +169,7 @@ func (pm *Manager) UpdateAll(proxyCfgs []v1.ProxyConfigurer) {
|
||||
for _, cfg := range proxyCfgs {
|
||||
name := cfg.GetBaseConfig().Name
|
||||
if _, ok := pm.proxies[name]; !ok {
|
||||
pxy := NewWrapper(pm.ctx, cfg, pm.clientCfg, pm.encryptionKey, pm.HandleEvent, pm.msgTransporter, pm.vnetController)
|
||||
pxy := NewWrapper(pm.ctx, cfg, pm.clientCfg, pm.encryptionKey, pm.HandleEvent, pm.msgTransporter, pm.vnetController, pm.udpPacketCodec)
|
||||
if pm.inWorkConnCallback != nil {
|
||||
pxy.SetInWorkConnCallback(pm.inWorkConnCallback)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
// Copyright 2026 The frp Authors
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//go:build !frps
|
||||
|
||||
package proxy
|
||||
|
||||
import (
|
||||
"io"
|
||||
"net"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
v1 "github.com/fatedier/frp/pkg/config/v1"
|
||||
"github.com/fatedier/frp/pkg/msg"
|
||||
"github.com/fatedier/frp/pkg/util/xlog"
|
||||
)
|
||||
|
||||
func TestHandleTCPWorkConnectionRejectsInvalidAddress(t *testing.T) {
|
||||
workConn, peerConn := net.Pipe()
|
||||
defer peerConn.Close()
|
||||
|
||||
pxy := &BaseProxy{
|
||||
baseCfg: &v1.ProxyBaseConfig{},
|
||||
xl: xlog.New(),
|
||||
}
|
||||
pxy.HandleTCPWorkConnection(workConn, &msg.StartWorkConn{
|
||||
SrcAddr: "[",
|
||||
SrcPort: 1,
|
||||
}, nil)
|
||||
|
||||
buffer := make([]byte, 1)
|
||||
_, err := peerConn.Read(buffer)
|
||||
require.ErrorIs(t, err, io.EOF)
|
||||
}
|
||||
@@ -99,6 +99,7 @@ func NewWrapper(
|
||||
eventHandler event.Handler,
|
||||
msgTransporter transport.MessageTransporter,
|
||||
vnetController *vnet.Controller,
|
||||
udpPacketCodec string,
|
||||
) *Wrapper {
|
||||
baseInfo := cfg.GetBaseConfig()
|
||||
xl := xlog.FromContextSafe(ctx).Spawn().AppendPrefix(baseInfo.Name)
|
||||
@@ -127,7 +128,7 @@ func NewWrapper(
|
||||
xl.Tracef("enable health check monitor")
|
||||
}
|
||||
|
||||
pw.pxy = NewProxy(pw.ctx, pw.Cfg, clientCfg, encryptionKey, pw.msgTransporter, pw.vnetController)
|
||||
pw.pxy = NewProxy(pw.ctx, pw.Cfg, clientCfg, encryptionKey, pw.msgTransporter, pw.vnetController, udpPacketCodec)
|
||||
return pw
|
||||
}
|
||||
|
||||
|
||||
@@ -87,7 +87,13 @@ func (pxy *SUDPProxy) InWorkConn(conn net.Conn, _ *msg.StartWorkConn) {
|
||||
}
|
||||
|
||||
workConn := netpkg.WrapReadWriteCloserToConn(remote, conn)
|
||||
payloadConn := msg.NewConn(workConn, msg.NewReadWriter(workConn, pxy.clientCfg.Transport.WireProtocol))
|
||||
payloadRW, err := msg.NewUDPPacketReadWriter(workConn, pxy.clientCfg.Transport.WireProtocol, pxy.udpPacketCodec)
|
||||
if err != nil {
|
||||
xl.Errorf("create SUDP packet read writer: %v", err)
|
||||
_ = workConn.Close()
|
||||
return
|
||||
}
|
||||
payloadConn := msg.NewConn(workConn, payloadRW)
|
||||
readCh := make(chan *msg.UDPPacket, 1024)
|
||||
sendCh := make(chan msg.Message, 1024)
|
||||
isClose := false
|
||||
|
||||
+10
-3
@@ -97,10 +97,17 @@ func (pxy *UDPProxy) InWorkConn(conn net.Conn, _ *msg.StartWorkConn) {
|
||||
return
|
||||
}
|
||||
|
||||
pxy.mu.Lock()
|
||||
pxy.workConn = netpkg.WrapReadWriteCloserToConn(remote, conn)
|
||||
workConn := netpkg.WrapReadWriteCloserToConn(remote, conn)
|
||||
// Plain UDP payload follows the configured wire protocol for message framing.
|
||||
payloadRW := msg.NewReadWriter(pxy.workConn, pxy.clientCfg.Transport.WireProtocol)
|
||||
payloadRW, err := msg.NewUDPPacketReadWriter(workConn, pxy.clientCfg.Transport.WireProtocol, pxy.udpPacketCodec)
|
||||
if err != nil {
|
||||
xl.Errorf("create UDP packet read writer: %v", err)
|
||||
workConn.Close()
|
||||
return
|
||||
}
|
||||
|
||||
pxy.mu.Lock()
|
||||
pxy.workConn = workConn
|
||||
pxy.readCh = make(chan *msg.UDPPacket, 1024)
|
||||
pxy.sendCh = make(chan msg.Message, 1024)
|
||||
pxy.closed = false
|
||||
|
||||
+16
-4
@@ -22,6 +22,7 @@ import (
|
||||
"net/http"
|
||||
"os"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/fatedier/golib/crypto"
|
||||
@@ -32,6 +33,7 @@ import (
|
||||
"github.com/fatedier/frp/pkg/config"
|
||||
"github.com/fatedier/frp/pkg/config/source"
|
||||
v1 "github.com/fatedier/frp/pkg/config/v1"
|
||||
"github.com/fatedier/frp/pkg/config/v1/validation"
|
||||
"github.com/fatedier/frp/pkg/msg"
|
||||
"github.com/fatedier/frp/pkg/policy/security"
|
||||
httppkg "github.com/fatedier/frp/pkg/util/http"
|
||||
@@ -109,6 +111,9 @@ func setServiceOptionsDefault(options *ServiceOptions) error {
|
||||
// Service is the client service that connects to frps and provides proxy services.
|
||||
type Service struct {
|
||||
ctlMu sync.RWMutex
|
||||
// Stores gracefulShutdownDuration independently from ctlMu, because the
|
||||
// graceful shutdown wait may hold ctlMu for an arbitrary duration.
|
||||
gracefulShutdownDuration atomic.Int64
|
||||
// manager control connection with server
|
||||
ctl *Control
|
||||
// Uniq id got from frps, it will be attached to loginMsg.
|
||||
@@ -149,8 +154,7 @@ type Service struct {
|
||||
// service context
|
||||
ctx context.Context
|
||||
// call cancel to stop service
|
||||
cancel context.CancelCauseFunc
|
||||
gracefulShutdownDuration time.Duration
|
||||
cancel context.CancelCauseFunc
|
||||
|
||||
connectorCreator func(context.Context, *v1.ClientCommonConfig) Connector
|
||||
handleWorkConnCb func(*v1.ProxyBaseConfig, net.Conn, *msg.StartWorkConn) bool
|
||||
@@ -412,7 +416,7 @@ func (svr *Service) Close() {
|
||||
}
|
||||
|
||||
func (svr *Service) GracefulClose(d time.Duration) {
|
||||
svr.gracefulShutdownDuration = d
|
||||
svr.gracefulShutdownDuration.Store(int64(d))
|
||||
svr.cancel(nil)
|
||||
}
|
||||
|
||||
@@ -429,7 +433,8 @@ func (svr *Service) stop() {
|
||||
svr.ctlMu.Lock()
|
||||
defer svr.ctlMu.Unlock()
|
||||
if svr.ctl != nil {
|
||||
svr.ctl.GracefulClose(svr.gracefulShutdownDuration)
|
||||
d := time.Duration(svr.gracefulShutdownDuration.Load())
|
||||
svr.ctl.GracefulClose(d)
|
||||
svr.ctl = nil
|
||||
}
|
||||
if svr.webServer != nil {
|
||||
@@ -506,6 +511,13 @@ func (svr *Service) reloadConfigFromSourcesLocked() error {
|
||||
proxies, visitors = config.FilterClientConfigurers(reloadCommon, proxies, visitors)
|
||||
proxies = config.CompleteProxyConfigurers(proxies)
|
||||
visitors = config.CompleteVisitorConfigurers(visitors)
|
||||
requirements := validation.GetClientConfigRequirements(reloadCommon, proxies, visitors)
|
||||
if svr.vnetController == nil && requirements.VirtualNet {
|
||||
return errors.New(
|
||||
"VirtualNet-dependent configuration requires a VirtualNet runtime enabled at startup; " +
|
||||
"restart frpc after configuring featureGates.VirtualNet and virtualNet.address",
|
||||
)
|
||||
}
|
||||
|
||||
// Atomically replace the entire configuration
|
||||
if err := svr.UpdateAllConfigurer(proxies, visitors); err != nil {
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
package client
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/fatedier/frp/client/proxy"
|
||||
"github.com/fatedier/frp/client/visitor"
|
||||
v1 "github.com/fatedier/frp/pkg/config/v1"
|
||||
"github.com/fatedier/frp/pkg/msg"
|
||||
)
|
||||
|
||||
type gracefulCloseTestConnector struct {
|
||||
conn net.Conn
|
||||
}
|
||||
|
||||
func (*gracefulCloseTestConnector) Connect() (*msg.Conn, error) { return nil, net.ErrClosed }
|
||||
func (c *gracefulCloseTestConnector) Close() error { return c.conn.Close() }
|
||||
|
||||
func newGracefulCloseTestService() *Service {
|
||||
ctx := context.Background()
|
||||
common := &v1.ClientCommonConfig{}
|
||||
serverConn, clientConn := net.Pipe()
|
||||
ctl := &Control{
|
||||
ctx: ctx,
|
||||
sessionCtx: &SessionContext{
|
||||
Common: common,
|
||||
RunID: "graceful-close-race",
|
||||
Conn: msg.NewConn(clientConn, msg.NewV1ReadWriter(clientConn)),
|
||||
Connector: &gracefulCloseTestConnector{conn: serverConn},
|
||||
},
|
||||
doneCh: make(chan struct{}),
|
||||
}
|
||||
ctl.pm = proxy.NewManager(ctx, common, nil, nil, nil, "")
|
||||
ctl.vm = visitor.NewManager(ctx, "graceful-close-race", common, nil, nil, nil, "")
|
||||
return &Service{ctl: ctl, cancel: context.CancelCauseFunc(func(error) {})}
|
||||
}
|
||||
|
||||
func TestGracefulCloseAndStopSynchronizeDuration(t *testing.T) {
|
||||
for i := range 10000 {
|
||||
svr := newGracefulCloseTestService()
|
||||
start := make(chan struct{})
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(2)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
svr.GracefulClose(time.Duration(i))
|
||||
}()
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
svr.stop()
|
||||
}()
|
||||
close(start)
|
||||
wg.Wait()
|
||||
}
|
||||
}
|
||||
|
||||
func TestGracefulCloseDoesNotBlockDuringStop(t *testing.T) {
|
||||
const gracefulDuration = 200 * time.Millisecond
|
||||
|
||||
svr := newGracefulCloseTestService()
|
||||
svr.GracefulClose(gracefulDuration)
|
||||
stopDone := make(chan struct{})
|
||||
go func() {
|
||||
svr.stop()
|
||||
close(stopDone)
|
||||
}()
|
||||
defer func() {
|
||||
select {
|
||||
case <-stopDone:
|
||||
case <-time.After(time.Second):
|
||||
t.Error("stop did not finish")
|
||||
}
|
||||
}()
|
||||
|
||||
deadline := time.Now().Add(time.Second)
|
||||
for svr.ctlMu.TryLock() {
|
||||
svr.ctlMu.Unlock()
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatal("stop did not acquire ctlMu")
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
svr.GracefulClose(0)
|
||||
if elapsed := time.Since(start); elapsed >= gracefulDuration/2 {
|
||||
t.Fatalf("GracefulClose blocked for %v while stop was waiting", elapsed)
|
||||
}
|
||||
}
|
||||
@@ -113,7 +113,13 @@ func (sv *SUDPVisitor) dispatcher() {
|
||||
func (sv *SUDPVisitor) worker(workConn net.Conn, firstPacket *msg.UDPPacket) {
|
||||
xl := xlog.FromContextSafe(sv.ctx)
|
||||
xl.Debugf("starting sudp proxy worker")
|
||||
payloadConn := msg.NewConn(workConn, msg.NewReadWriter(workConn, sv.clientCfg.Transport.WireProtocol))
|
||||
payloadRW, err := msg.NewUDPPacketReadWriter(workConn, sv.clientCfg.Transport.WireProtocol, udpPacketCodecFromHelper(sv.helper))
|
||||
if err != nil {
|
||||
xl.Errorf("create SUDP packet read writer: %v", err)
|
||||
_ = workConn.Close()
|
||||
return
|
||||
}
|
||||
payloadConn := msg.NewConn(workConn, payloadRW)
|
||||
|
||||
wg := &sync.WaitGroup{}
|
||||
wg.Add(2)
|
||||
|
||||
@@ -50,6 +50,17 @@ type Helper interface {
|
||||
RunID() string
|
||||
}
|
||||
|
||||
type udpPacketCodecProvider interface {
|
||||
UDPPacketCodec() string
|
||||
}
|
||||
|
||||
func udpPacketCodecFromHelper(helper Helper) string {
|
||||
if provider, ok := helper.(udpPacketCodecProvider); ok {
|
||||
return provider.UDPPacketCodec()
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Visitor is used for forward traffics from local port tot remote service.
|
||||
type Visitor interface {
|
||||
Run() error
|
||||
|
||||
@@ -53,7 +53,12 @@ func NewManager(
|
||||
connectServer func() (*msg.Conn, error),
|
||||
msgTransporter transport.MessageTransporter,
|
||||
vnetController *vnet.Controller,
|
||||
udpPacketCodecs ...string,
|
||||
) *Manager {
|
||||
udpPacketCodec := ""
|
||||
if len(udpPacketCodecs) > 0 {
|
||||
udpPacketCodec = udpPacketCodecs[0]
|
||||
}
|
||||
m := &Manager{
|
||||
clientCfg: clientCfg,
|
||||
cfgs: make(map[string]v1.VisitorConfigurer),
|
||||
@@ -68,6 +73,7 @@ func NewManager(
|
||||
vnetController: vnetController,
|
||||
transferConnFn: m.TransferConn,
|
||||
runID: runID,
|
||||
udpPacketCodec: udpPacketCodec,
|
||||
}
|
||||
return m
|
||||
}
|
||||
@@ -205,6 +211,7 @@ type visitorHelperImpl struct {
|
||||
vnetController *vnet.Controller
|
||||
transferConnFn func(name string, conn net.Conn) error
|
||||
runID string
|
||||
udpPacketCodec string
|
||||
}
|
||||
|
||||
func (v *visitorHelperImpl) ConnectServer() (*msg.Conn, error) {
|
||||
@@ -226,3 +233,7 @@ func (v *visitorHelperImpl) VNetController() *vnet.Controller {
|
||||
func (v *visitorHelperImpl) RunID() string {
|
||||
return v.runID
|
||||
}
|
||||
|
||||
func (v *visitorHelperImpl) UDPPacketCodec() string {
|
||||
return v.udpPacketCodec
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user