From d06cf792433554ca85f4770f05a475236e72c761 Mon Sep 17 00:00:00 2001 From: Fangliding Date: Thu, 8 Oct 2026 15:35:39 +0800 Subject: [PATCH] Hysteria: Remove closed instance conn --- core/xray.go | 11 +++--- transport/internet/hysteria/dialer.go | 50 +++++++++++++++++++-------- 2 files changed, 42 insertions(+), 19 deletions(-) diff --git a/core/xray.go b/core/xray.go index a70507068..3e785a823 100644 --- a/core/xray.go +++ b/core/xray.go @@ -4,6 +4,7 @@ import ( "context" "reflect" "sync" + "sync/atomic" "github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common/errors" @@ -84,7 +85,7 @@ type Instance struct { features []features.Feature pendingResolutions []resolution pendingOptionalResolutions []resolution - running bool + running atomic.Bool resolveLock sync.Mutex ctx context.Context @@ -92,7 +93,7 @@ type Instance struct { // Instance state func (server *Instance) IsRunning() bool { - return server.running + return server.running.Load() } func AddInboundHandler(server *Instance, config *InboundHandlerConfig) error { @@ -262,7 +263,7 @@ func (s *Instance) Close() error { s.statusLock.Lock() defer s.statusLock.Unlock() - s.running = false + s.running.Store(false) var errs []interface{} for _, f := range s.features { @@ -320,7 +321,7 @@ func (s *Instance) RequireFeatures(callback interface{}, optional bool) error { // AddFeature registers a feature into current Instance. func (s *Instance) AddFeature(feature features.Feature) error { - if s.running { + if s.running.Load() { if err := feature.Start(); err != nil { errors.LogInfoInner(s.ctx, err, "failed to start feature") } @@ -388,7 +389,7 @@ func (s *Instance) Start() error { s.statusLock.Lock() defer s.statusLock.Unlock() - s.running = true + s.running.Store(true) for _, f := range s.features { if err := f.Start(); err != nil { return err diff --git a/transport/internet/hysteria/dialer.go b/transport/internet/hysteria/dialer.go index 4b774c353..bfd92eabd 100644 --- a/transport/internet/hysteria/dialer.go +++ b/transport/internet/hysteria/dialer.go @@ -17,6 +17,7 @@ import ( "github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/net" "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/finalmask" "github.com/xtls/xray-core/transport/internet/hysteria/congestion" @@ -35,10 +36,11 @@ type client struct { finalMask *finalmask.FinalMask quicParams *internet.QuicParams - conn *quic.Conn - tr *quic.Transport - pktConn net.PacketConn - udpSM *udpSessionManager + conn *quic.Conn + tr *quic.Transport + pktConn net.PacketConn + udpSM *udpSessionManager + instance *core.Instance } func (c *client) status() status { @@ -54,9 +56,11 @@ func (c *client) status() status { } func (c *client) close() { - c.conn.CloseWithError(closeErrCodeOK, "") - c.tr.Close() - c.pktConn.Close() + if c.conn != nil { + c.conn.CloseWithError(closeErrCodeOK, "") + } + common.CloseIfExists(c.tr) + common.CloseIfExists(c.pktConn) c.conn = nil c.tr = nil c.pktConn = nil @@ -256,12 +260,18 @@ func (c *client) udp(ctx context.Context) (stat.Connection, error) { return c.udpSM.udp() } -func (c *client) clean() { +func (c *client) clean() (shouldDelete bool) { c.Lock() + defer c.Unlock() + if c.instance != nil && !c.instance.IsRunning() { + c.close() + return true + } if c.status() == StatusInactive { c.close() + return false } - c.Unlock() + return false } type dialerConf struct { @@ -277,11 +287,21 @@ type clientManager struct { func (m *clientManager) clean() { ticker := time.NewTicker(idleCleanupInterval) for range ticker.C { + var toDelete []dialerConf m.RLock() - for _, c := range m.m { - c.clean() + for k, c := range m.m { + if c.clean() { + toDelete = append(toDelete, k) + } } m.RUnlock() + if len(toDelete) > 0 { + m.Lock() + for _, k := range toDelete { + delete(m.m, k) + } + m.Unlock() + } } } @@ -306,15 +326,17 @@ func Dial(ctx context.Context, dest net.Destination, streamSettings *internet.Me go manager.clean() }) + dialerConfKey := dialerConf{dest, streamSettings} manager.RLock() - c := manager.m[dialerConf{dest, streamSettings}] + c := manager.m[dialerConfKey] manager.RUnlock() if c == nil { manager.Lock() - c = manager.m[dialerConf{dest, streamSettings}] + c = manager.m[dialerConfKey] if c == nil { c = &client{ + instance: core.FromContext(ctx), dest: dest, config: streamSettings.ProtocolSettings.(*Config), tlsConfig: tlsConfig.GetTLSConfig(tls.WithDestination(dest)), @@ -322,7 +344,7 @@ func Dial(ctx context.Context, dest net.Destination, streamSettings *internet.Me finalMask: streamSettings.FinalMask, quicParams: streamSettings.QuicParams, } - manager.m[dialerConf{dest, streamSettings}] = c + manager.m[dialerConfKey] = c } manager.Unlock() }