Files
xiaomihome/miot/miot_client_sub.go
T
4566704 2b7de69403 refactor(logger): 日志改用 logrus 全局包级调用,移除 Logger 注入接口
为什么:并入宿主项目后须与宿主日志体系一致——宿主统一用 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)
2026-09-15 20:16:26 +08:00

501 lines
14 KiB
Go

// Package miot provides MIoT core client for Xiaomi Home devices.
// It's a Golang port of MIoTClient from xiaomihome-py/miot/.
package miot
import (
"encoding/json"
"fmt"
"github.com/sirupsen/logrus"
)
// ============================================================================
// Subscription management (aligned with Python miot_client.py)
// ============================================================================
// SubProp subscribes to device property changes.
// Aligns with Python: sub_prop method.
// did: device ID
// handler: callback function called when property changes
// siid, piid: property identifiers
// Returns: subscription ID for unsubscribing.
func (c *MIoTClient) SubProp(did string, handler func(map[string]interface{}, interface{}), siid, piid int) string {
c.mu.Lock()
// Build topic with wildcard for 0 values
topic := buildPropTopic(did, siid, piid)
subID := fmt.Sprintf("%p", handler)
c.subTree.Sub(topic, handler, nil)
logrus.Debugf("[MIoTClient] Subscribed to property: %s (topic=%s, subID=%s)\n", did, topic, subID)
c.mu.Unlock()
// siid=0 && piid=0 为通配符订阅,展开为设备所有可读属性逐条刷新
if siid == 0 && piid == 0 {
c.RefreshDeviceAllProps(did)
} else {
c.RequestRefreshProp(did, siid, piid)
}
return subID
}
// UnsubProp unsubscribes from device property changes.
// Aligns with Python: unsub_prop method.
func (c *MIoTClient) UnsubProp(did string, subID string) {
c.mu.Lock()
defer c.mu.Unlock()
// Remove from subTree by subID
c.subTree.Unsub(subID)
// Check if there are any remaining subscriptions for this device
// TODO: Implement proper tracking of per-device subscription count
// For now, just log
logrus.Debugf("[MIoTClient] Unsubscribed from property: %s (subID=%s)\n", did, subID)
}
// SubEvent subscribes to device events.
// Aligns with Python: sub_event method.
// did: device ID
// handler: callback function called when event occurs
// siid, eiid: event identifiers
// Returns: subscription ID for unsubscribing.
func (c *MIoTClient) SubEvent(did string, handler func(map[string]interface{}, interface{}), siid, eiid int) string {
c.mu.Lock()
defer c.mu.Unlock()
topic := buildEventTopic(did, siid, eiid)
subID := fmt.Sprintf("%p", handler)
c.subTree.Sub(topic, handler, nil)
logrus.Debugf("[MIoTClient] Subscribed to event: %s (topic=%s, subID=%s)\n", did, topic, subID)
return subID
}
// UnsubEvent unsubscribes from device events.
// Aligns with Python: unsub_event method.
func (c *MIoTClient) UnsubEvent(did string, subID string) {
c.mu.Lock()
defer c.mu.Unlock()
// Remove from subTree by subID
c.subTree.Unsub(subID)
logrus.Debugf("[MIoTClient] Unsubscribed from event: %s (subID=%s)\n", did, subID)
}
// updateDeviceMsgSub updates device message subscription.
// Aligns with Python: __update_device_msg_sub.
// Determines the best subscription source (gateway > LAN > cloud) and switches if needed.
func (c *MIoTClient) updateDeviceMsgSub(did string) {
c.mu.Lock()
defer c.mu.Unlock()
// Read device state from three sources
cloudOnline := false
gwOnline := false
gwPushAvailable := false
lanOnline := false
lanPushAvailable := false
gwGroupID := ""
if dev, ok := c.deviceListCloud[did]; ok {
if o, ok := dev["online"].(bool); ok {
cloudOnline = o
}
}
if dev, ok := c.deviceListGateway[did]; ok {
if o, ok := dev["online"].(bool); ok {
gwOnline = o
}
if p, ok := dev["push_available"].(bool); ok {
gwPushAvailable = p
}
gwGroupID, _ = dev["group_id"].(string)
}
if dev, ok := c.deviceListLan[did]; ok {
if o, ok := dev["online"].(bool); ok {
lanOnline = o
}
if p, ok := dev["push_available"].(bool); ok {
lanPushAvailable = p
}
}
// Priority: gateway(online+push_available) > LAN(online+push_available) > cloud(online)
newSource := "none"
if gwOnline && gwPushAvailable {
newSource = gwGroupID // gateway source = groupID
} else if lanOnline && lanPushAvailable {
newSource = "lan"
} else if cloudOnline {
newSource = "cloud"
}
oldSource, existed := c.subSourceList[did]
if !existed {
oldSource = "none"
}
// No change
if oldSource == newSource {
return
}
// Source changed: unsubscribe from old, subscribe to new
if oldSource != "none" {
c.unsubFromSource(oldSource, did)
}
if newSource != "none" {
c.subFromSource(newSource, did)
}
// Update subscription source
c.subSourceList[did] = newSource
logrus.Debugf("[MIoTClient] updateDeviceMsgSub: %s source %s → %s (gw=%v/%v, lan=%v/%v, cloud=%v)\n",
did, oldSource, newSource, gwOnline, gwPushAvailable, lanOnline, lanPushAvailable, cloudOnline)
}
// unsubFromSource unsubscribes device messages from a given source.
func (c *MIoTClient) unsubFromSource(source, did string) {
switch source {
case "cloud":
if c.mipsCloud != nil {
c.mipsCloud.Unsubscribe(fmt.Sprintf("device/%s/up/properties_changed/#", did))
c.mipsCloud.Unsubscribe(fmt.Sprintf("device/%s/up/event_occured/#", did))
c.mipsCloud.Unsubscribe(fmt.Sprintf("device/%s/state/#", did))
}
case "lan":
if c.miotLan != nil {
lan, ok := c.miotLan.(*MIoTLan)
if ok {
lan.UnsubProp(did, 0, 0)
lan.UnsubEvent(did, 0, 0)
}
}
default:
// Gateway: source is the groupID
if mips, exists := c.mipsLocal[source]; exists {
mips.Unsubscribe(fmt.Sprintf("device/%s/up/properties_changed/#", did))
mips.Unsubscribe(fmt.Sprintf("device/%s/up/event_occured/#", did))
mips.Unsubscribe(fmt.Sprintf("device/%s/state/#", did))
}
}
}
// subFromSource subscribes device messages from a given source.
func (c *MIoTClient) subFromSource(source, did string) {
switch source {
case "cloud":
if c.mipsCloud != nil {
c.subscribeDeviceMessages(c.mipsCloud, did)
}
case "lan":
if c.miotLan != nil {
lan, ok := c.miotLan.(*MIoTLan)
if ok {
lan.SubProp(did, nil, 0, 0, nil)
lan.SubEvent(did, nil, 0, 0, nil)
}
}
default:
// Gateway: source is the groupID
if mips, exists := c.mipsLocal[source]; exists {
c.subscribeDeviceMessagesOnLocal(mips, did)
}
}
}
// subscribeDeviceMessagesOnLocal subscribes to device messages from a local MIPS client.
func (c *MIoTClient) subscribeDeviceMessagesOnLocal(client *MipsLocalClient, did string) {
propTopic := fmt.Sprintf("device/%s/up/properties_changed/#", did)
client.Subscribe(propTopic, func(topic string, payload []byte) {
c.handleMIPSMessage(did, payload)
})
eventTopic := fmt.Sprintf("device/%s/up/event_occured/#", did)
client.Subscribe(eventTopic, func(topic string, payload []byte) {
c.handleMIPSMessage(did, payload)
})
stateTopic := fmt.Sprintf("device/%s/state/#", did)
client.Subscribe(stateTopic, func(topic string, payload []byte) {
c.handleMIPSStateMessage(did, payload)
})
}
// subscribeDeviceMessages subscribes to device property/event/state messages.
func (c *MIoTClient) subscribeDeviceMessages(client *MipsCloudClient, did string) {
// Subscribe to property changes
propTopic := fmt.Sprintf("device/%s/up/properties_changed/#", did)
client.Subscribe(propTopic, func(topic string, payload []byte) {
c.handleMIPSMessage(did, payload)
})
// Subscribe to events
eventTopic := fmt.Sprintf("device/%s/up/event_occured/#", did)
client.Subscribe(eventTopic, func(topic string, payload []byte) {
c.handleMIPSMessage(did, payload)
})
// Subscribe to state changes
stateTopic := fmt.Sprintf("device/%s/state/#", did)
client.Subscribe(stateTopic, func(topic string, payload []byte) {
c.handleMIPSStateMessage(did, payload)
})
}
// handleMIPSMessage handles MIPS property/event messages.
func (c *MIoTClient) handleMIPSMessage(did string, payload []byte) {
// First try: direct JSON (cloud MIPS push messages are plain JSON)
var params map[string]interface{}
if err := json.Unmarshal(payload, &params); err == nil && len(params) > 0 {
c.dispatchMIPSMessage(did, params)
return
}
// Fallback: binary MIPS protocol (used by gateway mipsLocal)
msg := UnpackMipsMessage(payload)
if msg == nil || msg.Payload == "" {
logrus.Debugf("[MIoTClient] handleMIPSMessage: empty or invalid payload for %s (%d bytes)\n", did, len(payload))
return
}
// Parse JSON from unpacked payload
if err := json.Unmarshal([]byte(msg.Payload), &params); err != nil {
logrus.Errorf("[MIoTClient] handleMIPSMessage: JSON parse error for %s: %v\n", did, err)
return
}
c.dispatchMIPSMessage(did, params)
}
// dispatchMIPSMessage dispatches parsed MIPS params to OnPropMsg or OnEventMsg.
func (c *MIoTClient) dispatchMIPSMessage(did string, params map[string]interface{}) {
// Determine if it's a property or event message
if inner, ok := params["params"].(map[string]interface{}); ok {
// Cloud MIPS format: {"params": {"did":"...", "siid":..., "piid":..., "value":...}}
if _, hasEiid := inner["eiid"]; hasEiid {
c.OnEventMsg(params, nil)
} else {
c.OnPropMsg(params, nil)
}
} else if _, ok := params["siid"]; ok {
// Direct format: {"did":"...", "siid":..., "piid":..., "value":...}
c.OnPropMsg(params, nil)
} else if _, ok := params["eiid"]; ok {
c.OnEventMsg(params, nil)
} else {
logrus.Debugf("[MIoTClient] handleMIPSMessage: unknown message format for %s: %v\n", did, params)
}
}
// handleMIPSStateMessage handles MIPS state messages (online/offline).
func (c *MIoTClient) handleMIPSStateMessage(did string, payload []byte) {
// First try: direct JSON
var params map[string]interface{}
if err := json.Unmarshal(payload, &params); err != nil || len(params) == 0 {
// Fallback: binary MIPS protocol
msg := UnpackMipsMessage(payload)
if msg == nil || msg.Payload == "" {
logrus.Debugf("[MIoTClient] handleMIPSStateMessage: empty payload for %s (%d bytes)\n", did, len(payload))
return
}
if err := json.Unmarshal([]byte(msg.Payload), &params); err != nil {
logrus.Errorf("[MIoTClient] handleMIPSStateMessage: JSON parse error for %s: %v\n", did, err)
return
}
}
// Extract state fields
state := make(map[string]interface{})
// Online/offline from typical state formats
if online, ok := params["online"].(bool); ok {
state["online"] = online
} else if status, ok := params["status"].(string); ok {
state["online"] = status == "online"
}
// Push availability
if push, ok := params["push_available"].(bool); ok {
state["push_available"] = push
}
if len(state) == 0 {
logrus.Debugf("[MIoTClient] handleMIPSStateMessage: no state info in payload for %s: %v\n", did, params)
return
}
// Call OnDeviceStateChanged to update cache and notify subscribers
c.OnDeviceStateChanged(did, state)
}
// ============================================================================
// MQTT message handling (aligned with Python miot_client.py)
// ============================================================================
// OnPropMsg handles MIPS property messages.
// Aligns with Python: __on_prop_msg.
func (c *MIoTClient) OnPropMsg(params map[string]interface{}, ctx interface{}) {
// 1. Parse params
var did string
var siid, piid int
var value interface{}
// Check if wrapped in params.params (cloud MIPS format)
if inner, ok := params["params"].(map[string]interface{}); ok {
if d, ok := inner["did"].(string); ok {
did = d
}
if s, ok := inner["siid"].(float64); ok {
siid = int(s)
}
if p, ok := inner["piid"].(float64); ok {
piid = int(p)
}
value = inner["value"]
} else {
if d, ok := params["did"].(string); ok {
did = d
}
if s, ok := params["siid"].(float64); ok {
siid = int(s)
}
if p, ok := params["piid"].(float64); ok {
piid = int(p)
}
value = params["value"]
}
if did == "" {
logrus.Debugf("[MIoTClient] OnPropMsg: missing did in params=%v\n", params)
return
}
// 2. Match subscription and dispatch to handlers
topic := fmt.Sprintf("%s/p/%d/%d", did, siid, piid)
entries := c.subTree.Match(topic)
for _, entry := range entries {
// Call handler asynchronously
go entry.Handler(map[string]interface{}{
"did": did,
"siid": siid,
"piid": piid,
"value": value,
}, entry.Ctx)
}
}
// OnEventMsg handles MIPS event messages.
// Aligns with Python: __on_event_msg.
func (c *MIoTClient) OnEventMsg(params map[string]interface{}, ctx interface{}) {
// 1. Parse params
var did string
var siid, eiid int
var value interface{}
if inner, ok := params["params"].(map[string]interface{}); ok {
if d, ok := inner["did"].(string); ok {
did = d
}
if s, ok := inner["siid"].(float64); ok {
siid = int(s)
}
if e, ok := inner["eiid"].(float64); ok {
eiid = int(e)
}
value = inner["value"]
} else {
if d, ok := params["did"].(string); ok {
did = d
}
if s, ok := params["siid"].(float64); ok {
siid = int(s)
}
if e, ok := params["eiid"].(float64); ok {
eiid = int(e)
}
value = params["value"]
}
if did == "" {
logrus.Debugf("[MIoTClient] OnEventMsg: missing did in params=%v\n", params)
return
}
// 2. Match subscription and dispatch to handlers
topic := fmt.Sprintf("%s/e/%d/%d", did, siid, eiid)
entries := c.subTree.Match(topic)
for _, entry := range entries {
// Call handler asynchronously
go entry.Handler(map[string]interface{}{
"did": did,
"siid": siid,
"eiid": eiid,
"value": value,
}, entry.Ctx)
}
}
// checkDeviceState checks and updates device state from multiple sources.
// Aligns with Python: __check_device_state.
// NOTE: Caller must hold c.deviceListMu (read or write lock).
func (c *MIoTClient) checkDeviceState(did string) bool {
// Check cloud state
if dev, ok := c.deviceListCloud[did]; ok {
if online, ok := dev["online"].(bool); ok && online {
return true
}
}
// Check gateway state
if dev, ok := c.deviceListGateway[did]; ok {
if online, ok := dev["online"].(bool); ok && online {
return true
}
}
// Check LAN state
if dev, ok := c.deviceListLan[did]; ok {
if online, ok := dev["online"].(bool); ok && online {
return true
}
}
return false
}
// buildPropTopic builds a subscription topic, converting 0 to + wildcard.
// siid=0 → match any siid; piid=0 → match any piid.
func buildPropTopic(did string, siid, piid int) string {
s := "+"
p := "+"
if siid != 0 {
s = fmt.Sprintf("%d", siid)
}
if piid != 0 {
p = fmt.Sprintf("%d", piid)
}
return fmt.Sprintf("%s/p/%s/%s", did, s, p)
}
func buildEventTopic(did string, siid, eiid int) string {
s := "+"
e := "+"
if siid != 0 {
s = fmt.Sprintf("%d", siid)
}
if eiid != 0 {
e = fmt.Sprintf("%d", eiid)
}
return fmt.Sprintf("%s/e/%s/%s", did, s, e)
}