Compare commits

...
1 Commits
Author SHA1 Message Date
Fangliding d06cf79243 Hysteria: Remove closed instance conn 2026-10-08 15:35:39 +08:00
2 changed files with 42 additions and 19 deletions
+6 -5
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"reflect" "reflect"
"sync" "sync"
"sync/atomic"
"github.com/xtls/xray-core/common" "github.com/xtls/xray-core/common"
"github.com/xtls/xray-core/common/errors" "github.com/xtls/xray-core/common/errors"
@@ -84,7 +85,7 @@ type Instance struct {
features []features.Feature features []features.Feature
pendingResolutions []resolution pendingResolutions []resolution
pendingOptionalResolutions []resolution pendingOptionalResolutions []resolution
running bool running atomic.Bool
resolveLock sync.Mutex resolveLock sync.Mutex
ctx context.Context ctx context.Context
@@ -92,7 +93,7 @@ type Instance struct {
// Instance state // Instance state
func (server *Instance) IsRunning() bool { func (server *Instance) IsRunning() bool {
return server.running return server.running.Load()
} }
func AddInboundHandler(server *Instance, config *InboundHandlerConfig) error { func AddInboundHandler(server *Instance, config *InboundHandlerConfig) error {
@@ -262,7 +263,7 @@ func (s *Instance) Close() error {
s.statusLock.Lock() s.statusLock.Lock()
defer s.statusLock.Unlock() defer s.statusLock.Unlock()
s.running = false s.running.Store(false)
var errs []interface{} var errs []interface{}
for _, f := range s.features { 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. // 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.running.Load() {
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")
} }
@@ -388,7 +389,7 @@ func (s *Instance) Start() error {
s.statusLock.Lock() s.statusLock.Lock()
defer s.statusLock.Unlock() defer s.statusLock.Unlock()
s.running = true s.running.Store(true)
for _, f := range s.features { for _, f := range s.features {
if err := f.Start(); err != nil { if err := f.Start(); err != nil {
return err return err
+36 -14
View File
@@ -17,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"
@@ -35,10 +36,11 @@ 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
instance *core.Instance
} }
func (c *client) status() status { func (c *client) status() status {
@@ -54,9 +56,11 @@ func (c *client) status() status {
} }
func (c *client) close() { func (c *client) close() {
c.conn.CloseWithError(closeErrCodeOK, "") if c.conn != nil {
c.tr.Close() c.conn.CloseWithError(closeErrCodeOK, "")
c.pktConn.Close() }
common.CloseIfExists(c.tr)
common.CloseIfExists(c.pktConn)
c.conn = nil c.conn = nil
c.tr = nil c.tr = nil
c.pktConn = nil c.pktConn = nil
@@ -256,12 +260,18 @@ 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() (shouldDelete bool) {
c.Lock() c.Lock()
defer c.Unlock()
if c.instance != nil && !c.instance.IsRunning() {
c.close()
return true
}
if c.status() == StatusInactive { if c.status() == StatusInactive {
c.close() c.close()
return false
} }
c.Unlock() return false
} }
type dialerConf struct { type dialerConf struct {
@@ -277,11 +287,21 @@ 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 toDelete []dialerConf
m.RLock() m.RLock()
for _, c := range m.m { for k, c := range m.m {
c.clean() if c.clean() {
toDelete = append(toDelete, k)
}
} }
m.RUnlock() 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() 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)),
@@ -322,7 +344,7 @@ func Dial(ctx context.Context, dest net.Destination, streamSettings *internet.Me
finalMask: streamSettings.FinalMask, finalMask: streamSettings.FinalMask,
quicParams: streamSettings.QuicParams, quicParams: streamSettings.QuicParams,
} }
manager.m[dialerConf{dest, streamSettings}] = c manager.m[dialerConfKey] = c
} }
manager.Unlock() manager.Unlock()
} }