**命名统一(State → Props)** - 所有设备类型:XxxState → XxxProps,GetState() → GetProps(),FetchState() → FetchProps() - 回调接口:OnStateChanged → OnPropsChanged - Bridge 层类型别名、示例代码、文档同步更名 **新增设备类型** - KindPlug + Plug 驱动:智能插座/插头(TurnOn/TurnOff/IsOn/OnPropsChanged) - KindGateway + Gateway 驱动:网关设备类型识别 - KindOccupancySensor 从 KindSensor 独立,添加 OccupancyProps **设备分类系统改进** - classifyByModel 改为按 key 长度降序匹配(防止 sensor 先于 sensor-occupy) - 新增下划线变体识别(sensor_occupy、sensor_temp 等) - outlet 重新归类为 Plug(非 Switch) - 新增 Or() 组合过滤器 **连接桥增强** - 设备创建成功后自动订阅属性变更 - 新增 nil 防护检查 - 支持 Plug/Speaker/Gateway/Occupancy 类型的 FetchProps 刷新 - 日志级别 Info → Debug **AC 驱动增强** - 新增 HasPower()、GetElectricPower()、GetPowerConsumption() 功率接口 - ACProps 新增 ElectricPower、PowerConsumption 字段 - valMapper:Desc 为空时 fallback 用 Name,key 统一小写 **MIoT 核心** - 设备列表刷新改用 reflect.DeepEqual 全字段变更检测 - 移除 BLE/代理设备默认在线 hack - RefreshDeviceAllProps 增加离线检查 - 新增 miot_errors.go:MIoT 错误码定义 - 全局日志级别 Info→Debug/Warn(降低噪音) **示例与文档** - 所有示例代码同步 API 变更 - 传感器示例改用 Or 组合过滤器 - ARCH_PLAN 等迁移文档同步更新
434 lines
12 KiB
Go
434 lines
12 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 (
|
|
"fmt"
|
|
"reflect"
|
|
"regexp"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// subDevPattern matches .s\d+ suffix for sub-devices.
|
|
var subDevPattern = regexp.MustCompile(`\.s\d+$`)
|
|
|
|
// ============================================================================
|
|
// Device instance management
|
|
// ============================================================================
|
|
|
|
// getOrCreateDevice returns or creates a MIoTDevice instance for the given DID.
|
|
// Lazily parses SPEC on first access and caches the result.
|
|
func (c *MIoTClient) getOrCreateDevice(did string, info map[string]interface{}) *MIoTDevice {
|
|
c.mu.RLock()
|
|
if dev, ok := c.devices[did]; ok {
|
|
c.mu.RUnlock()
|
|
return dev
|
|
}
|
|
c.mu.RUnlock()
|
|
|
|
if c.specParser == nil {
|
|
return nil
|
|
}
|
|
|
|
urn, _ := info["urn"].(string)
|
|
model, _ := info["model"].(string)
|
|
if urn == "" && model == "" {
|
|
urn, _ = info["type"].(string)
|
|
}
|
|
|
|
spec, err := c.specParser.Parse(urn, model)
|
|
if err != nil {
|
|
c.lgr.Errorf("[MIoTClient] getOrCreateDevice: SPEC parse failed for %s (%s/%s): %v\n", did, urn, model, err)
|
|
return nil
|
|
}
|
|
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
// Double-check after acquiring write lock
|
|
if dev, ok := c.devices[did]; ok {
|
|
return dev
|
|
}
|
|
dev := NewMIoTDevice(c, info, spec)
|
|
if c.devices == nil {
|
|
c.devices = make(map[string]*MIoTDevice)
|
|
}
|
|
c.devices[did] = dev
|
|
return dev
|
|
}
|
|
|
|
// ============================================================================
|
|
// Device management (aligned with Python miot_client.py)
|
|
// ============================================================================
|
|
|
|
// RefreshDevices refreshes device list from cloud for the given home IDs (nil=all).
|
|
// Aligns with Python: __update_devices_from_cloud_async (line 1384).
|
|
func (c *MIoTClient) RefreshDevices(homeIDs []string) error {
|
|
if c.http == nil {
|
|
return fmt.Errorf("HTTP client not available")
|
|
}
|
|
|
|
// 1. Get full device list from cloud
|
|
result, err := c.http.GetDevices(homeIDs)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to get devices: %v", err)
|
|
}
|
|
|
|
devicesRaw, ok := result["devices"].(map[string]map[string]interface{})
|
|
if !ok {
|
|
// Try interface{} cast
|
|
dev, ok := result["devices"].(map[string]interface{})
|
|
if !ok {
|
|
return fmt.Errorf("invalid devices response")
|
|
}
|
|
devicesRaw = make(map[string]map[string]interface{})
|
|
for did, info := range dev {
|
|
if infoMap, ok := info.(map[string]interface{}); ok {
|
|
devicesRaw[did] = infoMap
|
|
}
|
|
}
|
|
}
|
|
|
|
c.deviceListMu.Lock()
|
|
defer c.deviceListMu.Unlock()
|
|
|
|
// 2. Build processed device maps
|
|
cloudDevices := make(map[string]map[string]interface{})
|
|
sharedDevices := make(map[string]map[string]interface{})
|
|
subDevices := make(map[string]map[string]map[string]interface{}) // parentDID -> subKey -> device
|
|
|
|
for did, device := range devicesRaw {
|
|
// Clone device info
|
|
info := make(map[string]interface{})
|
|
for k, v := range device {
|
|
info[k] = v
|
|
}
|
|
|
|
// Handle sub-devices: .s\d+ suffix
|
|
match := subDevPattern.FindString(did)
|
|
if match != "" {
|
|
parentDID := strings.TrimSuffix(did, match)
|
|
subKey := strings.TrimPrefix(match, ".")
|
|
if subDevices[parentDID] == nil {
|
|
subDevices[parentDID] = make(map[string]map[string]interface{})
|
|
}
|
|
subDevices[parentDID][subKey] = info
|
|
continue
|
|
}
|
|
|
|
// Handle shared devices: owner.userid exists
|
|
if isSharedDevice(info) {
|
|
sharedDevices[did] = info
|
|
continue
|
|
}
|
|
|
|
cloudDevices[did] = info
|
|
}
|
|
|
|
// Merge sub-devices into parents
|
|
for parentDID, subs := range subDevices {
|
|
if parent, exists := cloudDevices[parentDID]; exists {
|
|
if parent["sub_devices"] == nil {
|
|
parent["sub_devices"] = make(map[string]interface{})
|
|
}
|
|
existingSubs, _ := parent["sub_devices"].(map[string]interface{})
|
|
for subKey, subInfo := range subs {
|
|
existingSubs[subKey] = subInfo
|
|
}
|
|
}
|
|
}
|
|
|
|
// 3. Merge shared devices into cloud devices (python treats them the same)
|
|
for did, info := range sharedDevices {
|
|
cloudDevices[did] = info
|
|
}
|
|
|
|
// 4. Diff: compare cloudDevices with current deviceListCloud
|
|
changed := false
|
|
|
|
// Detect additions and updates
|
|
for did, info := range cloudDevices {
|
|
oldInfo, exists := c.deviceListCloud[did]
|
|
if !exists {
|
|
// New device
|
|
c.deviceListCloud[did] = info
|
|
if cacheInfo, cacheExists := c.deviceListCache[did]; cacheExists {
|
|
// Update existing cache entry
|
|
for k, v := range info {
|
|
cacheInfo[k] = v
|
|
}
|
|
} else {
|
|
c.deviceListCache[did] = cloneMap(info)
|
|
}
|
|
c.lgr.Debugf("[MIoTClient] Device added: %s (%v)\n", did, info["name"])
|
|
changed = true
|
|
} else {
|
|
// 检测并同步所有变更字段(name/model/nickname 等)
|
|
for k, v := range info {
|
|
if oldV, ok := oldInfo[k]; !ok || !reflect.DeepEqual(oldV, v) {
|
|
oldInfo[k] = v
|
|
changed = true
|
|
}
|
|
}
|
|
c.deviceListCloud[did] = oldInfo
|
|
if cacheInfo, ok := c.deviceListCache[did]; ok {
|
|
for k, v := range info {
|
|
cacheInfo[k] = v
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Detect deletions
|
|
for did := range c.deviceListCloud {
|
|
if _, exists := cloudDevices[did]; !exists {
|
|
delete(c.deviceListCloud, did)
|
|
if cacheInfo, ok := c.deviceListCache[did]; ok {
|
|
cacheInfo["online"] = false
|
|
cacheInfo["push_available"] = false
|
|
}
|
|
c.lgr.Infof("[MIoTClient] Device removed from cloud: %s\n", did)
|
|
changed = true
|
|
}
|
|
}
|
|
|
|
// 4. Aggregate three sources: cloud + gateway + lan
|
|
for did := range c.deviceListCache {
|
|
c.aggregateDeviceState(did)
|
|
}
|
|
|
|
// 5. Update subscription routing
|
|
for did := range cloudDevices {
|
|
c.updateDeviceMsgSub(did)
|
|
}
|
|
|
|
// 6. Save device list to storage
|
|
if c.storage != nil {
|
|
deviceListForStorage := make(map[string]interface{})
|
|
for did, info := range c.deviceListCache {
|
|
deviceListForStorage[did] = info
|
|
}
|
|
deviceKey := fmt.Sprintf("%s_%s", c.uid, c.cloudServer)
|
|
if err := c.storage.Save("miot_devices", deviceKey, deviceListForStorage); err != nil {
|
|
c.lgr.Warnf("[MIoTClient] Warning: failed to save devices: %v\n", err)
|
|
}
|
|
}
|
|
|
|
// 7. Notify device list changed
|
|
if changed {
|
|
c.ShowDevicesChangedNotify()
|
|
}
|
|
|
|
c.lgr.Infof("[MIoTClient] Devices refreshed: cloud=%d, cache=%d\n", len(cloudDevices), len(c.deviceListCache))
|
|
return nil
|
|
}
|
|
|
|
// ============================================================================
|
|
// Helper functions for device list management
|
|
// ============================================================================
|
|
|
|
// getConnectType extracts the connect_type from device info.
|
|
// Returns 0 if not found.
|
|
func getConnectType(info map[string]interface{}) int {
|
|
if ct, ok := info["connect_type"].(float64); ok {
|
|
return int(ct)
|
|
}
|
|
return 0
|
|
}
|
|
|
|
// getOnline extracts the online status from device info.
|
|
func getOnline(info map[string]interface{}) bool {
|
|
if online, ok := info["online"].(bool); ok {
|
|
return online
|
|
}
|
|
return false
|
|
}
|
|
|
|
// isSharedDevice checks if a device has an owner with a userid (shared device).
|
|
func isSharedDevice(info map[string]interface{}) bool {
|
|
owner, ok := info["owner"].(map[string]interface{})
|
|
if !ok || owner == nil {
|
|
return false
|
|
}
|
|
_, hasUserID := owner["userid"]
|
|
return hasUserID
|
|
}
|
|
|
|
// cloneMap performs a shallow clone of a map.
|
|
func cloneMap(src map[string]interface{}) map[string]interface{} {
|
|
dst := make(map[string]interface{}, len(src))
|
|
for k, v := range src {
|
|
dst[k] = v
|
|
}
|
|
return dst
|
|
}
|
|
|
|
// aggregateDeviceState merges cloud/gateway/lan state for a device into cache.
|
|
// Aligns with Python: __check_device_state.
|
|
func (c *MIoTClient) aggregateDeviceState(did string) {
|
|
// NOTE: caller must hold c.deviceListMu
|
|
|
|
cacheInfo, cacheExists := c.deviceListCache[did]
|
|
if !cacheExists {
|
|
return
|
|
}
|
|
|
|
// OR 聚合:任一来源显式 online=true → 在线;
|
|
// 否则任一来源显式 online=false → 离线;
|
|
// 否则(无来源提供 online 字段)→ 保留缓存旧值,缓存也没有则默认在线。
|
|
online := false
|
|
explicitTrue := false
|
|
explicitFalse := false
|
|
|
|
sources := []map[string]interface{}{
|
|
c.deviceListCloud[did],
|
|
c.deviceListGateway[did],
|
|
c.deviceListLan[did],
|
|
}
|
|
for _, info := range sources {
|
|
if o, ok := info["online"].(bool); ok {
|
|
if o {
|
|
explicitTrue = true
|
|
} else {
|
|
explicitFalse = true
|
|
}
|
|
} else if _, exists := info["online"]; exists {
|
|
c.lgr.Debugf("[MIoTClient] %s online type mismatch: %v (%T)\n", did, info["online"], info["online"])
|
|
}
|
|
}
|
|
|
|
switch {
|
|
case explicitTrue:
|
|
online = true
|
|
case explicitFalse:
|
|
online = false
|
|
default:
|
|
if oldOnline, ok := cacheInfo["online"].(bool); ok {
|
|
online = oldOnline
|
|
} else {
|
|
online = true // 云列表里有但没 online 字段,假设在线
|
|
}
|
|
c.lgr.Debugf("[MIoTClient] %s no online field from cloud/gw/lan, fallback online=%v (cache=%v)\n",
|
|
did, online, cacheInfo["online"])
|
|
}
|
|
|
|
cacheInfo["online"] = online
|
|
}
|
|
|
|
// LoadDevices loads devices from cache or cloud.
|
|
// Aligns with Python: __load_cache_device_async.
|
|
func (c *MIoTClient) LoadDevices() error {
|
|
// Try to load from cache first
|
|
if err := c.loadCacheDevice(); err != nil {
|
|
c.lgr.Errorf("[MIoTClient] Failed to load from cache: %v\n", err)
|
|
}
|
|
|
|
// Then refresh from cloud
|
|
return c.RefreshDevices(c.homeIDs)
|
|
}
|
|
|
|
// GetDeviceInfo returns device metadata for a given DID.
|
|
// Returns nil if device not found.
|
|
// Aligns with Python: device_list.get(did).
|
|
func (c *MIoTClient) GetDeviceInfo(did string) map[string]interface{} {
|
|
c.deviceListMu.RLock()
|
|
defer c.deviceListMu.RUnlock()
|
|
|
|
if info, ok := c.deviceListCache[did]; ok {
|
|
return info
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// GetDeviceInfos returns all device metadata.
|
|
// Aligns with Python: device_list property.
|
|
func (c *MIoTClient) GetDeviceInfos() map[string]map[string]interface{} {
|
|
c.deviceListMu.RLock()
|
|
defer c.deviceListMu.RUnlock()
|
|
|
|
// Return a copy
|
|
result := make(map[string]map[string]interface{})
|
|
for did, info := range c.deviceListCache {
|
|
result[did] = info
|
|
}
|
|
|
|
return result
|
|
}
|
|
|
|
// ScheduleRefreshDevices schedules a device refresh after delay.
|
|
// Aligns with Python: __request_refresh_cloud_devices.
|
|
func (c *MIoTClient) ScheduleRefreshDevices(delay time.Duration) {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
|
|
// Cancel previous timer if exists
|
|
if c.refreshCloudDevicesTimer != nil {
|
|
c.refreshCloudDevicesTimer.Stop()
|
|
}
|
|
|
|
c.refreshCloudDevicesTimer = time.AfterFunc(delay, func() {
|
|
c.lgr.Infof("[MIoTClient] Refreshing devices...\n")
|
|
if err := c.RefreshDevices(c.homeIDs); err != nil {
|
|
c.lgr.Errorf("[MIoTClient] Refresh failed, will retry: %v\n", err)
|
|
// Retry after 60s
|
|
c.ScheduleRefreshDevices(RefreshCloudDevicesRetryDelay * time.Millisecond)
|
|
}
|
|
})
|
|
|
|
c.lgr.Infof("[MIoTClient] Device refresh scheduled in %v\n", delay)
|
|
}
|
|
|
|
// Start starts the MIoTClient and schedules initial device load.
|
|
// Aligns with Python: init_async called at startup.
|
|
func (c *MIoTClient) Start() error {
|
|
// Initialize
|
|
if err := c.Init(); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Schedule initial device load (delayed, aligns with REFRESH_CLOUD_DEVICES_DELAY=6 in Python)
|
|
c.ScheduleRefreshDevices(RefreshCloudDevicesDelay * time.Millisecond)
|
|
|
|
c.lgr.Infof("[MIoTClient] Started\n")
|
|
return nil
|
|
}
|
|
|
|
// Stop stops the MIoTClient and all timers.
|
|
func (c *MIoTClient) Stop() error {
|
|
return c.Deinit()
|
|
}
|
|
|
|
// ============================================================================
|
|
// Device state change handling (aligned with Python miot_client.py)
|
|
// ============================================================================
|
|
|
|
// OnDeviceStateChanged handles device state changes.
|
|
// Aligns with Python: __on_cloud_device_state_changed.
|
|
func (c *MIoTClient) OnDeviceStateChanged(did string, state interface{}) {
|
|
c.deviceListMu.Lock()
|
|
defer c.deviceListMu.Unlock()
|
|
|
|
stateMap, ok := state.(map[string]interface{})
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
// Update or create entry in deviceListCloud (the source triggering this callback)
|
|
if dev, ok := c.deviceListCloud[did]; ok {
|
|
for k, v := range stateMap {
|
|
dev[k] = v
|
|
}
|
|
} else if _, ok := stateMap["online"]; ok {
|
|
c.deviceListCloud[did] = cloneMap(stateMap)
|
|
}
|
|
|
|
// Aggregate from three sources before updating cache
|
|
c.aggregateDeviceState(did)
|
|
|
|
// Propagate to subDeviceState subscribers
|
|
if handler, ok := c.subDeviceState[did]; ok {
|
|
go handler.Handler(did, state)
|
|
}
|
|
}
|