在米家 App 删除设备后设备列表仍能获取到:删除检测只 delete(deviceListCloud), deviceListCache 仅置 online=false,而对外 DeviceList()/GetDevices() 读的正是 cache;缓存又会在停机时落盘到 .dict、启动时原样恢复,已删设备跨重启复活。 - pruneRemovedDevicesLocked:候选集改为「缓存 ∪ 云列表」,云端不存在且网关/ 局域网非在线的设备从四张表彻底删除(旧实现只遍历 deviceListCloud,启动时它 为空,停机期间被删的设备永远检测不到) - 护栏:云端返回空列表时跳过剔除,避免接口异常清库 - homeScope:只刷新部分家庭时不再误删其它家庭的云列表 - 子设备(xxx.s1)不参与判定,它被归并到父设备的 sub_devices - 命中剔除的设备同时清理 c.devices 实例与 MQTT 订阅路由 - bridge: pruneControllers 剔除已删设备的控制器,AC()/Switch() 不再返回实例 - 新增 7 个 miot 用例(含 httptest 打桩的端到端)+ 3 个 bridge 用例
527 lines
16 KiB
Go
527 lines
16 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
|
||
}
|
||
}
|
||
}
|
||
|
||
// 剔除检测的家庭作用域:非空表示本次只刷新了部分家庭(GetDevices 只返回这些家庭),
|
||
// 未命中的设备一律不动,避免 RefreshDevices(homeA) 把其它家庭的设备误剔除。
|
||
homeScope := make(map[string]struct{}, len(homeIDs))
|
||
for _, hid := range homeIDs {
|
||
if hid != "" {
|
||
homeScope[hid] = struct{}{}
|
||
}
|
||
}
|
||
|
||
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:云端已不存在的设备要从三源 + 缓存 + 实例 + 订阅里彻底剔除
|
||
removed := c.pruneRemovedDevicesLocked(cloudDevices, homeScope)
|
||
if len(removed) > 0 {
|
||
// 设备实例由 c.mu 保护;此处持 deviceListMu,与下方 updateDeviceMsgSub 的
|
||
// 加锁顺序保持一致(deviceListMu → c.mu)
|
||
c.mu.Lock()
|
||
for _, did := range removed {
|
||
delete(c.devices, did)
|
||
c.lgr.Infof("[MIoTClient] Device removed: %s\n", did)
|
||
}
|
||
c.mu.Unlock()
|
||
// 取消该设备在云/网关/局域网的订阅(三源均已清空 → 内部会退订)
|
||
for _, did := range removed {
|
||
c.updateDeviceMsgSub(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
|
||
}
|
||
|
||
// deviceOnlineFlag reports whether info explicitly carries online=true.
|
||
func deviceOnlineFlag(info map[string]interface{}) bool {
|
||
if info == nil {
|
||
return false
|
||
}
|
||
online, _ := info["online"].(bool)
|
||
return online
|
||
}
|
||
|
||
// deviceHomeID returns the non-empty home_id recorded for did across all sources.
|
||
func (c *MIoTClient) deviceHomeID(did string) string {
|
||
for _, m := range []map[string]map[string]interface{}{
|
||
c.deviceListCache, c.deviceListCloud, c.deviceListGateway, c.deviceListLan,
|
||
} {
|
||
if info, ok := m[did]; ok {
|
||
if home, _ := info["home_id"].(string); home != "" {
|
||
return home
|
||
}
|
||
}
|
||
}
|
||
return ""
|
||
}
|
||
|
||
// pruneRemovedDevicesLocked 剔除「云端列表中已不存在」的设备,返回被剔除的 DID 列表。
|
||
// NOTE: caller must hold c.deviceListMu.
|
||
//
|
||
// 与上游 Python 的差异(有意为之):Python 只把设备置为离线、永不从
|
||
// _device_list_cache 删除,那是为 Home Assistant 的设备注册表服务的。本库的使用方
|
||
// 需要「设备列表 == 米家当前设备」,因此在米家删除的设备必须从缓存中真正移除,
|
||
// 否则 DeviceList()/GetDevices() 会一直返回已删除的设备(且随 .dict 跨重启复活)。
|
||
//
|
||
// 判定规则:
|
||
// 1. 云端返回空列表 → 视为接口异常(网络抖动、鉴权异常),不剔除任何设备;
|
||
// 2. 候选集 = 缓存 ∪ 云列表:必须遍历缓存,否则启动时 deviceListCloud 尚为空,
|
||
// 会漏掉「服务停机期间在米家被删除」的设备;
|
||
// 3. 网关或局域网仍报告在线 → 保留(云端列表可能还没同步到);
|
||
// 4. 子设备 DID(xxx.s1)不参与判定(它被归并到父设备的 sub_devices 中);
|
||
// 5. homeScope 非空时只处理该家庭下的设备。
|
||
func (c *MIoTClient) pruneRemovedDevicesLocked(cloudDevices map[string]map[string]interface{}, homeScope map[string]struct{}) []string {
|
||
if len(cloudDevices) == 0 {
|
||
c.lgr.Warnf("[MIoTClient] empty cloud device list, skip removal detection\n")
|
||
return nil
|
||
}
|
||
|
||
candidates := make(map[string]struct{}, len(c.deviceListCache)+len(c.deviceListCloud))
|
||
for did := range c.deviceListCache {
|
||
candidates[did] = struct{}{}
|
||
}
|
||
for did := range c.deviceListCloud {
|
||
candidates[did] = struct{}{}
|
||
}
|
||
|
||
var removed []string
|
||
for did := range candidates {
|
||
if _, ok := cloudDevices[did]; ok {
|
||
continue
|
||
}
|
||
// 子设备(xxx.s1)在 RefreshDevices 中被归并进父设备的 sub_devices,
|
||
// 本来就不会出现在 cloudDevices 里,不能据此判定为已删除
|
||
if subDevPattern.MatchString(did) {
|
||
continue
|
||
}
|
||
if deviceOnlineFlag(c.deviceListGateway[did]) || deviceOnlineFlag(c.deviceListLan[did]) {
|
||
continue
|
||
}
|
||
if len(homeScope) > 0 {
|
||
if _, inScope := homeScope[c.deviceHomeID(did)]; !inScope {
|
||
continue
|
||
}
|
||
}
|
||
delete(c.deviceListCloud, did)
|
||
delete(c.deviceListCache, did)
|
||
delete(c.deviceListGateway, did)
|
||
delete(c.deviceListLan, did)
|
||
removed = append(removed, did)
|
||
}
|
||
return removed
|
||
}
|
||
|
||
// 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)
|
||
}
|
||
}
|