- 新增 xiaomi/examples/ 下 8 个示例程序(用户/家庭查询、设备列表、开关/灯光/空调控制、属性订阅、高级过滤分类、SPEC 解析) - miot_client_sub.go: 重构 SubProp/SubEvent,使用 buildPropTopic/buildEventTopic 支持通配符订阅(siid/piid=0 → +),并修复锁顺序问题(将 RequestRefreshProp 移到 Lock 外) - spec_parser.go: 新增 downloadSpecFile 方法,本地 SPEC 文件缺失时自动从 miot-spec.org 下载 - ARCH_PLAN.md: 架构设计从 Proposed 更新为 Accepted(v1.0→v1.1),补充设备分类/工厂/SPEC 映射等模块设计 - xiaomi/: 新增 miot 上层强类型封装模块,包含 Client 主入口、用户/家庭/设备 API、属性读写、动作调用、订阅通知,以及 devices/ 设备控制抽象(Switch/Light/AirConditioner/Fan/Cover/Humidifier/Vacuum/WaterHeater/Thermostat)和 specs/ SPEC 查询辅助 - 更多 xiaomi 示例(风扇/窗帘/传感器控制)
495 lines
14 KiB
Go
495 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"
|
|
)
|
|
|
|
// ============================================================================
|
|
// 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)
|
|
|
|
fmt.Printf("[MIoTClient] Subscribed to property: %s (topic=%s, subID=%s)\n", did, topic, subID)
|
|
|
|
c.mu.Unlock()
|
|
|
|
// RequestRefreshProp must be outside the lock (it also acquires c.mu)
|
|
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
|
|
|
|
fmt.Printf("[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)
|
|
|
|
fmt.Printf("[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)
|
|
|
|
fmt.Printf("[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
|
|
|
|
fmt.Printf("[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, ¶ms); 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 == "" {
|
|
fmt.Printf("[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), ¶ms); err != nil {
|
|
fmt.Printf("[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 {
|
|
fmt.Printf("[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, ¶ms); err != nil || len(params) == 0 {
|
|
// Fallback: binary MIPS protocol
|
|
msg := UnpackMipsMessage(payload)
|
|
if msg == nil || msg.Payload == "" {
|
|
fmt.Printf("[MIoTClient] handleMIPSStateMessage: empty payload for %s (%d bytes)\n", did, len(payload))
|
|
return
|
|
}
|
|
if err := json.Unmarshal([]byte(msg.Payload), ¶ms); err != nil {
|
|
fmt.Printf("[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 {
|
|
fmt.Printf("[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 == "" {
|
|
fmt.Printf("[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 == "" {
|
|
fmt.Printf("[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)
|
|
}
|