为什么:并入宿主项目后须与宿主日志体系一致——宿主统一用 pkg/logger 初始化 logrus 全局实例并直接 logrus.Xxx 包级调用;包一层 Logger 接口会让 logrus 的 caller 定位(报告调用文件:行号)失准。 改动: - 删除 logger 包(Logger 接口 / Default / SetDefault,默认本就是 logrus.StandardLogger) - 11 个结构体移除 lgr logger.Logger 字段、SetLogger 方法及构造中的 lgr 初始化 - 192 处 lgr.Xxxf 调用改为 logrus.Xxxf,日志消息文本保持不变 - xiaomi.Client 移除 Logger() 访问器,air_conditioner 回读日志去掉判空调用直连 logrus - miot_i18n.go 两处标准库 log.Printf 改为 logrus.Errorf,统一日志出口 - go.mod:logrus 从 indirect 提升为直接依赖 - goimports/goformat 全树规范化:此前部分文件未 gofmt(单行 if、对齐),本次顺带 格式化,纯空白/换行变化,无逻辑改动 验证:go build / go vet / go test ./... 全部通过(bridge、miot、xiaomi、devices、specs)
1093 lines
29 KiB
Go
1093 lines
29 KiB
Go
// Package miot provides MIoT core client for Xiaomi Home devices.
|
|
// mips_client.go — full MIPS (MIoT Pub/Sub) client ported from py-miot/miot_mips.py.
|
|
package miot
|
|
|
|
import (
|
|
"encoding/binary"
|
|
"encoding/json"
|
|
"fmt"
|
|
"math"
|
|
"math/rand"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
mqtt "github.com/eclipse/paho.mqtt.golang"
|
|
"github.com/sirupsen/logrus"
|
|
)
|
|
|
|
// ============================================================================
|
|
// MIPS protocol constants (aligned with Python _MipsClient)
|
|
// ============================================================================
|
|
|
|
const (
|
|
mipsQoS = 2
|
|
mipsReconnectIntervalMin = 10.0 // seconds
|
|
mipsReconnectIntervalMax = 600.0 // seconds
|
|
mipsSubInterval = 1 // second
|
|
mipsMQTTInterval = 1 // second
|
|
uint32max = 0xFFFFFFFF
|
|
MipsRequestTimeoutDefault = 10000 // milliseconds
|
|
)
|
|
|
|
// ============================================================================
|
|
// MIPS message type options (aligned with Python _MipsMsgTypeOptions)
|
|
// ============================================================================
|
|
|
|
const (
|
|
mipsMsgTypeID = 0
|
|
mipsMsgTypeRetTopic = 1
|
|
mipsMsgTypePayload = 2
|
|
mipsMsgTypeFrom = 3
|
|
)
|
|
|
|
// ============================================================================
|
|
// MipsMessage — binary message packing/unpacking
|
|
// ============================================================================
|
|
|
|
// MipsMessage represents a packed MIPS protocol message.
|
|
// Protocol: each segment is [4-byte length LE][1-byte type][len bytes + 1 null]
|
|
type MipsMessage struct {
|
|
Mid int
|
|
MsgFrom string
|
|
RetTopic string
|
|
Payload string
|
|
}
|
|
|
|
// Pack packs a MIPS message into binary format.
|
|
// Format: <I (4 bytes len) B (1 byte type)> + data + null terminator
|
|
func PackMipsMessage(mid int, payload string, msgFrom, retTopic string) ([]byte, error) {
|
|
if payload == "" {
|
|
return nil, &MIoTMipsError{MIoTError{Code: CodeInvalidParams, Message: "invalid mid or payload"}}
|
|
}
|
|
|
|
var buf []byte
|
|
|
|
// mid segment
|
|
midBuf := make([]byte, 4)
|
|
binary.LittleEndian.PutUint32(midBuf, uint32(mid))
|
|
buf = append(buf, 0x04, 0x00, 0x00, 0x00) // length = 4
|
|
buf = append(buf, byte(mipsMsgTypeID))
|
|
buf = append(buf, midBuf...)
|
|
|
|
// msg_from segment (optional)
|
|
if msgFrom != "" {
|
|
data := []byte(msgFrom)
|
|
seg := buildMipsSegment(mipsMsgTypeFrom, data)
|
|
buf = append(buf, seg...)
|
|
}
|
|
|
|
// ret_topic segment (optional)
|
|
if retTopic != "" {
|
|
data := []byte(retTopic)
|
|
seg := buildMipsSegment(mipsMsgTypeRetTopic, data)
|
|
buf = append(buf, seg...)
|
|
}
|
|
|
|
// payload segment
|
|
data := []byte(payload)
|
|
seg := buildMipsSegment(mipsMsgTypePayload, data)
|
|
buf = append(buf, seg...)
|
|
|
|
return buf, nil
|
|
}
|
|
|
|
// buildMipsSegment builds a single protocol segment.
|
|
func buildMipsSegment(msgType int, data []byte) []byte {
|
|
totalLen := len(data) + 1 // +1 for null terminator
|
|
lenBuf := make([]byte, 4)
|
|
binary.LittleEndian.PutUint32(lenBuf, uint32(totalLen))
|
|
|
|
var seg []byte
|
|
seg = append(seg, lenBuf...)
|
|
seg = append(seg, byte(msgType))
|
|
seg = append(seg, data...)
|
|
seg = append(seg, 0) // null terminator
|
|
return seg
|
|
}
|
|
|
|
// UnpackMipsMessage unpacks binary data into a MipsMessage.
|
|
func UnpackMipsMessage(data []byte) *MipsMessage {
|
|
msg := &MipsMessage{}
|
|
dataLen := len(data)
|
|
pos := 0
|
|
|
|
for pos < dataLen {
|
|
if pos+5 > dataLen {
|
|
break
|
|
}
|
|
segLen := int(binary.LittleEndian.Uint32(data[pos : pos+4]))
|
|
segType := int(data[pos+4])
|
|
pos += 5
|
|
|
|
if segLen == 0 || pos+segLen > dataLen {
|
|
break
|
|
}
|
|
segData := data[pos : pos+segLen]
|
|
// Strip trailing null
|
|
segData = bytesTrimNull(segData)
|
|
|
|
switch segType {
|
|
case mipsMsgTypeID:
|
|
if len(segData) >= 4 {
|
|
msg.Mid = int(binary.LittleEndian.Uint32(segData))
|
|
}
|
|
case mipsMsgTypeRetTopic:
|
|
msg.RetTopic = string(segData)
|
|
case mipsMsgTypePayload:
|
|
msg.Payload = string(segData)
|
|
case mipsMsgTypeFrom:
|
|
msg.MsgFrom = string(segData)
|
|
}
|
|
pos += segLen
|
|
}
|
|
|
|
return msg
|
|
}
|
|
|
|
func bytesTrimNull(b []byte) []byte {
|
|
for len(b) > 0 && b[len(b)-1] == 0 {
|
|
b = b[:len(b)-1]
|
|
}
|
|
return b
|
|
}
|
|
|
|
// String returns a debug representation of the message.
|
|
func (m *MipsMessage) String() string {
|
|
return fmt.Sprintf("%d, %s, %s, %s", m.Mid, m.MsgFrom, m.RetTopic, m.Payload)
|
|
}
|
|
|
|
// ============================================================================
|
|
// MIoTDeviceState — device state enum (aligned with Python MIoTDeviceState)
|
|
// ============================================================================
|
|
|
|
// MIoTDeviceState represents the online/offline state of a device.
|
|
type MIoTDeviceState int
|
|
|
|
const (
|
|
// DeviceStateDisable means the device is disabled.
|
|
DeviceStateDisable MIoTDeviceState = 0
|
|
// DeviceStateOffline means the device is offline.
|
|
DeviceStateOffline MIoTDeviceState = 1
|
|
// DeviceStateOnline means the device is online.
|
|
DeviceStateOnline MIoTDeviceState = 2
|
|
)
|
|
|
|
// String returns the string representation of MIoTDeviceState.
|
|
func (s MIoTDeviceState) String() string {
|
|
switch s {
|
|
case DeviceStateDisable:
|
|
return "disable"
|
|
case DeviceStateOffline:
|
|
return "offline"
|
|
case DeviceStateOnline:
|
|
return "online"
|
|
default:
|
|
return "unknown"
|
|
}
|
|
}
|
|
|
|
// ============================================================================
|
|
// MipsRequest — pending request tracking
|
|
// ============================================================================
|
|
|
|
// MipsRequest tracks a pending MIPS request with timeout.
|
|
type MipsRequest struct {
|
|
Mid int
|
|
OnReply func(payload string, ctx interface{})
|
|
OnReplyCtx interface{}
|
|
Timer *time.Timer
|
|
}
|
|
|
|
// ============================================================================
|
|
// MipsBroadcast — subscription tracking
|
|
// ============================================================================
|
|
|
|
// MipsBroadcast represents a subscribed broadcast topic.
|
|
type MipsBroadcast struct {
|
|
Topic string
|
|
Handler func(topic string, payload string, ctx interface{})
|
|
HandlerCtx interface{}
|
|
}
|
|
|
|
// ============================================================================
|
|
// mipsClient — core MIPS client (internal)
|
|
// ============================================================================
|
|
|
|
// mipsClient is the concrete MIPS client using paho MQTT.
|
|
type mipsClient struct {
|
|
mu sync.Mutex
|
|
clientID string
|
|
host string
|
|
port int
|
|
username string
|
|
password string
|
|
caFile string
|
|
certFile string
|
|
keyFile string
|
|
|
|
// MQTT connection
|
|
mqttConn MQTTClient // interface to paho MQTT
|
|
mqttConnected bool
|
|
|
|
// Reconnect
|
|
reconnectActive bool
|
|
reconnectInterval float64
|
|
reconnectTimer *time.Timer
|
|
|
|
// Subscriptions
|
|
pendingSubs map[string]int // topic -> retry count
|
|
pendingSubMu sync.Mutex
|
|
subMu sync.Mutex
|
|
subscriptions map[string]*MipsBroadcast
|
|
|
|
// Requests
|
|
requestMu sync.Mutex
|
|
requests map[string]*MipsRequest // mid as string -> request
|
|
seedID int
|
|
|
|
// State subscribers
|
|
stateSubMu sync.Mutex
|
|
stateSubs map[string]func(key string, connected bool)
|
|
|
|
// Events
|
|
onConnect func()
|
|
onDisconnect func(err error)
|
|
|
|
// Shutdown
|
|
stopCh chan struct{}
|
|
doneCh chan struct{}
|
|
}
|
|
|
|
// MQTTClient is a small interface for MQTT operations.
|
|
// This allows decoupling from the concrete paho MQTT library.
|
|
type MQTTClient interface {
|
|
Connect() error
|
|
Disconnect()
|
|
IsConnected() bool
|
|
Subscribe(topic string, qos byte, handler func(topic string, payload []byte)) error
|
|
SubscribeMany(topics map[string]byte) error
|
|
Unsubscribe(topic string) error
|
|
Publish(topic string, qos byte, payload []byte) error
|
|
SetCredentials(username, password string)
|
|
SetTLS(caFile, certFile, keyFile string) error
|
|
}
|
|
|
|
// newMipsClient creates a new mipsClient (internal).
|
|
func newMipsClient(clientID, host string, port int, username, password, caFile, certFile, keyFile string) *mipsClient {
|
|
c := &mipsClient{
|
|
clientID: clientID,
|
|
host: host,
|
|
port: port,
|
|
username: username,
|
|
password: password,
|
|
caFile: caFile,
|
|
certFile: certFile,
|
|
keyFile: keyFile,
|
|
seedID: rand.Intn(uint32max),
|
|
pendingSubs: make(map[string]int),
|
|
subscriptions: make(map[string]*MipsBroadcast),
|
|
requests: make(map[string]*MipsRequest),
|
|
stateSubs: make(map[string]func(key string, connected bool)),
|
|
reconnectActive: true,
|
|
stopCh: make(chan struct{}),
|
|
doneCh: make(chan struct{}),
|
|
}
|
|
return c
|
|
}
|
|
|
|
// connect starts the MQTT connection.
|
|
func (c *mipsClient) connect(mqttClient MQTTClient) error {
|
|
c.mu.Lock()
|
|
c.mqttConn = mqttClient
|
|
c.mu.Unlock()
|
|
|
|
if err := c.mqttConn.Connect(); err != nil {
|
|
return fmt.Errorf("mqtt connect: %w", err)
|
|
}
|
|
|
|
c.mu.Lock()
|
|
c.mqttConnected = true
|
|
c.reconnectActive = true
|
|
c.reconnectInterval = 0
|
|
c.mu.Unlock()
|
|
|
|
// Notify connect callbacks
|
|
if c.onConnect != nil {
|
|
c.onConnect()
|
|
}
|
|
|
|
// Notify state subscribers
|
|
c.notifyStateSubs(true)
|
|
|
|
return nil
|
|
}
|
|
|
|
// disconnect stops the MQTT connection.
|
|
func (c *mipsClient) disconnect() {
|
|
c.mu.Lock()
|
|
c.reconnectActive = false
|
|
if c.reconnectTimer != nil {
|
|
c.reconnectTimer.Stop()
|
|
c.reconnectTimer = nil
|
|
}
|
|
c.mu.Unlock()
|
|
|
|
if c.mqttConn != nil {
|
|
c.mqttConn.Disconnect()
|
|
c.mu.Lock()
|
|
c.mqttConnected = false
|
|
c.mu.Unlock()
|
|
|
|
// Notify disconnect
|
|
if c.onDisconnect != nil {
|
|
c.onDisconnect(nil)
|
|
}
|
|
c.notifyStateSubs(false)
|
|
}
|
|
}
|
|
|
|
// isConnected returns the current connection state.
|
|
func (c *mipsClient) isConnected() bool {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
if c.mqttConn == nil {
|
|
return false
|
|
}
|
|
return c.mqttConn.IsConnected() && c.mqttConnected
|
|
}
|
|
|
|
// subscribe adds a subscription to a topic.
|
|
func (c *mipsClient) subscribe(topic string, handler func(topic string, payload string, ctx interface{}), ctx interface{}) error {
|
|
c.subMu.Lock()
|
|
defer c.subMu.Unlock()
|
|
|
|
if c.subscriptions[topic] != nil {
|
|
logrus.Debugf("[mips] re-register subscription: %s", topic)
|
|
}
|
|
c.subscriptions[topic] = &MipsBroadcast{
|
|
Topic: topic,
|
|
Handler: handler,
|
|
HandlerCtx: ctx,
|
|
}
|
|
|
|
// Add to pending subs for actual MQTT subscribe
|
|
c.pendingSubMu.Lock()
|
|
c.pendingSubs[topic] = 0
|
|
c.pendingSubMu.Unlock()
|
|
|
|
go c.flushPendingSubs()
|
|
return nil
|
|
}
|
|
|
|
// unsubscribe removes a subscription.
|
|
func (c *mipsClient) unsubscribe(topic string) error {
|
|
c.subMu.Lock()
|
|
delete(c.subscriptions, topic)
|
|
c.subMu.Unlock()
|
|
|
|
if c.mqttConn != nil && c.mqttConn.IsConnected() {
|
|
if err := c.mqttConn.Unsubscribe(topic); err != nil {
|
|
logrus.Debugf("[mips] unsubscribe error %s: %v", topic, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// publish sends a message on a topic.
|
|
func (c *mipsClient) publish(topic string, payload []byte) error {
|
|
if c.mqttConn == nil || !c.mqttConn.IsConnected() {
|
|
return fmt.Errorf("mips publish when not connected: %s", topic)
|
|
}
|
|
return c.mqttConn.Publish(topic, byte(mipsQoS), payload)
|
|
}
|
|
|
|
// subState subscribes to connection state changes.
|
|
func (c *mipsClient) subState(key string, handler func(key string, connected bool)) {
|
|
c.stateSubMu.Lock()
|
|
c.stateSubs[key] = handler
|
|
c.stateSubMu.Unlock()
|
|
}
|
|
|
|
// unsubState unsubscribes from connection state changes.
|
|
func (c *mipsClient) unsubState(key string) {
|
|
c.stateSubMu.Lock()
|
|
delete(c.stateSubs, key)
|
|
c.stateSubMu.Unlock()
|
|
}
|
|
|
|
// notifyStateSubs notifies all state subscribers.
|
|
func (c *mipsClient) notifyStateSubs(connected bool) {
|
|
c.stateSubMu.Lock()
|
|
subs := make(map[string]func(key string, connected bool), len(c.stateSubs))
|
|
for k, v := range c.stateSubs {
|
|
subs[k] = v
|
|
}
|
|
c.stateSubMu.Unlock()
|
|
|
|
for key, handler := range subs {
|
|
if handler != nil {
|
|
handler(key, connected)
|
|
}
|
|
}
|
|
}
|
|
|
|
// onMQTTConnect is called when the MQTT client connects.
|
|
func (c *mipsClient) onMQTTConnect() {
|
|
c.mu.Lock()
|
|
c.mqttConnected = true
|
|
c.reconnectInterval = 0 // reset backoff
|
|
c.mu.Unlock()
|
|
|
|
// Resubscribe all topics
|
|
c.subMu.Lock()
|
|
topics := make([]string, 0, len(c.subscriptions))
|
|
for topic := range c.subscriptions {
|
|
topics = append(topics, topic)
|
|
}
|
|
c.subMu.Unlock()
|
|
|
|
for _, topic := range topics {
|
|
c.pendingSubMu.Lock()
|
|
c.pendingSubs[topic] = 0
|
|
c.pendingSubMu.Unlock()
|
|
}
|
|
go c.flushPendingSubs()
|
|
}
|
|
|
|
// onMQTTDisconnect is called when the MQTT client disconnects.
|
|
func (c *mipsClient) onMQTTDisconnect() {
|
|
c.mu.Lock()
|
|
c.mqttConnected = false
|
|
// Clear pending subs
|
|
c.pendingSubMu.Lock()
|
|
c.pendingSubs = make(map[string]int)
|
|
c.pendingSubMu.Unlock()
|
|
c.mu.Unlock()
|
|
|
|
// Try reconnect
|
|
c.tryReconnect(false)
|
|
}
|
|
|
|
// onMQTTMessage handles an incoming MQTT message.
|
|
func (c *mipsClient) onMQTTMessage(topic string, payload []byte) {
|
|
payloadStr := string(payload)
|
|
|
|
// Try exact match first
|
|
c.subMu.Lock()
|
|
bc, ok := c.subscriptions[topic]
|
|
c.subMu.Unlock()
|
|
|
|
if ok && bc != nil && bc.Handler != nil {
|
|
bc.Handler(topic, payloadStr, bc.HandlerCtx)
|
|
return
|
|
}
|
|
|
|
// Wildcard match: check subscriptions ending with # or +
|
|
for subTopic, bc := range c.copySubs() {
|
|
if bc == nil || bc.Handler == nil {
|
|
continue
|
|
}
|
|
if matchWildcard(subTopic, topic) {
|
|
bc.Handler(topic, payloadStr, bc.HandlerCtx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// copySubs returns a copy of subscriptions for safe wildcard iteration.
|
|
func (c *mipsClient) copySubs() map[string]*MipsBroadcast {
|
|
c.subMu.Lock()
|
|
defer c.subMu.Unlock()
|
|
result := make(map[string]*MipsBroadcast, len(c.subscriptions))
|
|
for k, v := range c.subscriptions {
|
|
result[k] = v
|
|
}
|
|
return result
|
|
}
|
|
|
|
// matchWildcard checks if a topic matches a subscription pattern containing # or +.
|
|
func matchWildcard(pattern, topic string) bool {
|
|
// Remove trailing /# from pattern and check prefix
|
|
if strings.HasSuffix(pattern, "/#") {
|
|
prefix := strings.TrimSuffix(pattern, "/#")
|
|
return strings.HasPrefix(topic, prefix+"/") || topic == prefix
|
|
}
|
|
return pattern == topic
|
|
}
|
|
|
|
// flushPendingSubs processes pending subscriptions in batches.
|
|
// Aligns with Python: __mips_sub_internal_pending_handler.
|
|
func (c *mipsClient) flushPendingSubs() {
|
|
if c.mqttConn == nil || !c.mqttConn.IsConnected() {
|
|
return
|
|
}
|
|
|
|
// 收集所有待订阅 topic,一次 SUBSCRIBE 包批量订阅(大幅减少网络往返)
|
|
c.pendingSubMu.Lock()
|
|
batch := make(map[string]byte, len(c.pendingSubs))
|
|
for topic, count := range c.pendingSubs {
|
|
if count > 3 {
|
|
delete(c.pendingSubs, topic)
|
|
logrus.Debugf("[mips] retry sub exceeded: %s", topic)
|
|
continue
|
|
}
|
|
batch[topic] = byte(mipsQoS)
|
|
}
|
|
c.pendingSubMu.Unlock()
|
|
|
|
if len(batch) == 0 {
|
|
return
|
|
}
|
|
|
|
if err := c.mqttConn.SubscribeMany(batch); err == nil {
|
|
c.pendingSubMu.Lock()
|
|
for t := range batch {
|
|
delete(c.pendingSubs, t)
|
|
logrus.Debugf("[mips] sub success: %s", t)
|
|
}
|
|
hasMore := len(c.pendingSubs) > 0
|
|
c.pendingSubMu.Unlock()
|
|
if hasMore {
|
|
time.Sleep(time.Duration(mipsSubInterval) * time.Second)
|
|
c.flushPendingSubs()
|
|
}
|
|
} else {
|
|
logrus.Debugf("[mips] batch sub error: %v", err)
|
|
c.pendingSubMu.Lock()
|
|
for t := range batch {
|
|
c.pendingSubs[t] = c.pendingSubs[t] + 1
|
|
}
|
|
c.pendingSubMu.Unlock()
|
|
time.Sleep(time.Duration(mipsSubInterval) * time.Second)
|
|
c.flushPendingSubs()
|
|
}
|
|
}
|
|
|
|
// sendRequest sends a MIPS request and waits for a reply.
|
|
func (c *mipsClient) sendRequest(topic string, payload string, timeout time.Duration) (string, error) {
|
|
mid := c.nextMid()
|
|
replyCh := make(chan string, 1)
|
|
|
|
req := &MipsRequest{
|
|
Mid: mid,
|
|
OnReply: func(payload string, ctx interface{}) {
|
|
select {
|
|
case ctx.(chan string) <- payload:
|
|
default:
|
|
}
|
|
},
|
|
OnReplyCtx: replyCh,
|
|
}
|
|
|
|
c.requestMu.Lock()
|
|
c.requests[fmt.Sprintf("%d", mid)] = req
|
|
c.requestMu.Unlock()
|
|
|
|
// Build packed message
|
|
packed, err := PackMipsMessage(mid, payload, "", topic)
|
|
if err != nil {
|
|
c.requestMu.Lock()
|
|
delete(c.requests, fmt.Sprintf("%d", mid))
|
|
c.requestMu.Unlock()
|
|
return "", err
|
|
}
|
|
|
|
if err := c.publish(topic, packed); err != nil {
|
|
c.requestMu.Lock()
|
|
delete(c.requests, fmt.Sprintf("%d", mid))
|
|
c.requestMu.Unlock()
|
|
return "", err
|
|
}
|
|
|
|
// Wait for reply or timeout
|
|
select {
|
|
case reply := <-replyCh:
|
|
return reply, nil
|
|
case <-time.After(timeout):
|
|
c.requestMu.Lock()
|
|
delete(c.requests, fmt.Sprintf("%d", mid))
|
|
c.requestMu.Unlock()
|
|
return `{"error":{"code":-10006,"message":"timeout"}}`, nil
|
|
}
|
|
}
|
|
|
|
// handleReply processes a reply from mips.
|
|
func (c *mipsClient) handleReply(mid int, payload string) {
|
|
key := fmt.Sprintf("%d", mid)
|
|
c.requestMu.Lock()
|
|
req, ok := c.requests[key]
|
|
if ok {
|
|
delete(c.requests, key)
|
|
}
|
|
c.requestMu.Unlock()
|
|
|
|
if req != nil && req.OnReply != nil {
|
|
req.OnReply(payload, req.OnReplyCtx)
|
|
}
|
|
}
|
|
|
|
// nextMid generates the next MIPS message ID.
|
|
func (c *mipsClient) nextMid() int {
|
|
c.mu.Lock()
|
|
mid := c.seedID
|
|
c.seedID = int((c.seedID + 1) % uint32max)
|
|
c.mu.Unlock()
|
|
return mid
|
|
}
|
|
|
|
// tryReconnect schedules a reconnection attempt with exponential backoff.
|
|
func (c *mipsClient) tryReconnect(immediate bool) {
|
|
c.mu.Lock()
|
|
if !c.reconnectActive {
|
|
c.mu.Unlock()
|
|
return
|
|
}
|
|
|
|
if c.reconnectTimer != nil {
|
|
c.reconnectTimer.Stop()
|
|
}
|
|
|
|
var interval time.Duration
|
|
if !immediate {
|
|
if c.reconnectInterval < mipsReconnectIntervalMin {
|
|
c.reconnectInterval = mipsReconnectIntervalMin
|
|
} else {
|
|
c.reconnectInterval = math.Min(c.reconnectInterval*2, mipsReconnectIntervalMax)
|
|
}
|
|
interval = time.Duration(c.reconnectInterval) * time.Second
|
|
logrus.Debugf("[mips] reconnect after %v", interval)
|
|
}
|
|
|
|
c.reconnectTimer = time.AfterFunc(interval, func() {
|
|
if c.mqttConn != nil {
|
|
if err := c.mqttConn.Connect(); err != nil {
|
|
logrus.Debugf("[mips] reconnect failed: %v", err)
|
|
c.tryReconnect(false)
|
|
}
|
|
}
|
|
})
|
|
c.mu.Unlock()
|
|
}
|
|
|
|
// ============================================================================
|
|
// MipsCloudClient — cloud MIPS client (public API)
|
|
// ============================================================================
|
|
|
|
// MipsCloudClient wraps the MIPS client for cloud connections.
|
|
type MipsCloudClient struct {
|
|
client *mipsClient
|
|
onConn func()
|
|
onDisc func(error)
|
|
}
|
|
|
|
// SubState subscribes to MIPS cloud connection state changes.
|
|
// key is a unique subscriber identifier (e.g. "{uid}-{cloudServer}").
|
|
func (c *MipsCloudClient) SubState(key string, handler func(key string, connected bool)) {
|
|
c.client.subState(key, handler)
|
|
}
|
|
|
|
// UnsubState unsubscribes from MIPS cloud connection state changes.
|
|
func (c *MipsCloudClient) UnsubState(key string) {
|
|
c.client.unsubState(key)
|
|
}
|
|
|
|
// UpdateAccessToken updates the MQTT password (access token) and reconnects.
|
|
// Aligns with Python: __update_mips_access_token.
|
|
func (c *MipsCloudClient) UpdateAccessToken(accessToken string) {
|
|
c.client.mu.Lock()
|
|
c.client.password = accessToken
|
|
wasConnected := c.client.mqttConnected
|
|
c.client.mu.Unlock()
|
|
|
|
if wasConnected {
|
|
c.Disconnect()
|
|
c.Connect()
|
|
}
|
|
}
|
|
|
|
// NewMipsCloudClient creates a new MipsCloudClient.
|
|
// broker format: "ssl://cn-ha.mqtt.io.mi.com:8883"
|
|
func NewMipsCloudClient(broker, clientID, username, password string) *MipsCloudClient {
|
|
host, port := parseMQTTBroker(broker)
|
|
c := &MipsCloudClient{
|
|
client: newMipsClient(clientID, host, port, username, password, "", "", ""),
|
|
}
|
|
return c
|
|
}
|
|
|
|
// parseMQTTBroker parses a broker URL like "ssl://host:port".
|
|
func parseMQTTBroker(broker string) (string, int) {
|
|
// Strip protocol prefix
|
|
url := broker
|
|
if strings.HasPrefix(url, "ssl://") {
|
|
url = url[6:]
|
|
} else if strings.HasPrefix(url, "tcp://") {
|
|
url = url[6:]
|
|
}
|
|
|
|
host := url
|
|
port := 8883
|
|
if idx := strings.LastIndex(url, ":"); idx >= 0 {
|
|
host = url[:idx]
|
|
fmt.Sscanf(url[idx+1:], "%d", &port)
|
|
}
|
|
return host, port
|
|
}
|
|
|
|
// Connect connects to the cloud MQTT broker.
|
|
func (c *MipsCloudClient) Connect() error {
|
|
mqClient := newPahoMQTTClient(c.client.clientID, c.client.host, c.client.port,
|
|
c.client.username, c.client.password, "", "", "")
|
|
mqClient.onConnect = c.client.onMQTTConnect
|
|
mqClient.onDisconnect = c.client.onMQTTDisconnect
|
|
mqClient.onMessage = c.client.onMQTTMessage
|
|
|
|
c.client.onConnect = func() {
|
|
if c.onConn != nil {
|
|
c.onConn()
|
|
}
|
|
}
|
|
c.client.onDisconnect = func(err error) {
|
|
if c.onDisc != nil {
|
|
c.onDisc(err)
|
|
}
|
|
}
|
|
|
|
return c.client.connect(mqClient)
|
|
}
|
|
|
|
// Disconnect disconnects from the MQTT broker.
|
|
func (c *MipsCloudClient) Disconnect() {
|
|
c.client.disconnect()
|
|
}
|
|
|
|
// IsConnected returns true if connected.
|
|
func (c *MipsCloudClient) IsConnected() bool {
|
|
return c.client.isConnected()
|
|
}
|
|
|
|
// Subscribe subscribes to a topic with a handler.
|
|
func (c *MipsCloudClient) Subscribe(topic string, handler func(string, []byte)) error {
|
|
return c.client.subscribe(topic, func(_ string, payload string, _ interface{}) {
|
|
handler(topic, []byte(payload))
|
|
}, nil)
|
|
}
|
|
|
|
// Unsubscribe unsubscribes from a topic.
|
|
func (c *MipsCloudClient) Unsubscribe(topic string) error {
|
|
return c.client.unsubscribe(topic)
|
|
}
|
|
|
|
// Publish publishes a message to a topic.
|
|
func (c *MipsCloudClient) Publish(topic, payload string) error {
|
|
return c.client.publish(topic, []byte(payload))
|
|
}
|
|
|
|
// SubMipsState registers a handler for MIPS state changes.
|
|
func (c *MipsCloudClient) SubMipsState(handler func(string, interface{})) {
|
|
c.onConn = func() {
|
|
handler("connected", nil)
|
|
}
|
|
c.onDisc = func(err error) {
|
|
handler("disconnected", err)
|
|
}
|
|
}
|
|
|
|
// GetDevList sends a get_dev_list request through MIPS.
|
|
func (c *MipsCloudClient) GetDevList(payload string, timeout time.Duration) (string, error) {
|
|
return c.client.sendRequest("get_dev_list", payload, timeout)
|
|
}
|
|
|
|
// ============================================================================
|
|
// MipsLocalClient — local MIPS client (public API)
|
|
// ============================================================================
|
|
|
|
// MipsLocalClient wraps the MIPS client for local connections.
|
|
type MipsLocalClient struct {
|
|
client *mipsClient
|
|
did string
|
|
// GroupID is the home group ID this client belongs to.
|
|
GroupID string
|
|
// OnDevListChanged is called when the gateway's device list changes.
|
|
OnDevListChanged func()
|
|
}
|
|
|
|
// NewMipsLocalClient creates a new MipsLocalClient.
|
|
func NewMipsLocalClient(did, host string, port int, caFile, certFile, keyFile string) *MipsLocalClient {
|
|
return &MipsLocalClient{
|
|
client: newMipsClient(did, host, port, "", "", caFile, certFile, keyFile),
|
|
did: did,
|
|
}
|
|
}
|
|
|
|
// Connect connects to the local MQTT broker.
|
|
func (c *MipsLocalClient) Connect() error {
|
|
mqClient := newPahoMQTTClient(c.client.clientID, c.client.host, c.client.port,
|
|
"", "", c.client.caFile, c.client.certFile, c.client.keyFile)
|
|
mqClient.onConnect = c.client.onMQTTConnect
|
|
mqClient.onDisconnect = c.client.onMQTTDisconnect
|
|
mqClient.onMessage = func(topic string, payload []byte) {
|
|
// Try unpacking as MIPS message first
|
|
msg := UnpackMipsMessage(payload)
|
|
if msg != nil && msg.Mid != 0 {
|
|
c.client.handleReply(msg.Mid, msg.Payload)
|
|
}
|
|
// Also pass to generic handler
|
|
c.client.onMQTTMessage(topic, payload)
|
|
}
|
|
return c.client.connect(mqClient)
|
|
}
|
|
|
|
// Disconnect disconnects from the local MQTT broker.
|
|
func (c *MipsLocalClient) Disconnect() {
|
|
c.client.disconnect()
|
|
}
|
|
|
|
// IsConnected returns true if connected.
|
|
func (c *MipsLocalClient) IsConnected() bool {
|
|
return c.client.isConnected()
|
|
}
|
|
|
|
// Subscribe subscribes to a topic.
|
|
func (c *MipsLocalClient) Subscribe(topic string, handler func(string, []byte)) error {
|
|
return c.client.subscribe(topic, func(_ string, payload string, _ interface{}) {
|
|
handler(topic, []byte(payload))
|
|
}, nil)
|
|
}
|
|
|
|
// Unsubscribe unsubscribes from a topic.
|
|
func (c *MipsLocalClient) Unsubscribe(topic string) error {
|
|
return c.client.unsubscribe(topic)
|
|
}
|
|
|
|
// GetDevList returns the device list from the gateway.
|
|
func (c *MipsLocalClient) GetDevList(payload string, timeout time.Duration) (string, error) {
|
|
return c.client.sendRequest("getDevList", payload, timeout)
|
|
}
|
|
|
|
// GetProps fetches property values from the gateway via MIPS request.
|
|
// params is a JSON-encoded string of property parameters.
|
|
func (c *MipsLocalClient) GetProps(params []PropParam, timeout time.Duration) ([]PropResult, error) {
|
|
batch := make([]map[string]interface{}, len(params))
|
|
for i, p := range params {
|
|
batch[i] = map[string]interface{}{
|
|
"did": p.DID,
|
|
"siid": float64(p.Siid),
|
|
"piid": float64(p.Piid),
|
|
}
|
|
}
|
|
payloadBytes, _ := json.Marshal(map[string]interface{}{
|
|
"method": "properties/get",
|
|
"params": batch,
|
|
})
|
|
payload, err := c.client.sendRequest("properties/get", string(payloadBytes), timeout)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// Parse response
|
|
var result struct {
|
|
Result []PropResult `json:"result"`
|
|
}
|
|
if err := json.Unmarshal([]byte(payload), &result); err != nil {
|
|
return nil, fmt.Errorf("parse gateway props response: %w", err)
|
|
}
|
|
return result.Result, nil
|
|
}
|
|
|
|
// SubState subscribes to MIPS local connection state changes.
|
|
func (c *MipsLocalClient) SubState(key string, handler func(key string, connected bool)) {
|
|
c.client.subState(key, handler)
|
|
}
|
|
|
|
// UnsubState unsubscribes from MIPS local connection state changes.
|
|
func (c *MipsLocalClient) UnsubState(key string) {
|
|
c.client.unsubState(key)
|
|
}
|
|
|
|
// DID returns the device ID.
|
|
func (c *MipsLocalClient) DID() string {
|
|
return c.did
|
|
}
|
|
|
|
// ============================================================================
|
|
// MipsDeviceState — device state tracking
|
|
// ============================================================================
|
|
|
|
// MipsDeviceState represents the state of a device's MIPS connection.
|
|
type MipsDeviceState struct {
|
|
DID string
|
|
Online bool
|
|
Source string
|
|
Handler func(string, interface{})
|
|
}
|
|
|
|
// NewMipsDeviceState creates a new MipsDeviceState.
|
|
func NewMipsDeviceState(did string, online bool, source string) *MipsDeviceState {
|
|
return &MipsDeviceState{
|
|
DID: did,
|
|
Online: online,
|
|
Source: source,
|
|
}
|
|
}
|
|
|
|
// ============================================================================
|
|
// Paho MQTT Client adapter
|
|
// ============================================================================
|
|
|
|
// pahoMQTTClient adapts paho.mqtt.golang to our MQTTClient interface.
|
|
type pahoMQTTClient struct {
|
|
clientID string
|
|
host string
|
|
port int
|
|
username string
|
|
password string
|
|
caFile string
|
|
certFile string
|
|
keyFile string
|
|
|
|
onConnect func()
|
|
onDisconnect func()
|
|
onMessage func(topic string, payload []byte)
|
|
|
|
client mqtt.Client
|
|
connected bool
|
|
}
|
|
|
|
func newPahoMQTTClient(clientID, host string, port int, username, password, caFile, certFile, keyFile string) *pahoMQTTClient {
|
|
c := &pahoMQTTClient{
|
|
clientID: clientID,
|
|
host: host,
|
|
port: port,
|
|
username: username,
|
|
password: password,
|
|
caFile: caFile,
|
|
certFile: certFile,
|
|
keyFile: keyFile,
|
|
}
|
|
return c
|
|
}
|
|
|
|
func (p *pahoMQTTClient) Connect() error {
|
|
opts := mqtt.NewClientOptions()
|
|
opts.AddBroker(fmt.Sprintf("ssl://%s:%d", p.host, p.port))
|
|
opts.SetClientID(p.clientID)
|
|
opts.SetUsername(p.username)
|
|
opts.SetPassword(p.password)
|
|
opts.SetKeepAlive(time.Duration(MIHOME_MQTT_KEEPALIVE) * time.Second)
|
|
opts.SetAutoReconnect(false)
|
|
opts.SetCleanSession(true)
|
|
|
|
opts.OnConnect = func(c mqtt.Client) {
|
|
logrus.Debugf("[mips] MQTT connected to %s:%d", p.host, p.port)
|
|
p.connected = true
|
|
if p.onConnect != nil {
|
|
p.onConnect()
|
|
}
|
|
}
|
|
|
|
opts.OnConnectionLost = func(c mqtt.Client, err error) {
|
|
logrus.Debugf("[mips] MQTT connection lost: %v", err)
|
|
p.connected = false
|
|
if p.onDisconnect != nil {
|
|
p.onDisconnect()
|
|
}
|
|
}
|
|
|
|
opts.DefaultPublishHandler = func(c mqtt.Client, msg mqtt.Message) {
|
|
logrus.Debugf("[mips] MQTT message: %s (%d bytes) data:%s", msg.Topic(), len(msg.Payload()), string(msg.Payload()))
|
|
if p.onMessage != nil {
|
|
p.onMessage(msg.Topic(), msg.Payload())
|
|
}
|
|
}
|
|
|
|
p.client = mqtt.NewClient(opts)
|
|
token := p.client.Connect()
|
|
if token.Wait() && token.Error() != nil {
|
|
return token.Error()
|
|
}
|
|
|
|
logrus.Debugf("[mips] MQTT connected to %s:%d as %s", p.host, p.port, p.clientID)
|
|
return nil
|
|
}
|
|
|
|
func (p *pahoMQTTClient) Disconnect() {
|
|
if p.client != nil && p.client.IsConnected() {
|
|
p.client.Disconnect(250)
|
|
}
|
|
p.connected = false
|
|
if p.onDisconnect != nil {
|
|
p.onDisconnect()
|
|
}
|
|
}
|
|
|
|
func (p *pahoMQTTClient) IsConnected() bool {
|
|
return p.connected
|
|
}
|
|
|
|
func (p *pahoMQTTClient) Subscribe(topic string, qos byte, handler func(topic string, payload []byte)) error {
|
|
if p.client == nil || !p.client.IsConnected() {
|
|
return fmt.Errorf("not connected")
|
|
}
|
|
|
|
// Use DefaultPublishHandler for message delivery (set in Connect).
|
|
// Don't set a per-subscription callback to avoid double-firing with DefaultPublishHandler.
|
|
token := p.client.Subscribe(topic, qos, nil)
|
|
if token.Wait() && token.Error() != nil {
|
|
return token.Error()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SubscribeMany 用单个 SUBSCRIBE 包批量订阅多个 topic,显著减少网络往返。
|
|
// 消息分发仍走 DefaultPublishHandler,与 Subscribe 行为一致。
|
|
func (p *pahoMQTTClient) SubscribeMany(topics map[string]byte) error {
|
|
if p.client == nil || !p.client.IsConnected() {
|
|
return fmt.Errorf("not connected")
|
|
}
|
|
if len(topics) == 0 {
|
|
return nil
|
|
}
|
|
token := p.client.SubscribeMultiple(topics, nil)
|
|
if token.Wait() && token.Error() != nil {
|
|
return token.Error()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *pahoMQTTClient) Unsubscribe(topic string) error {
|
|
if p.client == nil || !p.client.IsConnected() {
|
|
return fmt.Errorf("not connected")
|
|
}
|
|
|
|
token := p.client.Unsubscribe(topic)
|
|
if token.Wait() && token.Error() != nil {
|
|
return token.Error()
|
|
}
|
|
logrus.Debugf("[mips] unsub success: %s", topic)
|
|
return nil
|
|
}
|
|
|
|
func (p *pahoMQTTClient) Publish(topic string, qos byte, payload []byte) error {
|
|
if p.client == nil || !p.client.IsConnected() {
|
|
return fmt.Errorf("not connected")
|
|
}
|
|
|
|
token := p.client.Publish(topic, qos, false, payload)
|
|
if token.Wait() && token.Error() != nil {
|
|
return token.Error()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (p *pahoMQTTClient) SetCredentials(username, password string) {
|
|
p.username = username
|
|
p.password = password
|
|
}
|
|
|
|
func (p *pahoMQTTClient) SetTLS(caFile, certFile, keyFile string) error {
|
|
p.caFile = caFile
|
|
p.certFile = certFile
|
|
p.keyFile = keyFile
|
|
return nil
|
|
}
|