145 lines
3.0 KiB
Go
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"`
|
|
}
|