Files

145 lines
3.0 KiB
Go

package nixmsg
import (
"context"
"encoding/json"
"fmt"
"sync"
"sync/atomic"
"time"
"github.com/google/uuid"
)
type pendingReq struct {
rid string
ch chan respFrame
}
type sendItem struct {
frame map[string]any
payload []byte
id string
sendAtMs *int64
result chan sendOutcome
inflight bool
}
type sendOutcome struct {
res SendResult
err error
}
type dedupState int
const (
dedupDelivered dedupState = iota + 1
dedupAcked
)
type dedupEntry struct {
state dedupState
from string
id string
}
// Client NixMsg 客户端。
type Client struct {
opts Options
endpointID string
url string
credKind string
transport transport
backoff *reconnectBackoff
mu sync.Mutex
state ConnectionState
handshook bool
limits HandshakeLimits
clockSkew int64
session string
stopReconnect bool
closed bool
ridSeq atomic.Uint64
pending map[string]*pendingReq
sendQ []*sendItem
inflight int
dedup map[string]*dedupEntry
dedupOrd []string
receiptSeen map[string]struct{}
cbMu sync.Mutex
onSession func(token string)
onMessage func(msg Message) error
onReceipt func(r Receipt)
onRevoked func(e RevokedEvent)
onPresence func(e PresenceEvent)
onGroupEvent func(e GroupEvent)
onConnection func(ev ConnectionEvent)
helloSentAt time.Time
ctx context.Context
cancel context.CancelFunc
// downCh 串行处理非 resp 下行,避免在 MQTT 收包回调里同步 request 死锁。
downCh chan []byte
}
// New 创建客户端(尚未连接)。
func New() *Client {
return &Client{
pending: make(map[string]*pendingReq),
dedup: make(map[string]*dedupEntry),
receiptSeen: make(map[string]struct{}),
state: StateOffline,
backoff: newReconnectBackoff(),
downCh: make(chan []byte, 256),
}
}
// OnSession 会话令牌回调。
func (c *Client) OnSession(h func(token string)) { c.onSession = h }
// OnMessage 消息回调。自动模式下返回 error 则不发 ack 并删去重记录。
func (c *Client) OnMessage(h func(msg Message) error) { c.onMessage = h }
// OnReceipt 回执回调。
func (c *Client) OnReceipt(h func(r Receipt)) { c.onReceipt = h }
// OnRevoked 撤回/作废回调。
func (c *Client) OnRevoked(h func(e RevokedEvent)) { c.onRevoked = h }
// OnPresence 上下线回调。
func (c *Client) OnPresence(h func(e PresenceEvent)) { c.onPresence = h }
// OnGroupEvent 群事件回调。
func (c *Client) OnGroupEvent(h func(e GroupEvent)) { c.onGroupEvent = h }
// OnConnection 连接状态回调。
func (c *Client) OnConnection(h func(ev ConnectionEvent)) { c.onConnection = h }
func newMessageID() string {
id, err := uuid.NewV7()
if err != nil {
return uuid.NewString()
}
return id.String()
}
func (c *Client) nextRID() string {
return fmt.Sprintf("%d", c.ridSeq.Add(1))
}
type respFrame struct {
OK bool `json:"ok"`
Data json.RawMessage `json:"data"`
Error *struct {
Code string `json:"code"`
Message string `json:"message"`
} `json:"error"`
}