fix: 完成 broker 复审 B-03 至 B-12

每连接异步下发与背压、写出后断开、校验当前连接与订阅、生命周期串行、登录条件更新、闲置按在线计、认证超时并发与 Shutdown 0x8B。
This commit is contained in:
Nixevol
2026-09-30 15:24:15 +08:00
parent 42160720bd
commit 8b845843d6
12 changed files with 1099 additions and 89 deletions
+5 -1
View File
@@ -153,7 +153,7 @@ func runServe(ctx context.Context, cfg config.Config) error {
Sessions: sessionTokens,
MaxScheduleSeconds: int64(cfg.Limits.MaxScheduleSeconds),
Logger: slog.Default(),
ConnControl: brk,
ConnControl: nil, // B-04:踢线走 Session 钩子,避免 identity 20ms 异步 Disconnect
Downlink: brk,
ClientIP: func(r *http.Request) string {
return httpx.ClientIP(r, trustedNets)
@@ -310,6 +310,10 @@ func runServe(ctx context.Context, cfg config.Config) error {
<-ctx.Done()
loopCancel()
// B-08:先对 MQTT 连接发 0x8B。HTTP Shutdown 与监听器完整停机顺序见 L-03。
shutCtx, shutCancel := context.WithTimeout(context.Background(), 5*time.Second)
_ = brk.Shutdown(shutCtx)
shutCancel()
_ = lnSrv.Close()
drainCtx, drainCancel := context.WithTimeout(context.Background(), 10*time.Second)
defer drainCancel()
+42 -16
View File
@@ -5,12 +5,14 @@ import (
"encoding/json"
"errors"
"log/slog"
"sync"
"git.asio.asia/nixevol/NixMsg/internal/app/group"
"git.asio.asia/nixevol/NixMsg/internal/app/identity"
"git.asio.asia/nixevol/NixMsg/internal/app/message"
"git.asio.asia/nixevol/NixMsg/internal/app/port"
"git.asio.asia/nixevol/NixMsg/internal/app/presence"
"git.asio.asia/nixevol/NixMsg/internal/broker"
"git.asio.asia/nixevol/NixMsg/internal/metrics"
"git.asio.asia/nixevol/NixMsg/internal/protocol"
)
@@ -25,9 +27,31 @@ type appUplink struct {
down port.Downlink
log *slog.Logger
metrics *metrics.Registry
lifeMu sync.Mutex
lifeLocks map[string]*sync.Mutex
hsMu sync.Mutex
handshake map[port.ConnID]string // 已 hello 的连接代号 → 端编号
}
func (u *appUplink) epLife(endpointID string) *sync.Mutex {
u.lifeMu.Lock()
defer u.lifeMu.Unlock()
if u.lifeLocks == nil {
u.lifeLocks = make(map[string]*sync.Mutex)
}
m := u.lifeLocks[endpointID]
if m == nil {
m = &sync.Mutex{}
u.lifeLocks[endpointID] = m
}
return m
}
func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo) error {
lk := u.epLife(conn.EndpointID)
lk.Lock()
defer lk.Unlock()
u.conns.Set(conn.EndpointID, message.LiveConn{
ConnID: conn.ConnID,
MaxPacketSize: conn.MaxPacketSize,
@@ -36,6 +60,15 @@ func (u *appUplink) OnSessionEstablished(ctx context.Context, conn port.ConnInfo
}
func (u *appUplink) OnHandshakeComplete(ctx context.Context, hs port.HandshakeInfo) error {
lk := u.epLife(hs.EndpointID)
lk.Lock()
defer lk.Unlock()
u.hsMu.Lock()
if u.handshake == nil {
u.handshake = make(map[port.ConnID]string)
}
u.handshake[hs.ConnID] = hs.EndpointID
u.hsMu.Unlock()
live := message.LiveConn{
ConnID: hs.ConnID,
MaxReceiveBytes: hs.MaxReceiveBytes,
@@ -50,11 +83,18 @@ func (u *appUplink) OnHandshakeComplete(ctx context.Context, hs port.HandshakeIn
}
func (u *appUplink) OnDisconnect(ctx context.Context, conn port.ConnInfo, reason port.DisconnectReason) {
lk := u.epLife(conn.EndpointID)
lk.Lock()
defer lk.Unlock()
if u.presence != nil {
u.presence.ClearWatch(conn.ConnID)
}
u.hsMu.Lock()
_, handshook := u.handshake[conn.ConnID]
delete(u.handshake, conn.ConnID)
u.hsMu.Unlock()
live, ok := u.conns.Current(conn.EndpointID)
isCurrent := ok && live.ConnID == conn.ConnID
isCurrent := ok && live.ConnID == conn.ConnID && handshook
if err := u.msg.OnDisconnect(ctx, conn.EndpointID, conn.ConnID, isCurrent); err != nil {
u.log.Error("message disconnect", "endpoint", conn.EndpointID, "err", err)
}
@@ -280,7 +320,7 @@ func (u *appUplink) publishResp(ctx context.Context, conn port.ConnInfo, resp pr
return
}
if live, ok := u.conns.Current(conn.EndpointID); ok && live.ConnID == conn.ConnID {
limit := respPayloadLimit(live.MaxPacketSize, live.MaxReceiveBytes)
limit := broker.EffectivePayloadLimit(live.MaxPacketSize, live.MaxReceiveBytes)
if limit > 0 && len(b) > limit {
tooLarge := protocol.Resp{
V: protocol.Version,
@@ -300,20 +340,6 @@ func (u *appUplink) publishResp(ctx context.Context, conn port.ConnInfo, resp pr
}
}
func respPayloadLimit(maxPacketSize uint32, maxRecvBytes int) int {
limit := 0
if maxRecvBytes > 0 {
limit = maxRecvBytes
}
if maxPacketSize > 0 {
n := int(maxPacketSize)
if limit == 0 || n < limit {
limit = n
}
}
return limit
}
func peekRID(payload []byte) string {
var peek struct {
RID string `json:"rid"`