mirror of
https://github.com/XTLS/Xray-core.git
synced 2026-10-08 16:50:07 +00:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2b0d898d08 | ||
|
|
8f705a1ea4 | ||
|
|
9897b5a9d3 | ||
|
|
6bcf81b0ab | ||
|
|
542030ca71 | ||
|
|
0d83d099ac | ||
|
|
d06cf79243 |
+3
-1
@@ -92,6 +92,8 @@ type Instance struct {
|
|||||||
|
|
||||||
// Instance state
|
// Instance state
|
||||||
func (server *Instance) IsRunning() bool {
|
func (server *Instance) IsRunning() bool {
|
||||||
|
server.statusLock.Lock()
|
||||||
|
defer server.statusLock.Unlock()
|
||||||
return server.running
|
return server.running
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -320,7 +322,7 @@ func (s *Instance) RequireFeatures(callback interface{}, optional bool) error {
|
|||||||
|
|
||||||
// AddFeature registers a feature into current Instance.
|
// AddFeature registers a feature into current Instance.
|
||||||
func (s *Instance) AddFeature(feature features.Feature) error {
|
func (s *Instance) AddFeature(feature features.Feature) error {
|
||||||
if s.running {
|
if s.IsRunning() {
|
||||||
if err := feature.Start(); err != nil {
|
if err := feature.Start(); err != nil {
|
||||||
errors.LogInfoInner(s.ctx, err, "failed to start feature")
|
errors.LogInfoInner(s.ctx, err, "failed to start feature")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,7 +5,6 @@ import (
|
|||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"io"
|
"io"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/apernet/quic-go"
|
"github.com/apernet/quic-go"
|
||||||
@@ -21,10 +20,8 @@ type interConn struct {
|
|||||||
local net.Addr
|
local net.Addr
|
||||||
remote net.Addr
|
remote net.Addr
|
||||||
|
|
||||||
client bool
|
client bool
|
||||||
user *protocol.MemoryUser
|
user *protocol.MemoryUser
|
||||||
closeOnce sync.Once
|
|
||||||
aliveTCP *atomic.Int64
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *interConn) User() *protocol.MemoryUser {
|
func (c *interConn) User() *protocol.MemoryUser {
|
||||||
@@ -49,11 +46,6 @@ func (c *interConn) Write(b []byte) (int, error) {
|
|||||||
|
|
||||||
func (c *interConn) Close() error {
|
func (c *interConn) Close() error {
|
||||||
c.stream.CancelRead(0)
|
c.stream.CancelRead(0)
|
||||||
if c.aliveTCP != nil {
|
|
||||||
c.closeOnce.Do(func() {
|
|
||||||
c.aliveTCP.Add(-1)
|
|
||||||
})
|
|
||||||
}
|
|
||||||
return c.stream.Close()
|
return c.stream.Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ import (
|
|||||||
"runtime"
|
"runtime"
|
||||||
"strconv"
|
"strconv"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/apernet/quic-go"
|
"github.com/apernet/quic-go"
|
||||||
@@ -18,6 +17,7 @@ import (
|
|||||||
"github.com/xtls/xray-core/common/errors"
|
"github.com/xtls/xray-core/common/errors"
|
||||||
"github.com/xtls/xray-core/common/net"
|
"github.com/xtls/xray-core/common/net"
|
||||||
"github.com/xtls/xray-core/common/net/cnc"
|
"github.com/xtls/xray-core/common/net/cnc"
|
||||||
|
"github.com/xtls/xray-core/core"
|
||||||
"github.com/xtls/xray-core/transport/internet"
|
"github.com/xtls/xray-core/transport/internet"
|
||||||
"github.com/xtls/xray-core/transport/internet/finalmask"
|
"github.com/xtls/xray-core/transport/internet/finalmask"
|
||||||
"github.com/xtls/xray-core/transport/internet/hysteria/congestion"
|
"github.com/xtls/xray-core/transport/internet/hysteria/congestion"
|
||||||
@@ -29,6 +29,8 @@ import (
|
|||||||
type client struct {
|
type client struct {
|
||||||
sync.Mutex
|
sync.Mutex
|
||||||
|
|
||||||
|
instance *core.Instance
|
||||||
|
forced bool
|
||||||
dest net.Destination
|
dest net.Destination
|
||||||
config *Config
|
config *Config
|
||||||
tlsConfig *gotls.Config
|
tlsConfig *gotls.Config
|
||||||
@@ -36,12 +38,10 @@ type client struct {
|
|||||||
finalMask *finalmask.FinalMask
|
finalMask *finalmask.FinalMask
|
||||||
quicParams *internet.QuicParams
|
quicParams *internet.QuicParams
|
||||||
|
|
||||||
conn *quic.Conn
|
conn *quic.Conn
|
||||||
tr *quic.Transport
|
tr *quic.Transport
|
||||||
pktConn net.PacketConn
|
pktConn net.PacketConn
|
||||||
udpSM *udpSessionManager
|
udpSM *udpSessionManager
|
||||||
aliveTCP *atomic.Int64
|
|
||||||
lastAlive time.Time
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *client) status() status {
|
func (c *client) status() status {
|
||||||
@@ -67,11 +67,14 @@ func (c *client) close() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (c *client) dial(ctx context.Context) error {
|
func (c *client) dial(ctx context.Context) error {
|
||||||
status := c.status()
|
if c.forced {
|
||||||
if status == StatusActive {
|
return errors.New("client is closed")
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
if status == StatusInactive {
|
|
||||||
|
switch c.status() {
|
||||||
|
case StatusActive:
|
||||||
|
return nil
|
||||||
|
case StatusInactive:
|
||||||
c.close()
|
c.close()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -93,7 +96,7 @@ func (c *client) dial(ctx context.Context) error {
|
|||||||
ChromeParrot: !quicParams.DisableChromeParrot,
|
ChromeParrot: !quicParams.DisableChromeParrot,
|
||||||
EnableDatagrams: true,
|
EnableDatagrams: true,
|
||||||
MaxDatagramFrameSize: MaxDatagramFrameSize,
|
MaxDatagramFrameSize: MaxDatagramFrameSize,
|
||||||
OmitMaxDatagramFrameSize: time.Now().After(time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC)),
|
OmitMaxDatagramFrameSize: true,
|
||||||
DisablePathManager: true,
|
DisablePathManager: true,
|
||||||
}
|
}
|
||||||
if quicParams.InitStreamReceiveWindow == 0 {
|
if quicParams.InitStreamReceiveWindow == 0 {
|
||||||
@@ -238,15 +241,12 @@ func (c *client) tcp(ctx context.Context) (stat.Connection, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
c.lastAlive = time.Now()
|
|
||||||
c.aliveTCP.Add(1)
|
|
||||||
return &interConn{
|
return &interConn{
|
||||||
stream: stream,
|
stream: stream,
|
||||||
local: c.conn.LocalAddr(),
|
local: c.conn.LocalAddr(),
|
||||||
remote: c.conn.RemoteAddr(),
|
remote: c.conn.RemoteAddr(),
|
||||||
|
|
||||||
client: true,
|
client: true,
|
||||||
aliveTCP: c.aliveTCP,
|
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -262,30 +262,15 @@ func (c *client) udp(ctx context.Context) (stat.Connection, error) {
|
|||||||
return c.udpSM.udp()
|
return c.udpSM.udp()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *client) clean() {
|
func (c *client) clean(force bool) {
|
||||||
c.Lock()
|
c.Lock()
|
||||||
defer c.Unlock()
|
if force {
|
||||||
switch c.status() {
|
c.forced = true
|
||||||
case StatusInactive:
|
}
|
||||||
|
if status := c.status(); force && status != StatusNull || status == StatusInactive {
|
||||||
c.close()
|
c.close()
|
||||||
return
|
|
||||||
case StatusNull:
|
|
||||||
return
|
|
||||||
}
|
|
||||||
var udpSessions int
|
|
||||||
if c.udpSM != nil {
|
|
||||||
c.udpSM.RLock()
|
|
||||||
udpSessions = len(c.udpSM.m)
|
|
||||||
c.udpSM.RUnlock()
|
|
||||||
}
|
|
||||||
if udpSessions == 0 && c.aliveTCP.Load() == 0 {
|
|
||||||
if c.lastAlive.Add(net.ConnIdleTimeout).Before(time.Now()) {
|
|
||||||
c.close()
|
|
||||||
return
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
c.lastAlive = time.Now()
|
|
||||||
}
|
}
|
||||||
|
c.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
type dialerConf struct {
|
type dialerConf struct {
|
||||||
@@ -301,11 +286,23 @@ type clientManager struct {
|
|||||||
func (m *clientManager) clean() {
|
func (m *clientManager) clean() {
|
||||||
ticker := time.NewTicker(idleCleanupInterval)
|
ticker := time.NewTicker(idleCleanupInterval)
|
||||||
for range ticker.C {
|
for range ticker.C {
|
||||||
|
var forced []dialerConf
|
||||||
|
|
||||||
m.RLock()
|
m.RLock()
|
||||||
for _, c := range m.m {
|
for k, c := range m.m {
|
||||||
c.clean()
|
force := c.instance != nil && !c.instance.IsRunning()
|
||||||
|
c.clean(force)
|
||||||
|
if force {
|
||||||
|
forced = append(forced, k)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
m.RUnlock()
|
m.RUnlock()
|
||||||
|
|
||||||
|
for i := range forced {
|
||||||
|
m.Lock()
|
||||||
|
delete(m.m, forced[i])
|
||||||
|
m.Unlock()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -330,25 +327,26 @@ func Dial(ctx context.Context, dest net.Destination, streamSettings *internet.Me
|
|||||||
go manager.clean()
|
go manager.clean()
|
||||||
})
|
})
|
||||||
|
|
||||||
|
dialerConfKey := dialerConf{dest, streamSettings}
|
||||||
|
|
||||||
manager.RLock()
|
manager.RLock()
|
||||||
c := manager.m[dialerConf{dest, streamSettings}]
|
c := manager.m[dialerConfKey]
|
||||||
manager.RUnlock()
|
manager.RUnlock()
|
||||||
|
|
||||||
if c == nil {
|
if c == nil {
|
||||||
manager.Lock()
|
manager.Lock()
|
||||||
c = manager.m[dialerConf{dest, streamSettings}]
|
c = manager.m[dialerConfKey]
|
||||||
if c == nil {
|
if c == nil {
|
||||||
c = &client{
|
c = &client{
|
||||||
|
instance: core.FromContext(ctx),
|
||||||
dest: dest,
|
dest: dest,
|
||||||
config: streamSettings.ProtocolSettings.(*Config),
|
config: streamSettings.ProtocolSettings.(*Config),
|
||||||
tlsConfig: tlsConfig.GetTLSConfig(tls.WithDestination(dest)),
|
tlsConfig: tlsConfig.GetTLSConfig(tls.WithDestination(dest)),
|
||||||
socketConfig: streamSettings.SocketSettings,
|
socketConfig: streamSettings.SocketSettings,
|
||||||
finalMask: streamSettings.FinalMask,
|
finalMask: streamSettings.FinalMask,
|
||||||
quicParams: streamSettings.QuicParams,
|
quicParams: streamSettings.QuicParams,
|
||||||
aliveTCP: &atomic.Int64{},
|
|
||||||
lastAlive: time.Now(),
|
|
||||||
}
|
}
|
||||||
manager.m[dialerConf{dest, streamSettings}] = c
|
manager.m[dialerConfKey] = c
|
||||||
}
|
}
|
||||||
manager.Unlock()
|
manager.Unlock()
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user