feat: 初始化小米 IoT (MIoT) 智能家居 Go 库
- 实现 MIoT 客户端核心功能(MQTT 连接、设备管理、属性读写) - 支持云端 API 调用与局域网设备发现(mDNS) - 集成国际化(i18n)多语言支持 - 添加 MIoT 设备规约解析器(spec_parser) - 包含单元测试与使用示例
This commit is contained in:
@@ -0,0 +1,479 @@
|
||||
// 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()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
// 1. Build topic
|
||||
topic := fmt.Sprintf("%s/p/%d/%d", did, siid, piid)
|
||||
|
||||
// 2. Generate subscription ID
|
||||
subID := fmt.Sprintf("%p", handler)
|
||||
|
||||
// 3. Store handler in subTree
|
||||
c.subTree.Sub(topic, handler, nil)
|
||||
|
||||
// 4. Check if this is the first subscription for this device
|
||||
// TODO: Implement proper tracking
|
||||
|
||||
fmt.Printf("[MIoTClient] Subscribed to property: %s (topic=%s, subID=%s)\n", did, topic, subID)
|
||||
|
||||
// 5. Request refresh property from cloud
|
||||
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()
|
||||
|
||||
// 1. Build topic
|
||||
topic := fmt.Sprintf("%s/e/%d/%d", did, siid, eiid)
|
||||
|
||||
// 2. Generate subscription ID
|
||||
subID := fmt.Sprintf("%p", handler)
|
||||
|
||||
// 3. Store handler in subTree
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user