Compare commits

..
7 Commits
Author SHA1 Message Date
null 2b0d898d08 use RLock 2026-10-08 23:38:14 +08:00
null 8f705a1ea4 fix logic plus 2026-10-08 22:35:35 +08:00
null 9897b5a9d3 refine logic 2026-10-08 21:12:51 +08:00
null 6bcf81b0ab fix logic 2026-10-08 20:38:06 +08:00
null 542030ca71 chore 2026-10-08 19:30:46 +08:00
Fangliding 0d83d099ac make client unavailable if instance stopped 2026-10-08 19:03:48 +08:00
Fangliding d06cf79243 Hysteria: Remove closed instance conn 2026-10-08 15:35:39 +08:00
3 changed files with 47 additions and 55 deletions
+3 -1
View File
@@ -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")
} }
+2 -10
View File
@@ -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()
} }
+42 -44
View File
@@ -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()
} }