Compare commits

...
1 Commits
Author SHA1 Message Date
Fangliding dbee0f6561 Track alive hy conn 2026-10-08 04:00:54 +08:00
2 changed files with 43 additions and 9 deletions
+8
View File
@@ -5,6 +5,7 @@ import (
"encoding/binary" "encoding/binary"
"io" "io"
"sync" "sync"
"sync/atomic"
"time" "time"
"github.com/apernet/quic-go" "github.com/apernet/quic-go"
@@ -22,6 +23,8 @@ type interConn struct {
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 {
@@ -46,6 +49,11 @@ 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()
} }
+28 -2
View File
@@ -9,6 +9,7 @@ import (
"runtime" "runtime"
"strconv" "strconv"
"sync" "sync"
"sync/atomic"
"time" "time"
"github.com/apernet/quic-go" "github.com/apernet/quic-go"
@@ -39,6 +40,8 @@ type client struct {
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 {
@@ -235,12 +238,15 @@ 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
} }
@@ -258,10 +264,28 @@ func (c *client) udp(ctx context.Context) (stat.Connection, error) {
func (c *client) clean() { func (c *client) clean() {
c.Lock() c.Lock()
if c.status() == StatusInactive { defer c.Unlock()
switch c.status() {
case 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 {
@@ -321,6 +345,8 @@ func Dial(ctx context.Context, dest net.Destination, streamSettings *internet.Me
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[dialerConf{dest, streamSettings}] = c
} }