mirror of
https://github.com/XTLS/Xray-core.git
synced 2026-10-08 08:40:01 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d06cf79243 |
+6
-5
@@ -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
|
||||||
|
|||||||
@@ -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()
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user