Files
notify/client_session_runtime.go
T
b612 1f2e74acca fix(notify): 修复传输生命周期竞态,完善背压与协议边界
- 完善 stream/bulk DataID 分配、预留和双向命名空间,修复并发打开及 dedicated/shared 回退时的 ID 冲突
- 将收发、回复、恢复任务和 sidecar 绑定原始会话与物理连接,防止重连后的旧消息误操作新连接
- 加强 close/reset 身份校验及实例移除检查,修复 dedicated attach 失败、通道引用和资源回收竞态
- 收紧批量发送器停止准入,确保在途入队完成后统一清理请求、缓冲区和等待者
- 修复 record 满队列死锁、取消时序号消耗及关闭竞态,确保关闭有界并返回真实错误
- 增加协商式 record 逻辑半关闭,保留反向 ACK;通过 reset 传递 RecordFailure,避免背压掩盖原始失败原因
- 补齐帧长度、批次数量、序号溢出和未确认窗口校验,提前拒绝超限数据并按字节预算拆批
- 为入站分发增加全局及单连接的条数、字节预算和阻塞背压,关闭时唤醒等待者,消除正常断连日志噪音
- 完善 bulk 窗口释放失败处理与传输诊断,补充并发、重连、背压、协议边界及真实 TCP 回归覆盖
2026-09-23 15:33:17 +08:00

338 lines
8.3 KiB
Go

package notify
import (
"b612.me/stario"
"context"
"errors"
"net"
"sync/atomic"
)
type clientSessionRuntime struct {
transport *transportBinding
transportAttached bool
conn net.Conn
stopCtx context.Context
stopFn context.CancelFunc
transportStopCtx context.Context
transportStopFn context.CancelFunc
queue *stario.StarQueue
inboundDispatcher *inboundDispatcher
epoch uint64
suppressGoodByeOnStop *atomic.Bool
}
func newClientSessionRuntimeBase(stopCtx context.Context, stopFn context.CancelFunc) *clientSessionRuntime {
return &clientSessionRuntime{
stopCtx: stopCtx,
stopFn: stopFn,
inboundDispatcher: newInboundDispatcher(),
suppressGoodByeOnStop: &atomic.Bool{},
}
}
func prepareClientSessionRuntime(rt *clientSessionRuntime) *clientSessionRuntime {
if rt == nil {
return nil
}
if rt.inboundDispatcher == nil {
rt.inboundDispatcher = newInboundDispatcher()
}
if rt.suppressGoodByeOnStop == nil {
rt.suppressGoodByeOnStop = &atomic.Bool{}
}
if rt.transport == nil && rt.conn != nil {
rt.transport = newTransportBinding(rt.conn, rt.queue)
}
normalizeClientSessionRuntimeTransportState(rt)
ensureClientSessionRuntimeTransportLifecycle(rt)
return rt
}
func (c *ClientCommon) setClientSessionRuntime(rt *clientSessionRuntime) {
c.setClientSessionRuntimeWithCloseOld(rt, false)
}
func (c *ClientCommon) setClientSessionRuntimeWithCloseOld(rt *clientSessionRuntime, closeOld bool) {
if c == nil || rt == nil {
return
}
var oldBinding *transportBinding
if prev := c.clientSessionRuntimeSnapshot(); prev != nil && prev.transport != nil && prev.transport != rt.transport {
oldBinding = prev.transport
}
rt = prepareClientSessionRuntime(rt)
c.sessionRuntime.Store(rt)
c.stopCtx = rt.stopCtx
c.stopFn = rt.stopFn
if rt.transport != nil {
c.queue = rt.transport.queueSnapshot()
c.conn = rt.transport.connSnapshot()
} else {
c.queue = rt.queue
c.conn = rt.conn
}
if oldBinding != nil {
stopReplacedTransportBinding(oldBinding, rt.transport, closeOld)
}
}
func (c *ClientCommon) resetClientSessionRuntimeBase() {
if c == nil {
return
}
stopCtx, stopFn := context.WithCancel(context.Background())
c.sessionRuntime.Store(newClientSessionRuntimeBase(stopCtx, stopFn))
c.conn = nil
c.queue = nil
c.stopCtx = stopCtx
c.stopFn = stopFn
}
func (c *ClientCommon) cleanupFailedClientStart() {
if c == nil {
return
}
rt := c.clientSessionRuntimeSnapshot()
if rt != nil && rt.stopFn != nil {
rt.stopFn()
}
c.cleanupClientSessionResources()
c.rollbackClientSessionStart()
c.resetClientSessionRuntimeBase()
}
func newClientSessionRuntime(conn net.Conn, stopCtx context.Context, stopFn context.CancelFunc, queue *stario.StarQueue, epoch uint64) *clientSessionRuntime {
return prepareClientSessionRuntime(&clientSessionRuntime{
transport: newTransportBinding(conn, queue),
transportAttached: conn != nil,
conn: conn,
stopCtx: stopCtx,
stopFn: stopFn,
queue: queue,
inboundDispatcher: newInboundDispatcher(),
epoch: epoch,
suppressGoodByeOnStop: &atomic.Bool{},
})
}
func (rt *clientSessionRuntime) runtimeShouldSuppressGoodByeOnStop() bool {
if rt == nil || rt.suppressGoodByeOnStop == nil {
return false
}
return rt.suppressGoodByeOnStop.Load()
}
func (rt *clientSessionRuntime) markRuntimeSuppressGoodByeOnStop() {
if rt == nil || rt.suppressGoodByeOnStop == nil {
return
}
rt.suppressGoodByeOnStop.Store(true)
}
func (c *ClientCommon) retireClientSessionRuntime(rt *clientSessionRuntime, suppressGoodBye bool) {
if c == nil || rt == nil {
return
}
if suppressGoodBye {
rt.markRuntimeSuppressGoodByeOnStop()
}
if rt.transportStopFn != nil {
rt.transportStopFn()
}
}
func (c *ClientCommon) clearClientSessionRuntimeTransport() {
if c == nil {
return
}
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return
}
if rt.transportStopFn != nil {
rt.transportStopFn()
}
next := *rt
next.transport = nil
next.transportAttached = false
next.conn = nil
next.transportStopCtx = nil
next.transportStopFn = nil
c.setClientSessionRuntimeWithCloseOld(&next, true)
}
func (c *ClientCommon) clearClientSessionRuntimeQueue() {
if c == nil {
return
}
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return
}
next := *rt
next.queue = nil
if next.transport != nil {
next.transport = newTransportBinding(next.transport.connSnapshot(), nil)
}
c.setClientSessionRuntime(&next)
}
func (c *ClientCommon) attachClientSessionTransport(conn net.Conn) error {
if c == nil {
return errors.New("client is nil")
}
if conn == nil {
return errors.New("conn is nil")
}
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return errors.New("client session runtime is nil")
}
if rt.queue == nil {
return errClientSessionQueueUnavailable
}
oldBinding := rt.transport
if rt.transportStopFn != nil {
rt.transportStopFn()
}
oldRoute := clientSessionRouteFromRuntime(rt)
if streamRuntime := c.getStreamRuntime(); streamRuntime != nil {
streamRuntime.closeClientRoute(oldRoute, errTransportDetached)
}
if bulkRuntime := c.getBulkRuntime(); bulkRuntime != nil {
bulkRuntime.closeClientRoute(oldRoute, errTransportDetached)
}
// A sidecar is physically tied to the old primary transport. Retire it
// before publishing the replacement route so a new bulk cannot reuse it.
c.closeClientDedicatedSidecarWithError(errTransportDetached)
next := *rt
next.transport = newTransportBinding(conn, rt.queue)
next.transportAttached = true
next.conn = conn
next.transportStopCtx = nil
next.transportStopFn = nil
next.suppressGoodByeOnStop = &atomic.Bool{}
c.setClientSessionRuntimeWithCloseOld(&next, true)
if oldConn := oldBinding.connSnapshot(); oldConn != nil && oldConn != conn {
_ = oldConn.Close()
}
return c.startClientTransportRuntime(c.clientSessionRuntimeSnapshot())
}
func (c *ClientCommon) clientSessionRuntimeSnapshot() *clientSessionRuntime {
if c == nil {
return nil
}
return c.sessionRuntime.Load()
}
func normalizeClientSessionRuntimeTransportState(rt *clientSessionRuntime) {
if rt == nil {
return
}
if rt.transport != nil {
rt.transportAttached = rt.transport.connSnapshot() != nil
return
}
rt.transportAttached = rt.conn != nil
}
func ensureClientSessionRuntimeTransportLifecycle(rt *clientSessionRuntime) {
if rt == nil {
return
}
if rt.conn == nil {
rt.transportStopCtx = nil
rt.transportStopFn = nil
return
}
if rt.transportStopCtx != nil && rt.transportStopFn != nil {
return
}
parent := rt.stopCtx
if parent == nil {
parent = context.Background()
}
rt.transportStopCtx, rt.transportStopFn = context.WithCancel(parent)
}
func (c *ClientCommon) clientTransportConnSnapshot() net.Conn {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
if rt.transport != nil {
return rt.transport.connSnapshot()
}
return rt.conn
}
func (c *ClientCommon) clientInboundDispatcherSnapshot() *inboundDispatcher {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
return rt.inboundDispatcher
}
func (c *ClientCommon) clientStopContextSnapshot() context.Context {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
return rt.stopCtx
}
func (c *ClientCommon) clientStopFuncSnapshot() context.CancelFunc {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
return rt.stopFn
}
func (c *ClientCommon) clientQueueSnapshot() *stario.StarQueue {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
if rt.transport != nil {
return rt.transport.queueSnapshot()
}
return rt.queue
}
func (c *ClientCommon) clientTransportBindingSnapshot() *transportBinding {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
if rt.transport != nil {
return rt.transport
}
if rt.conn == nil {
return nil
}
return newTransportBinding(rt.conn, rt.queue)
}
func (c *ClientCommon) clientTransportStopContextSnapshot() context.Context {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return nil
}
if rt.transportStopCtx != nil {
return rt.transportStopCtx
}
return rt.stopCtx
}
func (c *ClientCommon) clientTransportAttachedSnapshot() bool {
rt := c.clientSessionRuntimeSnapshot()
if rt == nil {
return false
}
return rt.transportAttached
}