feat: AC 三合一设备支持 + 新增 Plug/Gateway 驱动 + Cover 结构化重构

**AC 驱动增强**
- 支持 thermostat + air-conditioner + air-fresh 多服务三合一设备,属性跨服务合并解析
- 新增 specResolver.findValListInService/findInService/hasService/serviceShortName 方法
- 移除 KindThermostat,heater/electric-blanket/thermostat 重新归类为 KindAirConditioner

**Cover 窗帘重构**
- CoverProps 结构化:Position/Status/Fault 分离
- SetPosition 优先写 target-position,回退 current-position
- OnPropsChanged/FetchProps 支持多属性同步刷新

**新增设备类型**
- Plug 驱动:插座/插头,支持开合、在线状态、属性变更订阅
- Gateway 驱动:网关设备类型识别
- miot/miot_errors.go:新增 MIoT 云端错误码定义

**MIoT 客户端健壮性**
- SetProp 空结果视为失败;设备被移除或离线时自动刷新设备列表
- GetDevices 的 has_more 改用 parseBool 严格解析
- MIPS SubscribeMany 批量订阅减少往返
- matchURN 改为按段精确匹配,避免 tofan 误匹配 fan

**Bridge**
- 移除 Thermostat case,Cover 改用结构化 CoverProps
- 新增设备创建后自动订阅属性变更

**其他**
- 日志整体 Info → Debug 降噪
- .gitignore 忽略 /.codebuddy/memory
- 更新 devices_test 与 08_adv_filter_classify 示例
This commit is contained in:
2026-08-23 19:42:55 +08:00
parent 04333750ab
commit 86c21ce91c
25 changed files with 1129 additions and 155 deletions
+17
View File
@@ -0,0 +1,17 @@
{
"auth_info": {
"access_token": "V3_R0Je3kB-S6Mf1O7kKDPm3EDvgHiwsJMze4mVkPG0rRGHYv0KQEUVUyNRJ-ittoSSOmbNxso1syxYrKjOIAGck6EatzRIJsyULI5maduCo-GUyfHu4A6ZYa8zWD6EXp6pmnLim8KT8i1mDbczyOsFTQhNxWY1tq4CM54xlVhFqQwGS65SVV1DhrXvW9XU3nPa",
"refresh_token": "R3_WZYXokg6Nnmj2Ceib_bcPtI2fr548hNrjCnK5aG_Lm8-8X5PJeK5SHpicsOBy4Z5cWP4672pviFfWORSsSCpH3tF7UXY8Xuib8uyQP4ppxRiPssybXb67VLFfEHa7raQYtVH_2NHkp53flOZKiRtcFHOX4HAZoC7_AKuIdfWcIa4P3k03XmQF1WG0Ccp4BMH12P7eOJBHyU3LScJbFqVLQ",
"mac_key": "0K_213_6BrQPCpmmfolHVZpTR1Y",
"expires_in": 259200,
"expires_ts": 1783011905
},
"uid": "2897121707",
"cloud_server": "cn",
"uuid": "84259d29c352cad945d1f5e3dab3d89b",
"redirect_url": "http://homeassistant.local:8123",
"language": "en",
"ctrl_mode": "auto",
"storage_path": "",
"nick_name": "Xiaomi"
}
+90
View File
@@ -0,0 +1,90 @@
// Command diag_control 直接控制米家设备(纯 HTTP,不初始化 MQTT)。
//
// 控制序列:开空调 → 设温度 → 设模式 → 回读验证。
// 与生产 SDK 同一云端链路(miotspec/prop/set),不占 MQTT client id。
//
// 用法:token/did 改下面常量,在生产环境 go run .(Linux 下 dataDir 填生产路径)。
package main
import (
"fmt"
"os"
"xiaomihome/miot"
)
// ── 配置:按需修改 ──
const (
accessToken = "V3_1GfW88WpawoKACRqwdYBb86In_I5e8r8stZdRPXI1PR4zZgIEm4cbjNwn0wYa6-ZQtIliJJerf-iI0W9Lx8t6fm7Ln16w0RVAmYO3e9XubJ0dK7D3eM_zBlCcVmQCla6Sx0BKQv_YXI9E2zXXopKaRjIUWyx0lb0YcE7oEopLuvfA0sEFoTohVKt543KnHJH" // 米家 OAuth access_token
targetDID = "2188673403" // 目标设备 did
dataDir = "/home/ubuntu/server/zoSmartUser/data" // spec 缓存目录(生产 Linux)
// ──────────────────────
)
func main() {
if accessToken == "" {
fmt.Fprintln(os.Stderr, "请先填写 accessToken")
os.Exit(1)
}
httpClient, err := miot.NewMIoTHttpClient("cn", miot.OAUTH2_CLIENT_ID, accessToken)
if err != nil {
fmt.Fprintf(os.Stderr, "创建 HTTP 客户端失败: %v\n", err)
os.Exit(1)
}
fmt.Printf("HTTP 客户端 OK(未连接 MQTT),设备 %s\n", targetDID)
// 读取当前状态(控制前)
readAll("控制前状态", httpClient)
// ── 控制序列 ──
set("开空调 siid=4/piid=1=true", 4, 1, true, httpClient)
set("设温度 siid=4/piid=2=26", 4, 2, float64(26), httpClient)
set("设模式 siid=2/piid=2=1(Cool)", 2, 2, float64(1), httpClient)
set("关空调 siid=4/piid=1=false", 4, 1, false, httpClient)
// 回读验证(控制后)
readAll("控制后状态", httpClient)
}
// set 写属性并回读验证。
func set(label string, siid, piid int, value interface{}, c *miot.MIoTHttpClient) {
fmt.Printf("\n=== %s ===\n", label)
rs, err := c.SetProp([]miot.PropParam{{DID: targetDID, Siid: siid, Piid: piid, Value: value}})
if err != nil {
fmt.Printf(" 写失败: %v\n", err)
return
}
for _, r := range rs {
fmt.Printf(" 写返回 code=%d\n", r.Code)
}
// 回读
gr, err := c.GetProps([]miot.PropParam{{DID: targetDID, Siid: siid, Piid: piid}})
if err != nil {
fmt.Printf(" 回读失败: %v\n", err)
return
}
if len(gr) > 0 {
fmt.Printf(" 回读 siid=%d/piid=%d = %v\n", siid, piid, gr[0].Value)
}
}
// readAll 读取关键属性。
func readAll(label string, c *miot.MIoTHttpClient) {
fmt.Printf("\n--- %s ---\n", label)
read := func(siid, piid int, name string) {
rs, err := c.GetProps([]miot.PropParam{{DID: targetDID, Siid: siid, Piid: piid}})
if err != nil {
fmt.Printf(" %-28s 错误: %v\n", name, err)
return
}
if len(rs) > 0 {
fmt.Printf(" %-28s = %v\n", name, rs[0].Value)
}
}
read(4, 1, "空调开关 siid=4/piid=1")
read(4, 2, "空调温度 siid=4/piid=2")
read(2, 2, "模式 siid=2/piid=2")
read(2, 4, "地暖温度 siid=2/piid=4")
read(5, 2, "环境温度 siid=5/piid=2")
}
+222
View File
@@ -0,0 +1,222 @@
// Command diag_device 诊断指定米家设备的云端属性(纯 HTTP,不初始化 MQTT)。
//
// 用途:对比 SDK 读到属性值 vs 米家 App 显示,验证 siid/piid 是否正确、
// 控制是否生效、设备连接类型等。可在生产环境旁路运行(不与生产 MQTT 冲突)。
//
// 用法:token/did 直接改下面常量,然后 go run . 即可。
package main
import (
"fmt"
"os"
"sort"
"strconv"
"strings"
"xiaomihome/miot"
)
// ── 配置:按需修改 ──
const (
accessToken = "V3_1GfW88WpawoKACRqwdYBb86In_I5e8r8stZdRPXI1PR4zZgIEm4cbjNwn0wYa6-ZQtIliJJerf-iI0W9Lx8t6fm7Ln16w0RVAmYO3e9XubJ0dK7D3eM_zBlCcVmQCla6Sx0BKQv_YXI9E2zXXopKaRjIUWyx0lb0YcE7oEopLuvfA0sEFoTohVKt543KnHJH" // 米家 OAuth access_token
targetDID = "2188673403" // 目标设备 did
homeID = "619001787902" // 家庭 ID(留空=全部家庭;先跑一遍看家庭列表再填)
dataDir = "/home/ubuntu/server/zoSmartUser/data" // spec 缓存目录(生产 Linux 路径)
setTest = "4,1,false|4,2,26" // 写属性测试:"siid,piid,value",多组用 | 分隔;留空跳过
// ──────────────────────
)
func main() {
if accessToken == "" {
fmt.Fprintln(os.Stderr, "请先在 main.go 顶部 const 里填写 accessToken")
os.Exit(1)
}
// 纯 HTTP 客户端(不连 MQTT,不占 mqtt client id)
httpClient, err := miot.NewMIoTHttpClient("cn", miot.OAUTH2_CLIENT_ID, accessToken)
if err != nil {
fmt.Fprintf(os.Stderr, "创建 HTTP 客户端失败: %v\n", err)
os.Exit(1)
}
fmt.Println("HTTP 客户端 OK(未连接 MQTT)")
// 1. 按 did 精确拉设备信息(不依赖家庭列表)
devMap, err := httpClient.GetDevicesWithDIDs([]string{targetDID})
if err != nil {
fmt.Fprintf(os.Stderr, "获取设备信息失败: %v\n", err)
os.Exit(1)
}
info, ok := devMap[targetDID]
if !ok {
fmt.Fprintf(os.Stderr, "设备 %s 未返回(共 %d 台)\n", targetDID, len(devMap))
os.Exit(1)
}
infoMap := info
fmt.Printf("\n=== 设备信息 %s ===\n", targetDID)
printDeviceInfo(infoMap)
urn, _ := infoMap["urn"].(string)
model, _ := infoMap["model"].(string)
fmt.Printf("URN: %s\nMODEL: %s\n", urn, model)
// 2. 解析 SPEC
parser := miot.NewMIoTSpecParser("en", miot.NewMIoTStorage(dataDir), dataDir+"/specs")
spec, err := parser.Parse(urn, model)
if err != nil {
fmt.Fprintf(os.Stderr, "SPEC 解析失败: %v\n", err)
os.Exit(1)
}
fmt.Printf("\n=== SPEC 服务/属性(siid 布局) ===\n")
for _, svc := range spec.Services {
if svc == nil {
continue
}
fmt.Printf("svc[%d] %s\n", svc.IID, short(svc.Type))
names := make([]string, 0, len(svc.Properties))
for _, p := range svc.Properties {
if p == nil {
continue
}
names = append(names, fmt.Sprintf("%s(%d)", short(p.Type), p.PIID))
}
sort.Strings(names)
fmt.Printf(" %s\n", strings.Join(names, " "))
}
// 3. 批量读取全部可读属性
fmt.Printf("\n=== 读取全部可读属性 ===\n")
var params []miot.PropParam
keys := make([]string, 0)
for _, svc := range spec.Services {
if svc == nil {
continue
}
for _, p := range svc.Properties {
if p == nil {
continue
}
params = append(params, miot.PropParam{DID: targetDID, Siid: svc.IID, Piid: p.PIID})
keys = append(keys, fmt.Sprintf("%d.%d %s", svc.IID, p.PIID, short(p.Type)))
}
}
results, err := httpClient.GetProps(params)
if err != nil {
fmt.Fprintf(os.Stderr, "读取属性失败: %v\n", err)
} else {
for i, r := range results {
name := ""
if i < len(keys) {
name = keys[i]
}
fmt.Printf(" %-24s code=%d value=%v\n", name, r.Code, r.Value)
}
}
// 4. 三合一设备对比读取(空调/地暖双温度)
fmt.Printf("\n=== 温度属性对比(三合一设备) ===\n")
read := func(siid, piid int, label string) {
rs, err := httpClient.GetProps([]miot.PropParam{{DID: targetDID, Siid: siid, Piid: piid}})
if err != nil {
fmt.Printf(" %s: 错误 %v\n", label, err)
return
}
if len(rs) > 0 {
fmt.Printf(" %-28s code=%d value=%v\n", label, rs[0].Code, rs[0].Value)
}
}
read(2, 1, "siid=2/piid=1 on(空调)")
read(2, 2, "siid=2/piid=2 mode")
read(2, 4, "siid=2/piid=4 target-temperature")
read(4, 2, "siid=4/piid=2 target-temperature")
read(5, 2, "siid=5/piid=2 temperature(环境)")
// 5. 写属性测试(可选,多组用 | 分隔,如 "2,4,26|4,2,26")
if setTest != "" {
for _, spec := range strings.Split(setTest, "|") {
doSetAndVerify(httpClient, strings.TrimSpace(spec))
}
}
}
// doSetAndVerify 执行一次写属性并回读验证。
func doSetAndVerify(httpClient *miot.MIoTHttpClient, spec string) {
parts := strings.Split(spec, ",")
if len(parts) < 3 {
fmt.Fprintf(os.Stderr, "setTest 格式错误: %s (应为 siid,piid,value)\n", spec)
return
}
siid, _ := strconv.Atoi(strings.TrimSpace(parts[0]))
piid, _ := strconv.Atoi(strings.TrimSpace(parts[1]))
var val interface{}
if b, err := strconv.ParseBool(parts[2]); err == nil {
val = b
} else if f, err := strconv.ParseFloat(strings.TrimSpace(parts[2]), 64); err == nil {
val = f
} else {
val = strings.TrimSpace(parts[2])
}
fmt.Printf("\n=== 写属性 siid=%d piid=%d value=%v ===\n", siid, piid, val)
rs, err := httpClient.SetProp([]miot.PropParam{{DID: targetDID, Siid: siid, Piid: piid, Value: val}})
if err != nil {
fmt.Printf("写属性失败: %v\n", err)
return
}
for _, r := range rs {
fmt.Printf(" code=%d (0/1=成功,非0=失败)\n", r.Code)
}
// 写后立即回读验证(区分"写错服务" vs "云端假成功")
fmt.Println("--- 写后回读 ---")
verify := func(vsiid, vpiid int, label string) {
vrs, err := httpClient.GetProps([]miot.PropParam{{DID: targetDID, Siid: vsiid, Piid: vpiid}})
if err != nil {
fmt.Printf(" %s: 错误 %v\n", label, err)
return
}
if len(vrs) > 0 {
fmt.Printf(" %-30s code=%d value=%v\n", label, vrs[0].Code, vrs[0].Value)
}
}
verify(siid, piid, "回读 siid=目标")
verify(2, 1, "siid=2/piid=1 on(地暖)")
verify(2, 2, "siid=2/piid=2 mode(地暖)")
verify(2, 4, "siid=2/piid=4 target-temperature(地暖)")
verify(4, 1, "siid=4/piid=1 on(空调)")
verify(4, 2, "siid=4/piid=2 target-temperature(空调)")
}
func printDeviceInfo(m map[string]interface{}) {
fields := []struct{ k, label string }{
{"name", "名称"}, {"model", "型号"}, {"connect_type", "connect_type"},
{"pid", "pid"}, {"online", "online"}, {"isOnline", "isOnline"},
{"rssi", "rssi"}, {"local_ip", "local_ip"}, {"parent_id", "parent_id"},
{"urn", "urn"},
}
for _, f := range fields {
if v, ok := m[f.k]; ok {
fmt.Printf(" %-12s = %v\n", f.label, v)
}
}
}
func short(urn string) string {
parts := strings.Split(urn, ":")
if len(parts) >= 4 {
return parts[3]
}
return urn
}
// toStrSlice 将 interface{}([]interface{})转为 []string。
func toStrSlice(v interface{}) []string {
arr, ok := v.([]interface{})
if !ok {
return nil
}
out := make([]string, 0, len(arr))
for _, item := range arr {
if s, ok := item.(string); ok {
out = append(out, s)
}
}
return out
}
+130
View File
@@ -0,0 +1,130 @@
// Command diag_sdk 用完整 SDK 封装(xiaomi.NewClient + devices.Create)重现生产控制路径。
//
// 目的:验证 SDK 封装层解析出的属性绑定(on/target-temperature 绑到哪个 siid)。
// 与生产同链路:ac.TurnOn / ac.SetTargetTemp → SetProp(绑定siid) → 云端。
//
// 注意:会初始化 MQTT,但用自定义 uuid(client id 唯一,不和生产进程冲突)。
// 在开发机运行;dataDir 指向与生产相同的 spec 缓存(拷一份或指定生产目录)。
package main
import (
"context"
"fmt"
"os"
"time"
"xiaomihome/xiaomi"
"xiaomihome/xiaomi/devices"
)
// ── 配置:按需修改 ──
const (
accessToken = "V3_1GfW88WpawoKACRqwdYBb86In_I5e8r8stZdRPXI1PR4zZgIEm4cbjNwn0wYa6-ZQtIliJJerf-iI0W9Lx8t6fm7Ln16w0RVAmYO3e9XubJ0dK7D3eM_zBlCcVmQCla6Sx0BKQv_YXI9E2zXXopKaRjIUWyx0lb0YcE7oEopLuvfA0sEFoTohVKt543KnHJH" // 米家 OAuth access_token
uid = "2897121707" // 小米帐号 UID
homeID = int64(619001787902) // 家庭 ID
targetDID = "2188673403" // 目标设备 did
dataDir = "D:/xiaomihome/go-xiaomihome/miot/examples/diag_sdk/data" // 本 demo 自己的 spec 缓存目录
// ──────────────────────
)
func main() {
ctx := context.Background()
cfg := xiaomi.Config{
AuthInfo: xiaomi.AuthInfo{AccessToken: accessToken},
UUID: "diag-sdk-" + time.Now().Format("150405"), // 唯一 uuid,不占生产 client id
UID: uid,
DataDir: dataDir,
HomeIDs: []int64{homeID},
}
client, err := xiaomi.NewClient(ctx, cfg)
if err != nil {
fmt.Fprintf(os.Stderr, "NewClient 失败: %v\n", err)
os.Exit(1)
}
defer client.Close()
fmt.Println("SDK 客户端 OK")
info, err := client.GetDevice(ctx, targetDID)
if err != nil || info == nil {
fmt.Fprintf(os.Stderr, "获取设备 %s 失败: %v\n", targetDID, err)
os.Exit(1)
}
fmt.Printf("设备: %s (%s)\nURN: %s\n", info.Name, info.Model, info.URN)
dev, err := devices.NewDevice(client, info)
if err != nil {
fmt.Fprintf(os.Stderr, "创建设备控制失败: %v\n", err)
os.Exit(1)
}
if dev == nil {
fmt.Fprintf(os.Stderr, "设备类型不支持\n")
os.Exit(1)
}
ac, ok := dev.(devices.AirConditioner)
if !ok {
fmt.Fprintf(os.Stderr, "设备不是 AirConditioner(got %T)\n", dev)
os.Exit(1)
}
fmt.Println("AirConditioner 封装 OK")
// 读取当前状态(封装层)
if props, err := ac.FetchProps(ctx); err != nil {
fmt.Printf("FetchProps 失败: %v\n", err)
} else {
fmt.Printf("封装层状态: on=%v mode=%s target=%.0f current=%.0f fan=%s hum=%d\n",
props.On, props.Mode, props.TargetTemperature, props.Temperature,
props.FanLevel, props.CurrentHumidity)
}
// 控制(观察日志里 SetProp 写的 siid)
fmt.Println("\n=== 控制序列(看 SetProp 日志的 siid) ===")
if err := ac.TurnOn(ctx); err != nil {
fmt.Printf("TurnOn 失败: %v\n", err)
} else {
fmt.Println("TurnOn OK")
}
time.Sleep(500 * time.Millisecond)
if err := ac.SetTargetTemp(ctx, 26); err != nil {
fmt.Printf("SetTargetTemp 失败: %v\n", err)
} else {
fmt.Println("SetTargetTemp(26) OK")
}
time.Sleep(500 * time.Millisecond)
if ac.HasFanSpeed() {
if err := ac.SetFanSpeed(ctx, "medium"); err != nil {
fmt.Printf("SetFanSpeed 失败: %v\n", err)
} else {
fmt.Println("SetFanSpeed(medium) OK")
}
} else {
fmt.Println("HasFanSpeed = false")
}
time.Sleep(500 * time.Millisecond)
if err := ac.SetMode(ctx, "cool"); err != nil {
fmt.Printf("SetMode 失败: %v\n", err)
} else {
fmt.Println("SetMode(cool) OK")
}
time.Sleep(500 * time.Millisecond)
if err := ac.TurnOff(ctx); err != nil {
fmt.Printf("TurnOff 失败: %v\n", err)
} else {
fmt.Println("TurnOff OK")
}
time.Sleep(500 * time.Millisecond)
// 控制后回读
if props, err := ac.FetchProps(ctx); err != nil {
fmt.Printf("FetchProps 失败: %v\n", err)
} else {
fmt.Printf("控制后状态: on=%v mode=%s target=%.0f current=%.0f fan=%s hum=%d\n",
props.On, props.Mode, props.TargetTemperature, props.Temperature,
props.FanLevel, props.CurrentHumidity)
}
}
+30 -30
View File
@@ -156,7 +156,7 @@ func (c *MIoTClient) Init() error {
if et, ok := authInfo["expires_ts"].(float64); ok {
c.entryData["expires_ts"] = et
}
c.lgr.Infof("[MIoTClient] Loaded token from storage\n")
c.lgr.Debugf("[MIoTClient] Loaded token from storage\n")
} else if accessToken != "" {
// First run: write entryData token to storage
expiresTS := float64(time.Now().Unix()) + 3600*24*30 // default 30 days
@@ -168,7 +168,7 @@ func (c *MIoTClient) Init() error {
"refresh_token": refreshToken,
"expires_ts": expiresTS,
}, true)
c.lgr.Infof("[MIoTClient] Saved initial token to storage\n")
c.lgr.Debugf("[MIoTClient] Saved initial token to storage\n")
}
}
@@ -187,7 +187,7 @@ func (c *MIoTClient) Init() error {
oauthRedirectURL, _ := c.entryData["oauth_redirect_url"].(string)
uuid, _ := c.entryData["uuid"].(string)
c.oauth = NewMIoTOauthClient(OAUTH2_CLIENT_ID, oauthRedirectURL, c.cloudServer, uuid)
c.lgr.Infof("[MIoTClient] OAuth client initialized\n")
c.lgr.Debugf("[MIoTClient] OAuth client initialized\n")
}
// 3.2 Initialize Cert client
@@ -197,7 +197,7 @@ func (c *MIoTClient) Init() error {
} else {
c.cert = cert
c.cert.VerifyCACert()
c.lgr.Infof("[MIoTClient] Cert client initialized\n")
c.lgr.Debugf("[MIoTClient] Cert client initialized\n")
}
// 4. Initialize HTTP client (use token from storage if available)
@@ -207,7 +207,7 @@ func (c *MIoTClient) Init() error {
c.lgr.Warnf("[MIoTClient] Warning: failed to create HTTP client: %v\n", err)
} else {
c.http = httpClient
c.lgr.Infof("[MIoTClient] HTTP client created\n")
c.lgr.Debugf("[MIoTClient] HTTP client created\n")
}
}
@@ -225,7 +225,7 @@ func (c *MIoTClient) Init() error {
if err := mips.Connect(); err != nil {
c.lgr.Warnf("[MIoTClient] Warning: MIPS cloud connect failed: %v\n", err)
} else {
c.lgr.Infof("[MIoTClient] MIPS cloud connected\n")
c.lgr.Debugf("[MIoTClient] MIPS cloud connected\n")
}
c.mu.Lock()
}
@@ -239,7 +239,7 @@ func (c *MIoTClient) Init() error {
nw.SubNetworkStatus(key, c.onNetworkStatusChanged)
status := nw.GetNetworkStatus()
c.onNetworkStatusChanged(status)
c.lgr.Infof("[MIoTClient] Network status subscribed (online=%v)\n", status)
c.lgr.Debugf("[MIoTClient] Network status subscribed (online=%v)\n", status)
}
}
@@ -258,12 +258,12 @@ func (c *MIoTClient) Init() error {
groupID, _ := info["group_id"].(string)
if groupID != "" {
svc.SubServiceChange(key, groupID, c.onMipsServiceStateChange)
c.lgr.Infof("[MIoTClient] mDNS service subscribed: home=%s, group=%s\n", homeID, groupID)
c.lgr.Debugf("[MIoTClient] mDNS service subscribed: home=%s, group=%s\n", homeID, groupID)
// Check if service already discovered
serviceData := svc.GetServices(groupID)
if sd, exists := serviceData[groupID]; exists && sd.validService() {
c.lgr.Infof("[MIoTClient] Central mips service scanned: %s, %v\n", homeID, sd)
c.lgr.Debugf("[MIoTClient] Central mips service scanned: %s, %v\n", homeID, sd)
mips := NewMipsLocalClient(
getVirtualDID(c.entryData),
sd.Addresses[0],
@@ -294,7 +294,7 @@ func (c *MIoTClient) Init() error {
if lan.InitDone() {
c.onMiotLanStateChange(true)
}
c.lgr.Infof("[MIoTClient] LAN control initialized\n")
c.lgr.Debugf("[MIoTClient] LAN control initialized\n")
}
}
} else {
@@ -324,7 +324,7 @@ func (c *MIoTClient) Init() error {
}
c.i18n = NewMIoTI18n(lang)
c.i18n.Init()
c.lgr.Infof("[MIoTClient] i18n initialized (lang=%s)\n", lang)
c.lgr.Debugf("[MIoTClient] i18n initialized (lang=%s)\n", lang)
c.lgr.Infof("[MIoTClient] Initialized\n")
return nil
@@ -358,7 +358,7 @@ func (c *MIoTClient) Deinit() error {
if c.mipsCloud != nil {
c.mipsCloud.UnsubState(key)
c.mipsCloud.Disconnect()
c.lgr.Infof("[MIoTClient] MIPS cloud disconnected\n")
c.lgr.Debugf("[MIoTClient] MIPS cloud disconnected\n")
}
// 4. Cancel refresh cloud devices timer
@@ -493,7 +493,7 @@ func (c *MIoTClient) loadCacheDevice() error {
c.deviceListCache[did] = infoMap
}
}
c.lgr.Infof("[MIoTClient] Loaded %d devices from cache\n", len(deviceList))
c.lgr.Debugf("[MIoTClient] Loaded %d devices from cache\n", len(deviceList))
return nil
}
@@ -647,7 +647,7 @@ func (c *MIoTClient) RefreshAuthInfo() (map[string]interface{}, error) {
c.mipsCloud.UpdateAccessToken(ac)
}
c.lgr.Infof("[MIoTClient] RefreshAuthInfo: token refreshed\n")
c.lgr.Debugf("[MIoTClient] RefreshAuthInfo: token refreshed\n")
return result, nil
}
@@ -690,7 +690,7 @@ func (c *MIoTClient) refreshUserCert() {
nextDelay = 60
}
c.RequestRefreshUserCert(nextDelay)
c.lgr.Infof("[MIoTClient] Cert still valid for %ds, recheck in %ds\n", remainSec, nextDelay)
c.lgr.Debugf("[MIoTClient] Cert still valid for %ds, recheck in %ds\n", remainSec, nextDelay)
return
}
@@ -720,7 +720,7 @@ func (c *MIoTClient) refreshUserCert() {
}
// Save cert
if err := c.cert.SaveUserCert(certPEM); err != nil {
if err = c.cert.SaveUserCert(certPEM); err != nil {
c.lgr.Errorf("[MIoTClient] Cert refresh: save cert FAILED: %v\n", err)
return
}
@@ -783,13 +783,13 @@ func (c *MIoTClient) ShowCentralStateChangedNotify(connected bool) {
// Aligns with Python: __on_mips_cloud_state_changed.
func (c *MIoTClient) onMipsCloudStateChanged(key string, connected bool) {
if connected {
c.lgr.Infof("[MIoTClient] MIPS cloud connected\n")
c.lgr.Debugf("[MIoTClient] MIPS cloud connected\n")
// Refresh devices on reconnect
c.ScheduleRefreshDevices(RefreshCloudDevicesDelay * time.Millisecond)
return
}
c.lgr.Infof("[MIoTClient] MIPS cloud disconnected\n")
c.lgr.Debugf("[MIoTClient] MIPS cloud disconnected\n")
// Disconnect: mark all cloud devices offline
c.deviceListMu.Lock()
@@ -821,7 +821,7 @@ func (c *MIoTClient) onMipsCloudStateChanged(key string, connected bool) {
// onNetworkStatusChanged handles network connectivity changes.
// Aligns with Python: __on_network_status_changed.
func (c *MIoTClient) onNetworkStatusChanged(status bool) {
c.lgr.Infof("[MIoTClient] Network status changed: online=%v\n", status)
c.lgr.Debugf("[MIoTClient] Network status changed: online=%v\n", status)
if status {
// Network is back; trigger device refresh and reconnect MIPS
if c.mipsCloud != nil && !c.mipsCloud.IsConnected() {
@@ -834,15 +834,15 @@ func (c *MIoTClient) onNetworkStatusChanged(status bool) {
// onMipsServiceStateChange handles mDNS service discoveries for central gateways.
// Aligns with Python: __on_mips_service_state_change.
func (c *MIoTClient) onMipsServiceStateChange(groupID string, state MipsServiceState, data map[string]interface{}) {
c.lgr.Infof("[MIoTClient] MIPS service change: group=%s, state=%s\n", groupID, state)
c.lgr.Debugf("[MIoTClient] MIPS service change: group=%s, state=%s\n", groupID, state)
virtualDID := getVirtualDID(c.entryData)
if virtualDID == "" {
c.lgr.Infof("[MIoTClient] MIPS service: missing virtual_did, skip\n")
c.lgr.Debugf("[MIoTClient] MIPS service: missing virtual_did, skip\n")
return
}
if c.cert == nil {
c.lgr.Infof("[MIoTClient] MIPS service: cert not initialized, skip\n")
c.lgr.Debugf("[MIoTClient] MIPS service: cert not initialized, skip\n")
return
}
@@ -851,7 +851,7 @@ func (c *MIoTClient) onMipsServiceStateChange(groupID string, state MipsServiceS
// Create new MipsLocalClient
host, _ := data["addresses"].([]interface{})
if len(host) == 0 {
c.lgr.Infof("[MIoTClient] MIPS service: no addresses for %s\n", groupID)
c.lgr.Debugf("[MIoTClient] MIPS service: no addresses for %s\n", groupID)
return
}
hostStr, _ := host[0].(string)
@@ -872,7 +872,7 @@ func (c *MIoTClient) onMipsServiceStateChange(groupID string, state MipsServiceS
mips.SubState(groupID, c.onMipsLocalStateChanged)
mips.Connect()
c.mipsLocal[groupID] = mips
c.lgr.Infof("[MIoTClient] MIPS local client created: %s @ %s:%d\n", groupID, hostStr, int(port))
c.lgr.Debugf("[MIoTClient] MIPS local client created: %s @ %s:%d\n", groupID, hostStr, int(port))
case MipsServiceUpdated:
// Disconnect old, connect new
@@ -902,7 +902,7 @@ func (c *MIoTClient) onMipsServiceStateChange(groupID string, state MipsServiceS
mips.SubState(groupID, c.onMipsLocalStateChanged)
mips.Connect()
c.mipsLocal[groupID] = mips
c.lgr.Infof("[MIoTClient] MIPS local client updated: %s @ %s:%d\n", groupID, hostStr, int(port))
c.lgr.Debugf("[MIoTClient] MIPS local client updated: %s @ %s:%d\n", groupID, hostStr, int(port))
case MipsServiceRemoved:
if mips, exists := c.mipsLocal[groupID]; exists {
@@ -910,7 +910,7 @@ func (c *MIoTClient) onMipsServiceStateChange(groupID string, state MipsServiceS
mips.UnsubState(groupID)
mips.Disconnect()
delete(c.mipsLocal, groupID)
c.lgr.Infof("[MIoTClient] MIPS local client removed: %s\n", groupID)
c.lgr.Debugf("[MIoTClient] MIPS local client removed: %s\n", groupID)
}
}
}
@@ -919,7 +919,7 @@ func (c *MIoTClient) onMipsServiceStateChange(groupID string, state MipsServiceS
// Aligns with Python: __on_mips_local_state_changed.
// key is the groupID of the gateway.
func (c *MIoTClient) onMipsLocalStateChanged(key string, connected bool) {
c.lgr.Infof("[MIoTClient] MIPS local state changed: group=%s, connected=%v\n", key, connected)
c.lgr.Debugf("[MIoTClient] MIPS local state changed: group=%s, connected=%v\n", key, connected)
if connected {
// Reconnect: pull gateway device list
@@ -966,7 +966,7 @@ func (c *MIoTClient) onMipsLocalStateChanged(key string, connected bool) {
// onMiotLanStateChange handles LAN controller state transitions.
// Aligns with Python: __on_miot_lan_state_change.
func (c *MIoTClient) onMiotLanStateChange(initialized bool) {
c.lgr.Infof("[MIoTClient] LAN state changed: initialized=%v\n", initialized)
c.lgr.Debugf("[MIoTClient] LAN state changed: initialized=%v\n", initialized)
if initialized {
c.deviceListMu.Lock()
lan, ok := c.miotLan.(*MIoTLan)
@@ -1026,7 +1026,7 @@ func (c *MIoTClient) onLanDeviceStateChanged(did string, state map[string]interf
// Aligns with Python: __on_gw_device_list_changed.
// Called when a MipsLocalClient reports that its device list has changed.
func (c *MIoTClient) onGWDeviceListChanged() {
c.lgr.Infof("[MIoTClient] GW device list changed\n")
c.lgr.Debugf("[MIoTClient] GW device list changed\n")
for groupID, mips := range c.mipsLocal {
if mips == nil || !mips.IsConnected() {
@@ -1090,7 +1090,7 @@ func (c *MIoTClient) updateDevicesFromGW(groupID string, devices map[string]inte
}
if changed {
c.lgr.Infof("[MIoTClient] GW devices updated: group=%s, count=%d\n", groupID, len(devices))
c.lgr.Debugf("[MIoTClient] GW devices updated: group=%s, count=%d\n", groupID, len(devices))
}
}
+19 -1
View File
@@ -31,8 +31,26 @@ func (c *MIoTClient) SetProp(did string, siid, piid int, value interface{}) erro
return err
}
// Aligns with Python: 空 result 视为失败(Python 会 raise device_exec_error)
if len(results) == 0 {
return fmt.Errorf("set property failed: empty result (device may be offline)")
}
// code=0 或 code=1 均视为成功(0=读取成功,1=设置成功)
if len(results) > 0 && results[0].Code != 0 && results[0].Code != 1 {
if results[0].Code != 0 && results[0].Code != 1 {
// 设备移除或离线:刷新设备列表(Aligns with Python)
if results[0].Code == -704010000 || results[0].Code == -704042011 {
c.lgr.Errorf("[MIoTClient] device may be removed or offline, %s\n", did)
if devs, rerr := c.http.GetDevices(nil); rerr == nil {
c.deviceListMu.Lock()
for k := range c.deviceListCloud {
if _, ok := devs[k]; !ok {
delete(c.deviceListCloud, k)
}
}
c.deviceListMu.Unlock()
}
}
return fmt.Errorf("set property failed, code=%d", results[0].Code)
}
+1 -1
View File
@@ -155,7 +155,7 @@ func (c *MIoTClient) updateDeviceMsgSub(did string) {
// Update subscription source
c.subSourceList[did] = newSource
c.lgr.Infof("[MIoTClient] updateDeviceMsgSub: %s source %s → %s (gw=%v/%v, lan=%v/%v, cloud=%v)\n",
c.lgr.Debugf("[MIoTClient] updateDeviceMsgSub: %s source %s → %s (gw=%v/%v, lan=%v/%v, cloud=%v)\n",
did, oldSource, newSource, gwOnline, gwPushAvailable, lanOnline, lanPushAvailable, cloudOnline)
}
+57 -36
View File
@@ -477,7 +477,7 @@ func (c *MIoTHttpClient) getDevRoomPage(maxID string) (map[string]map[string]int
}
// Check for more pages
hasMore, _ := result["has_more"].(bool)
hasMore := parseBool(result["has_more"])
nextMaxID, _ := result["max_id"].(string)
if hasMore && nextMaxID != "" {
nextList, err := c.getDevRoomPage(nextMaxID)
@@ -582,42 +582,43 @@ func (c *MIoTHttpClient) GetHomeInfos() (map[string]interface{}, error) {
}
// Handle pagination
hasMore, _ := result["has_more"].(bool)
hasMore := parseBool(result["has_more"])
maxID, _ := result["max_id"].(string)
if hasMore && maxID != "" {
moreList, err := c.getDevRoomPage(maxID)
if err == nil {
for _, deviceSource := range []string{"homelist", "share_home_list"} {
for hidStr, info := range moreList {
if _, exists := homeInfos[deviceSource][hidStr]; !exists {
continue
}
h, _ := homeInfos[deviceSource][hidStr].(map[string]interface{})
if h == nil {
continue
}
h["dids"] = append(toStringSlice(h["dids"]), toStringSlice(info["dids"])...)
if err != nil {
return nil, err
}
for _, deviceSource := range []string{"homelist", "share_home_list"} {
for hidStr, info := range moreList {
if _, exists := homeInfos[deviceSource][hidStr]; !exists {
continue
}
h, _ := homeInfos[deviceSource][hidStr].(map[string]interface{})
if h == nil {
continue
}
h["dids"] = append(toStringSlice(h["dids"]), toStringSlice(info["dids"])...)
for rid, rinfo := range info["room_info"].(map[string]map[string]interface{}) {
hRoomInfo, _ := h["room_info"].(map[string]interface{})
if hRoomInfo == nil {
hRoomInfo = make(map[string]interface{})
h["room_info"] = hRoomInfo
}
if hRoomInfo[rid] == nil {
hRoomInfo[rid] = map[string]interface{}{
"room_id": rid,
"room_name": "",
"dids": []string{},
}
}
hrEntry, _ := hRoomInfo[rid].(map[string]interface{})
if hrEntry != nil {
hrEntry["dids"] = append(
toStringSlice(hrEntry["dids"]),
toStringSlice(rinfo["dids"])...)
for rid, rinfo := range info["room_info"].(map[string]map[string]interface{}) {
hRoomInfo, _ := h["room_info"].(map[string]interface{})
if hRoomInfo == nil {
hRoomInfo = make(map[string]interface{})
h["room_info"] = hRoomInfo
}
if hRoomInfo[rid] == nil {
hRoomInfo[rid] = map[string]interface{}{
"room_id": rid,
"room_name": "",
"dids": []string{},
}
}
hrEntry, _ := hRoomInfo[rid].(map[string]interface{})
if hrEntry != nil {
hrEntry["dids"] = append(
toStringSlice(hrEntry["dids"]),
toStringSlice(rinfo["dids"])...)
}
}
}
}
@@ -725,13 +726,14 @@ func (c *MIoTHttpClient) getDeviceListPage(dids []string, startDID string) (map[
// Handle pagination
nextStartDID, _ := result["next_start_did"].(string)
hasMore, _ := result["has_more"].(bool)
hasMore := parseBool(result["has_more"])
if hasMore && nextStartDID != "" {
nextInfos, err := c.getDeviceListPage(dids, nextStartDID)
if err == nil {
for k, v := range nextInfos {
deviceInfos[k] = v
}
if err != nil {
return nil, err
}
for k, v := range nextInfos {
deviceInfos[k] = v
}
}
@@ -1196,3 +1198,22 @@ func containsString(slice []string, val string) bool {
}
return false
}
// parseBool 兼容云端多种布尔表达(bool / 0|1 / "true"|"false")。
func parseBool(v interface{}) bool {
switch n := v.(type) {
case bool:
return n
case float64:
return n != 0
case float32:
return n != 0
case int:
return n != 0
case int64:
return n != 0
case string:
return n == "true" || n == "1"
}
return false
}
+142
View File
@@ -0,0 +1,142 @@
package miot
// MIoT 云端错误码
// 来源: https://iot.mi.com/v2/new/doc/help-support/faq/spec/error-code
const (
ErrCodeOK = -702000000 // OK
ErrCodeAccept = -702010000 // accept
ErrCodeProcessing = -702022036 // 操作正在处理中
ErrCodeBadRequest = -704000000 // 错误的请求
ErrCodeBadBody = -704000001 // 错误的请求体
ErrCodeScheduleOnlyMi = -704001000 // 定时只支持小米设备
ErrCodeDeviceError = -704002000 // 设备错误(通用)
ErrCodeUnauthenticated = -704010000 // 未认证
ErrCodeTokenExpired = -704012901 // token 不存在或过期
ErrCodeTokenInvalid = -704012902 // token 非法
ErrCodeAuthExpired = -704012903 // 授权过期
ErrCodeNotAuthToXiaoAi = -704012904 // 设备未授权控制能力给小爱
ErrCodeDeviceUnbound = -704012905 // 设备未绑定
ErrCodeAuthFailed = -704012906 // 认证失败
ErrCodeIRNotSupported = -704013101 // 红外设备不支持此操作
ErrCodePropertyNotRead = -704030013 // 属性不可读
ErrCodePropertyNotWrite = -704030023 // 属性不可写
ErrCodePropertyNotNotify = -704030033 // 属性不可上报
ErrCodeTooFrequent = -704030992 // 请求过于频繁
ErrCodeServiceNotFound = -704040002 // 服务不存在
ErrCodePropertyNotFound = -704040003 // 属性不存在
ErrCodeEventNotFound = -704040004 // 事件不存在
ErrCodeMethodNotFound = -704040005 // 方法不存在
ErrCodeFeatureOffline = -704040999 // 功能未上线
ErrCodeCloudNotFound = -704041007 // cloud not found
ErrCodeDeviceNotFound2 = -704042001 // 未找到设备
ErrCodeSceneNotFound = -704042009 // 未找到场景
ErrCodeSceneError = -704042010 // 触发场景异常
ErrCodeDeviceOffline = -704042011 // 设备离线
ErrCodeSceneNoPermission = -704042012 // 场景无权限
ErrCodeSpecNotFound = -704044006 // 未找到功能定义
ErrCodeCannotExecute = -704053100 // 无法执行此操作
ErrCodeCameraSleeping = -704053101 // 摄像机休眠中
ErrCodeTimeout = -704083036 // 操作超时
ErrCodeDeviceNotFound = -704090001 // 未找到设备
ErrCodeInvalidID = -704220008 // 非法的 ID
ErrCodeInvalidUID = -704220009 // 非法的 uid
ErrCodeInvalidSubID = -704220010 // 非法的订阅 ID
ErrCodeActionParamCount = -704220025 // 方法输入参数数量不匹配
ErrCodeActionParamError = -704220035 // 方法输入参数错误
ErrCodeInvalidPropValue = -704220043 // 属性值不正确
ErrCodeEventParamCount = -704222034 // 事件参数数量不匹配
ErrCodeActionOutParamError = -704222035 // 方法输出参数错误
ErrCodeOpenPlatformError = -705001000 // 开放平台服务器内部错误
ErrCodeSpecServerError = -705004000 // MIoT Spec 服务器内部错误
ErrCodeSubIDEncodeError = -705005000 // 订阅 ID 编码错误
ErrCodeNotifyError = -705006000 // notify error
ErrCodeReadPropFailed = -705201013 // 读属性失败
ErrCodeActionFailed = -705201015 // 方法执行失败
ErrCodeWritePropFailed = -705201023 // 写属性失败
ErrCodeReportPropFailed = -705201033 // 上报属性失败
ErrCodeInvalidSpec = -705204006 // 非法的功能定义
ErrCodeInstanceNotLoaded = -705204007 // 实例未加载
ErrCodeRemoteServiceError = -706010002 // 远程服务异常
ErrCodeMIoTError = -706010004 // MIoT 错误
ErrCodePropCacheFailed = -706010005 // 属性缓存失败
ErrCodeThirdPartyError = -706012000 // 三方云服务器内部错误
ErrCodeThirdReadPropFailed = -706012013 // 三方云读属性失败
ErrCodeThirdActionFailed = -706012015 // 三方云方法执行失败
ErrCodeThirdWritePropFailed = -706012023 // 三方云写属性失败
ErrCodeThirdReportPropFailed = -706012033 // 三方云上报属性失败
ErrCodeThirdSubFailed = -706012043 // 三方云订阅失败
)
// errDescriptions 错误码说明
var errDescriptions = map[int]string{
ErrCodeOK: "OK",
ErrCodeAccept: "accept",
ErrCodeProcessing: "操作正在处理中",
ErrCodeBadRequest: "错误的请求",
ErrCodeBadBody: "错误的请求体",
ErrCodeScheduleOnlyMi: "定时只支持小米设备",
ErrCodeDeviceError: "设备错误(通用)",
ErrCodeUnauthenticated: "未认证",
ErrCodeTokenExpired: "token 不存在或过期",
ErrCodeTokenInvalid: "token 非法",
ErrCodeAuthExpired: "授权过期",
ErrCodeNotAuthToXiaoAi: "设备未授权控制能力给小爱",
ErrCodeDeviceUnbound: "设备未绑定",
ErrCodeAuthFailed: "认证失败",
ErrCodeIRNotSupported: "红外设备不支持此操作",
ErrCodePropertyNotRead: "属性不可读",
ErrCodePropertyNotWrite: "属性不可写",
ErrCodePropertyNotNotify: "属性不可上报",
ErrCodeTooFrequent: "请求过于频繁",
ErrCodeServiceNotFound: "服务不存在",
ErrCodePropertyNotFound: "属性不存在",
ErrCodeEventNotFound: "事件不存在",
ErrCodeMethodNotFound: "方法不存在",
ErrCodeFeatureOffline: "功能未上线",
ErrCodeCloudNotFound: "cloud not found",
ErrCodeDeviceNotFound2: "未找到设备",
ErrCodeSceneNotFound: "未找到场景",
ErrCodeSceneError: "触发场景异常",
ErrCodeDeviceOffline: "设备离线",
ErrCodeSceneNoPermission: "场景无权限",
ErrCodeSpecNotFound: "未找到功能定义",
ErrCodeCannotExecute: "无法执行此操作",
ErrCodeCameraSleeping: "摄像机休眠中",
ErrCodeTimeout: "操作超时",
ErrCodeDeviceNotFound: "未找到设备",
ErrCodeInvalidID: "非法的 ID",
ErrCodeInvalidUID: "非法的 uid",
ErrCodeInvalidSubID: "非法的订阅 ID",
ErrCodeActionParamCount: "方法输入参数数量不匹配",
ErrCodeActionParamError: "方法输入参数错误",
ErrCodeInvalidPropValue: "属性值不正确",
ErrCodeEventParamCount: "事件参数数量不匹配",
ErrCodeActionOutParamError: "方法输出参数错误",
ErrCodeOpenPlatformError: "开放平台服务器内部错误",
ErrCodeSpecServerError: "MIoT Spec 服务器内部错误",
ErrCodeSubIDEncodeError: "订阅 ID 编码错误",
ErrCodeNotifyError: "notify error",
ErrCodeReadPropFailed: "读属性失败",
ErrCodeActionFailed: "方法执行失败",
ErrCodeWritePropFailed: "写属性失败",
ErrCodeReportPropFailed: "上报属性失败",
ErrCodeInvalidSpec: "非法的功能定义",
ErrCodeInstanceNotLoaded: "实例未加载",
ErrCodeRemoteServiceError: "远程服务异常",
ErrCodeMIoTError: "MIoT 错误",
ErrCodePropCacheFailed: "属性缓存失败",
ErrCodeThirdPartyError: "三方云服务器内部错误",
ErrCodeThirdReadPropFailed: "三方云读属性失败",
ErrCodeThirdActionFailed: "三方云方法执行失败",
ErrCodeThirdWritePropFailed: "三方云写属性失败",
ErrCodeThirdReportPropFailed: "三方云上报属性失败",
ErrCodeThirdSubFailed: "三方云订阅失败",
}
// ErrDesc 返回错误码的中文说明。
func ErrDesc(code int) string {
if s, ok := errDescriptions[code]; ok {
return s
}
return "未知错误"
}
+43 -19
View File
@@ -25,7 +25,6 @@ const (
mipsQoS = 2
mipsReconnectIntervalMin = 10.0 // seconds
mipsReconnectIntervalMax = 600.0 // seconds
mipsSubPatch = 300
mipsSubInterval = 1 // second
mipsMQTTInterval = 1 // second
uint32max = 0xFFFFFFFF
@@ -269,6 +268,7 @@ type MQTTClient interface {
Disconnect()
IsConnected() bool
Subscribe(topic string, qos byte, handler func(topic string, payload []byte)) error
SubscribeMany(topics map[string]byte) error
Unsubscribe(topic string) error
Publish(topic string, qos byte, payload []byte) error
SetCredentials(username, password string)
@@ -528,34 +528,42 @@ func (c *mipsClient) flushPendingSubs() {
return
}
// 收集所有待订阅 topic,一次 SUBSCRIBE 包批量订阅(大幅减少网络往返)
c.pendingSubMu.Lock()
subbed := 0
batch := make(map[string]byte, len(c.pendingSubs))
for topic, count := range c.pendingSubs {
if subbed >= mipsSubPatch {
break
}
if count > 3 {
delete(c.pendingSubs, topic)
c.lgr.Debugf("[mips] retry sub exceeded: %s", topic)
continue
}
subbed++
if err := c.mqttConn.Subscribe(topic, byte(mipsQoS), func(t string, p []byte) {
c.onMQTTMessage(t, p)
}); err == nil {
delete(c.pendingSubs, topic)
c.lgr.Debugf("[mips] sub success: %s", topic)
} else {
c.pendingSubs[topic] = count + 1
c.lgr.Debugf("[mips] retry sub %d: %s %v", count, topic, err)
}
batch[topic] = byte(mipsQoS)
}
hasMore := len(c.pendingSubs) > 0
c.pendingSubMu.Unlock()
if hasMore {
if len(batch) == 0 {
return
}
if err := c.mqttConn.SubscribeMany(batch); err == nil {
c.pendingSubMu.Lock()
for t := range batch {
delete(c.pendingSubs, t)
c.lgr.Debugf("[mips] sub success: %s", t)
}
hasMore := len(c.pendingSubs) > 0
c.pendingSubMu.Unlock()
if hasMore {
time.Sleep(time.Duration(mipsSubInterval) * time.Second)
c.flushPendingSubs()
}
} else {
c.lgr.Debugf("[mips] batch sub error: %v", err)
c.pendingSubMu.Lock()
for t := range batch {
c.pendingSubs[t] = c.pendingSubs[t] + 1
}
c.pendingSubMu.Unlock()
time.Sleep(time.Duration(mipsSubInterval) * time.Second)
c.flushPendingSubs()
}
@@ -1042,6 +1050,22 @@ func (p *pahoMQTTClient) Subscribe(topic string, qos byte, handler func(topic st
return nil
}
// SubscribeMany 用单个 SUBSCRIBE 包批量订阅多个 topic,显著减少网络往返。
// 消息分发仍走 DefaultPublishHandler,与 Subscribe 行为一致。
func (p *pahoMQTTClient) SubscribeMany(topics map[string]byte) error {
if p.client == nil || !p.client.IsConnected() {
return fmt.Errorf("not connected")
}
if len(topics) == 0 {
return nil
}
token := p.client.SubscribeMultiple(topics, nil)
if token.Wait() && token.Error() != nil {
return token.Error()
}
return nil
}
func (p *pahoMQTTClient) Unsubscribe(topic string) error {
if p.client == nil || !p.client.IsConnected() {
return fmt.Errorf("not connected")
-3
View File
@@ -183,9 +183,6 @@ func TestMipsConstants(t *testing.T) {
if mipsQoS != 2 {
t.Errorf("mipsQoS = %d, want 2", mipsQoS)
}
if mipsSubPatch != 300 {
t.Errorf("mipsSubPatch = %d, want 300", mipsSubPatch)
}
if MipsRequestTimeoutDefault != 10000 {
t.Errorf("MipsRequestTimeoutDefault = %d, want 10000", MipsRequestTimeoutDefault)
}