0826e17063
- 为控制消息增加优先级、公平调度、队列字节预算和自适应批处理 - 支持可取消的写门等待,收紧 shared/dedicated bulk、stream 和 Reply 写入边界 - 修复 bulk reset/close、连接 handoff 和安全 profile 切换时序 - 保留旧取消与超时错误契约,新增阶段化 TransportSendError - 增加 ReplyCtx、ReplyObjCtx、写超时配置及黑洞连接和竞态回归测试
312 lines
7.9 KiB
Go
312 lines
7.9 KiB
Go
package notify
|
|
|
|
import (
|
|
cryptorand "crypto/rand"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
const (
|
|
systemPeerAttachKey = "_notify_peer_attach"
|
|
peerAttachTimeout = 5 * time.Second
|
|
peerAttachTransitionFallbackTTL = peerAttachTimeout
|
|
)
|
|
|
|
type peerAttachRequest struct {
|
|
PeerID string
|
|
Features uint64
|
|
ClientNonce []byte
|
|
ClientECDHEPublicKey []byte
|
|
AuthTag []byte
|
|
}
|
|
|
|
type peerAttachResponse struct {
|
|
PeerID string
|
|
Accepted bool
|
|
Reused bool
|
|
Error string
|
|
Features uint64
|
|
KeyMode string
|
|
ServerNonce []byte
|
|
ServerECDHEPublicKey []byte
|
|
AuthTag []byte
|
|
}
|
|
|
|
func newClientPeerIdentity() string {
|
|
var buf [16]byte
|
|
if _, err := cryptorand.Read(buf[:]); err == nil {
|
|
return "peer-" + hex.EncodeToString(buf[:])
|
|
}
|
|
return fmt.Sprintf("peer-%d", time.Now().UnixNano())
|
|
}
|
|
|
|
func (c *ClientCommon) ensureClientPeerIdentity() string {
|
|
if c == nil {
|
|
return ""
|
|
}
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
if strings.TrimSpace(c.peerIdentity) == "" {
|
|
c.peerIdentity = newClientPeerIdentity()
|
|
}
|
|
return c.peerIdentity
|
|
}
|
|
|
|
func (c *ClientCommon) setClientPeerIdentity(peerID string) {
|
|
if c == nil {
|
|
return
|
|
}
|
|
peerID = strings.TrimSpace(peerID)
|
|
if peerID == "" {
|
|
return
|
|
}
|
|
c.mu.Lock()
|
|
c.peerIdentity = peerID
|
|
c.mu.Unlock()
|
|
}
|
|
|
|
func decodePeerAttachRequest(decodeFn func([]byte) (interface{}, error), data []byte) (peerAttachRequest, error) {
|
|
if decodeFn == nil {
|
|
decodeFn = Decode
|
|
}
|
|
value, err := decodeFn(data)
|
|
if err != nil {
|
|
return peerAttachRequest{}, err
|
|
}
|
|
switch req := value.(type) {
|
|
case peerAttachRequest:
|
|
return req, nil
|
|
case *peerAttachRequest:
|
|
if req == nil {
|
|
return peerAttachRequest{}, errors.New("peer attach request is nil")
|
|
}
|
|
return *req, nil
|
|
default:
|
|
return peerAttachRequest{}, fmt.Errorf("unexpected peer attach request type %T", value)
|
|
}
|
|
}
|
|
|
|
func decodePeerAttachResponse(decodeFn func([]byte) (interface{}, error), data []byte) (peerAttachResponse, error) {
|
|
if decodeFn == nil {
|
|
decodeFn = Decode
|
|
}
|
|
value, err := decodeFn(data)
|
|
if err != nil {
|
|
return peerAttachResponse{}, err
|
|
}
|
|
switch resp := value.(type) {
|
|
case peerAttachResponse:
|
|
return resp, nil
|
|
case *peerAttachResponse:
|
|
if resp == nil {
|
|
return peerAttachResponse{}, errors.New("peer attach response is nil")
|
|
}
|
|
return *resp, nil
|
|
default:
|
|
return peerAttachResponse{}, fmt.Errorf("unexpected peer attach response type %T", value)
|
|
}
|
|
}
|
|
|
|
func (c *ClientCommon) announceClientPeerIdentity() error {
|
|
if c == nil {
|
|
return errors.New("client is nil")
|
|
}
|
|
peerID := c.ensureClientPeerIdentity()
|
|
if peerID == "" {
|
|
return errors.New("peer identity is empty")
|
|
}
|
|
req, requestState, err := c.buildPeerAttachRequest(peerID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
encoded, err := c.sequenceEn(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
reply, err := c.sendWait(TransferMsg{
|
|
Key: systemPeerAttachKey,
|
|
Value: encoded,
|
|
Type: MSG_SYS_WAIT,
|
|
}, peerAttachTimeout)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
resp, err := decodePeerAttachResponse(c.sequenceDe, reply.Value)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if resp.PeerID != "" {
|
|
c.setClientPeerIdentity(resp.PeerID)
|
|
}
|
|
if !resp.Accepted {
|
|
if strings.TrimSpace(resp.Error) != "" {
|
|
return errors.New(resp.Error)
|
|
}
|
|
return errors.New("peer attach rejected")
|
|
}
|
|
verifyResult, err := c.verifyPeerAttachResponse(req, resp, requestState)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
c.setClientNegotiatedSteadyTransportProtection(verifyResult.steadyProfile)
|
|
c.markClientPeerAttachAuthenticated(verifyResult.authFallback, time.Now())
|
|
return nil
|
|
}
|
|
|
|
func (s *ServerCommon) bindAcceptedClientIdentity(current *LogicalConn, peerID string) (*LogicalConn, bool, error) {
|
|
if s == nil {
|
|
return nil, false, errors.New("server is nil")
|
|
}
|
|
if current == nil {
|
|
return nil, false, errors.New("client is nil")
|
|
}
|
|
peerID = strings.TrimSpace(peerID)
|
|
if peerID == "" {
|
|
return nil, false, errors.New("peer id is empty")
|
|
}
|
|
if current.ID() == peerID {
|
|
current.markIdentityBound()
|
|
return current, false, nil
|
|
}
|
|
existing := s.GetLogicalConn(peerID)
|
|
if existing == nil {
|
|
if err := s.renameAcceptedLogical(current, peerID); err != nil {
|
|
return nil, false, err
|
|
}
|
|
current.markIdentityBound()
|
|
return current, false, nil
|
|
}
|
|
if existing == current {
|
|
existing.markIdentityBound()
|
|
return existing, false, nil
|
|
}
|
|
if err := s.handoffAcceptedLogicalTransport(existing, current); err != nil {
|
|
return nil, true, err
|
|
}
|
|
existing.markIdentityBound()
|
|
return existing, true, nil
|
|
}
|
|
|
|
func (s *ServerCommon) replyPeerAttach(client *LogicalConn, message Message, resp peerAttachResponse) error {
|
|
if s == nil {
|
|
return errors.New("server is nil")
|
|
}
|
|
if client == nil {
|
|
return errors.New("client is nil")
|
|
}
|
|
encoded, err := s.sequenceEn(resp)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
reply := TransferMsg{
|
|
ID: message.ID,
|
|
Key: systemPeerAttachKey,
|
|
Value: encoded,
|
|
Type: MSG_SYS_REPLY,
|
|
}
|
|
transport := messageTransportConnSnapshot(&message)
|
|
profile := messageInboundTransportProtectionSnapshot(&message)
|
|
if message.inboundConn != nil {
|
|
return s.sendTransferInbound(client, transport, message.inboundConn, profile, reply)
|
|
}
|
|
env, err := wrapTransferMsgEnvelope(reply, s.sequenceEn)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
env.transportProfile = profile
|
|
return s.sendSignalEnvelopeMaybeReliableTransport(transport, env, reply)
|
|
}
|
|
|
|
func (s *ServerCommon) handlePeerAttachSystemMessage(message Message) bool {
|
|
if message.Key != systemPeerAttachKey {
|
|
return false
|
|
}
|
|
message = hydrateServerMessagePeerFields(message)
|
|
current := messageLogicalConnSnapshot(&message)
|
|
if message.inboundTransportProfile == nil && current != nil {
|
|
profile := current.transportProtectionProfileSnapshot()
|
|
message.inboundTransportProfile = &profile
|
|
}
|
|
transport := message.inboundConn
|
|
if transport == nil && current != nil {
|
|
transport = current.transportSnapshot()
|
|
}
|
|
req, err := decodePeerAttachRequest(s.sequenceDe, message.Value)
|
|
if err != nil {
|
|
if current != nil {
|
|
_ = s.replyPeerAttach(current, message, peerAttachResponse{
|
|
Accepted: false,
|
|
Error: err.Error(),
|
|
})
|
|
}
|
|
return true
|
|
}
|
|
auth, err := s.validatePeerAttachRequestAuth(current, transport, req)
|
|
if err != nil {
|
|
classifyPeerAttachRejectCounter(s, err)
|
|
if current != nil {
|
|
_ = s.replyPeerAttach(current, message, peerAttachResponse{
|
|
PeerID: req.PeerID,
|
|
Accepted: false,
|
|
Error: err.Error(),
|
|
})
|
|
}
|
|
return true
|
|
}
|
|
bound, reused, err := s.bindAcceptedClientIdentity(current, req.PeerID)
|
|
if err != nil {
|
|
if current != nil {
|
|
_ = s.replyPeerAttach(current, message, peerAttachResponse{
|
|
PeerID: req.PeerID,
|
|
Accepted: false,
|
|
Error: err.Error(),
|
|
})
|
|
}
|
|
return true
|
|
}
|
|
resp := peerAttachResponse{
|
|
PeerID: bound.ID(),
|
|
Accepted: true,
|
|
Reused: reused,
|
|
}
|
|
steadyProfile, err := s.preparePeerAttachSteadyTransportProfile(bound, req, &resp, auth)
|
|
if err != nil {
|
|
if bound != nil {
|
|
_ = s.replyPeerAttach(bound, message, peerAttachResponse{
|
|
PeerID: req.PeerID,
|
|
Accepted: false,
|
|
Error: err.Error(),
|
|
})
|
|
}
|
|
return true
|
|
}
|
|
s.signPeerAttachResponse(bound, req, &resp, auth)
|
|
if bound != nil {
|
|
bound.markPeerAttachAuthenticated(s.securityAuthMode, auth.fallback, time.Now())
|
|
if auth.explicit {
|
|
s.peerAttachExplicitCount.Add(1)
|
|
} else if auth.fallback {
|
|
s.peerAttachAuthFallbackCount.Add(1)
|
|
}
|
|
}
|
|
var transitionProfile *transportProtectionProfile
|
|
if bound != nil && s.securityConfigured {
|
|
if message.inboundTransportProfile != nil {
|
|
transitionProfile = bound.installInboundTransitionProfile(*message.inboundTransportProfile)
|
|
}
|
|
bound.applyTransportProtectionProfile(steadyProfile)
|
|
}
|
|
replyErr := s.replyPeerAttach(bound, message, resp)
|
|
if transitionProfile != nil {
|
|
bound.clearInboundTransitionProfile(transitionProfile)
|
|
}
|
|
if replyErr != nil && bound != nil {
|
|
s.stopLogicalSession(bound, "peer attach reply failed", replyErr)
|
|
return true
|
|
}
|
|
return true
|
|
}
|