feat: 完善 AC"set 后回读"机制 — 目标聚合 + 防抖重试验证

**回读目标聚合(recordReadback)**
- 所有 set 方法(TurnOn/TurnOff/SetMode/SetTargetTemp/SetFanSpeed/
  SetSwing/SetTargetHumidity/SetAuxHeat/SetEco)成功后记录目标值 rbTarget
- 同属性目标覆盖、不同属性追加,连续 set 自动聚合

**防抖 + 重试验证(readbackLoop)**
- 静默窗口 readbackIdle=3s:等待最后一次 set 无新操作才开始验证
- 未命中目标时按 readbackRetryInterval=3s 重试,上限 readbackMaxRetries=5
- 整体超时 readbackTimeout=30s,防云端卡住导致回读永久挂起
- 命中后逐字段回调 OnPropsChanged 注册的 handler(与 MQTT 推送同语义)
- rbValueEqual 数值比较用 epsilon=0.01,bool/string 直接相等

**白名单配置透传**
- xiaomi.Config 新增 ReadbackModels 白名单(留空 = 全部不回读,默认)
- Client 新增 NeedReadback(model) 判定 + Logger() 访问器
- bridge.Pool.NewPool(readbackModels...) 透传;newBridge 传入 Client 配置

**其他**
- bridge 设备创建失败(Warn)与不支持跳过(Debug)分别记录日志
This commit is contained in:
2026-08-23 23:33:32 +08:00
parent 86c21ce91c
commit a0ef854726
4 changed files with 352 additions and 21 deletions
+16 -5
View File
@@ -47,11 +47,16 @@ type OccupancyProps struct {
type Pool struct {
mu sync.RWMutex
items map[int64]*Bridge
readbackModels []string // 需要 set 后回读的型号列表(透传到 xiaomi.Config)
}
// NewPool 创建 Pool。
func NewPool() *Pool {
return &Pool{items: make(map[int64]*Bridge)}
// NewPool 创建 Pool。readbackModels 为需要"set 后回读"的设备型号列表(可选,默认空)。
func NewPool(readbackModels ...[]string) *Pool {
var models []string
if len(readbackModels) > 0 {
models = readbackModels[0]
}
return &Pool{items: make(map[int64]*Bridge), readbackModels: models}
}
// Get 按 id 查找。
@@ -69,7 +74,7 @@ func (p *Pool) Add(ctx context.Context, id int64, uuid string, auth xiaomi.AuthI
if _, exists := p.items[id]; exists {
return nil, fmt.Errorf("bridge: %d already exists", id)
}
b, err := newBridge(ctx, id, uuid, auth, dataDir, homeIDs)
b, err := newBridge(ctx, id, uuid, auth, dataDir, homeIDs, p.readbackModels)
if err != nil {
return nil, err
}
@@ -115,7 +120,7 @@ type Bridge struct {
onOnlineChanged func(id int64, did string, online bool)
}
func newBridge(ctx context.Context, id int64, uuid string, auth xiaomi.AuthInfo, dataDir string, homeIDs []int64) (*Bridge, error) {
func newBridge(ctx context.Context, id int64, uuid string, auth xiaomi.AuthInfo, dataDir string, homeIDs []int64, readbackModels []string) (*Bridge, error) {
idStr := fmt.Sprintf("%d", id)
client, err := xiaomi.NewClient(ctx, xiaomi.Config{
AuthInfo: auth,
@@ -123,6 +128,7 @@ func newBridge(ctx context.Context, id int64, uuid string, auth xiaomi.AuthInfo,
UUID: uuid,
UID: idStr,
HomeIDs: homeIDs,
ReadbackModels: readbackModels,
})
if err != nil {
return nil, fmt.Errorf("bridge: %w", err)
@@ -144,6 +150,11 @@ func newBridge(ctx context.Context, id int64, uuid string, auth xiaomi.AuthInfo,
for _, info := range devList {
d, createErr := devices.Create(client, info)
if createErr != nil || d == nil {
if createErr != nil {
b.lgr.Warnf("[bridge] %d: create device %s (%s) failed: %v", id, info.DID, info.Model, createErr)
} else {
b.lgr.Debugf("[bridge] %d: device %s (%s) unsupported, skipped", id, info.DID, info.Model)
}
continue
}
b.dv[info.DID] = d
+20
View File
@@ -25,6 +25,8 @@ type Client struct {
propSubs map[string]string
eventSubs map[string]string
mu sync.RWMutex
readbackModels map[string]struct{} // 需要 set 后回读的型号集合
}
// NewClient 创建并初始化客户端。
@@ -103,13 +105,31 @@ func NewClient(ctx context.Context, cfg Config) (*Client, error) {
devices: make(map[string]*miot.MIoTDevice),
propSubs: make(map[string]string),
eventSubs: make(map[string]string),
readbackModels: make(map[string]struct{}, len(cfg.ReadbackModels)),
}
for _, m := range cfg.ReadbackModels {
if m != "" {
c.readbackModels[m] = struct{}{}
}
}
return c, nil
}
// NeedReadback 判断指定型号是否需要"set 成功后回读状态"。
func (c *Client) NeedReadback(model string) bool {
if model == "" || len(c.readbackModels) == 0 {
return false
}
_, ok := c.readbackModels[model]
return ok
}
// SetLogger 设置日志记录器。
func (c *Client) SetLogger(l logger.Logger) { c.lgr = l }
// Logger 返回日志记录器。
func (c *Client) Logger() logger.Logger { return c.lgr }
// SetHomeIDs 动态设置 miot 层的家庭过滤,供后台定时 RefreshDevices 使用。
func (c *Client) SetHomeIDs(homeIDs []string) { c.inner.SetHomeIDs(homeIDs) }
+6
View File
@@ -18,4 +18,10 @@ type Config struct {
RedirectURL string // OAuth 授权回调 URL
HomeIDs []int64 // 需要管理的家庭 ID 列表,nil/空=全部
// ReadbackModels 需要"set 成功后回读状态"的设备型号列表。
// 用于设备固件不广播云端指令的场景(如部分 BLE/Mesh 设备):
// 控制成功后 SDK 防抖延迟回读真实状态,并通过 OnPropsChanged 送达业务层。
// 留空 = 所有设备都不做回读(默认)。
ReadbackModels []string
}
+303 -9
View File
@@ -3,6 +3,8 @@ package devices
import (
"context"
"strings"
"sync"
"time"
"xiaomihome/miot"
"xiaomihome/xiaomi"
@@ -66,6 +68,19 @@ const (
ACFanHigh = "high"
ACFanAuto = "auto"
// readbackIdle set 后回读的防抖窗口:最后一次 set 之后等待该时长无新操作才开始验证。
// 用于固件不广播云端指令的设备(BLE/Mesh),防抖合并连续操作。
readbackIdle = 3 * time.Second
// readbackRetryInterval 回读未命中目标时的重试间隔。
readbackRetryInterval = 3 * time.Second
// readbackMaxRetries 回读最大重试次数(含首次尝试)。
readbackMaxRetries = 5
// readbackTimeout 回读整体超时(含防抖等待 + 重试),防止云端卡住导致回读永久失效。
readbackTimeout = 30 * time.Second
ACSwingOff = "off"
ACSwingVertical = "vertical"
ACSwingHorizontal = "horizontal"
@@ -165,6 +180,19 @@ type acDevice struct {
swingMap *valMapper // SPEC value→name for swing
propsCache ACProps // MQTT 推送维护的本地缓存
// set 后回读(验证设置是否成功):记录目标值,防抖 + 重试比较
rbMu sync.Mutex
rbPending bool // 是否有回读调度在跑
rbLastSet time.Time // 最近一次 set 时间
rbTargets []rbTarget // 本次待验证的目标(聚合连续 set)
rbHandler func(props ACProps, changed string) // OnPropsChanged 注册的 handler
}
// rbTarget 回读验证目标:某属性本次 set 的目标值。
type rbTarget struct {
name string // "on"/"mode"/"target-temperature"/"fan-level"/"swing"/...
want interface{} // 目标值:bool / float64 / string(与 ACProps 字段类型一致)
}
// NewAirConditioner creates an AirConditioner control for the given device.
@@ -302,14 +330,22 @@ func (d *acDevice) TurnOn(ctx context.Context) error {
if !d.propOn.valid {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propOn.siid, d.propOn.piid, true)
if err := d.SetProp(ctx, d.propOn.siid, d.propOn.piid, true); err != nil {
return err
}
d.recordReadback("on", true)
return nil
}
func (d *acDevice) TurnOff(ctx context.Context) error {
if !d.propOn.valid {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propOn.siid, d.propOn.piid, false)
if err := d.SetProp(ctx, d.propOn.siid, d.propOn.piid, false); err != nil {
return err
}
d.recordReadback("on", false)
return nil
}
func (d *acDevice) IsOn(ctx context.Context) (bool, error) {
@@ -341,7 +377,11 @@ func (d *acDevice) SetMode(ctx context.Context, mode string) error {
if !ok {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propMode.siid, d.propMode.piid, v)
if err := d.SetProp(ctx, d.propMode.siid, d.propMode.piid, v); err != nil {
return err
}
d.recordReadback("mode", d.modeMap.name(v))
return nil
}
func (d *acDevice) GetMode(ctx context.Context) (string, error) {
@@ -359,7 +399,11 @@ func (d *acDevice) SetTargetTemp(ctx context.Context, celsius float64) error {
if !d.propTargetTemp.valid {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propTargetTemp.siid, d.propTargetTemp.piid, celsius)
if err := d.SetProp(ctx, d.propTargetTemp.siid, d.propTargetTemp.piid, celsius); err != nil {
return err
}
d.recordReadback("target-temperature", celsius)
return nil
}
func (d *acDevice) GetTargetTemp(ctx context.Context) (float64, error) {
@@ -392,7 +436,11 @@ func (d *acDevice) SetFanSpeed(ctx context.Context, speed string) error {
if !ok {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propFanSpeed.siid, d.propFanSpeed.piid, v)
if err := d.SetProp(ctx, d.propFanSpeed.siid, d.propFanSpeed.piid, v); err != nil {
return err
}
d.recordReadback("fan-level", d.fanMap.name(v))
return nil
}
func (d *acDevice) GetFanSpeed(ctx context.Context) (string, error) {
@@ -414,7 +462,11 @@ func (d *acDevice) SetSwing(ctx context.Context, mode string) error {
if !ok {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propSwing.siid, d.propSwing.piid, v)
if err := d.SetProp(ctx, d.propSwing.siid, d.propSwing.piid, v); err != nil {
return err
}
d.recordReadback("swing", d.swingMap.name(v))
return nil
}
func (d *acDevice) GetSwing(ctx context.Context) (string, error) {
@@ -432,7 +484,11 @@ func (d *acDevice) SetTargetHumidity(ctx context.Context, pct int) error {
if !d.propTargetHumidity.valid {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propTargetHumidity.siid, d.propTargetHumidity.piid, pct)
if err := d.SetProp(ctx, d.propTargetHumidity.siid, d.propTargetHumidity.piid, pct); err != nil {
return err
}
d.recordReadback("target-humidity", float64(pct))
return nil
}
func (d *acDevice) GetTargetHumidity(ctx context.Context) (int, error) {
@@ -461,7 +517,11 @@ func (d *acDevice) SetAuxHeat(ctx context.Context, on bool) error {
if !d.propAuxHeat.valid {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propAuxHeat.siid, d.propAuxHeat.piid, on)
if err := d.SetProp(ctx, d.propAuxHeat.siid, d.propAuxHeat.piid, on); err != nil {
return err
}
d.recordReadback("heater", on)
return nil
}
func (d *acDevice) GetAuxHeat(ctx context.Context) (bool, error) {
@@ -482,7 +542,11 @@ func (d *acDevice) SetEco(ctx context.Context, on bool) error {
if !d.propEco.valid {
return xiaomi.ErrNotSupported
}
return d.SetProp(ctx, d.propEco.siid, d.propEco.piid, on)
if err := d.SetProp(ctx, d.propEco.siid, d.propEco.piid, on); err != nil {
return err
}
d.recordReadback("eco", on)
return nil
}
func (d *acDevice) GetEco(ctx context.Context) (bool, error) {
@@ -534,6 +598,10 @@ func (d *acDevice) OnPropsChanged(handler func(props ACProps, changed string)) (
if len(fields) == 0 {
return "", xiaomi.ErrNotSupported
}
// 保存 handler 供 set 后回读使用
d.rbMu.Lock()
d.rbHandler = handler
d.rbMu.Unlock()
// 每个属性精确订阅 siid/piid,避免 lumi.acpartner.mcn02 等设备
// siid/piid 冲突时把风速 int 路由到 on 槽位
var firstID string
@@ -564,6 +632,232 @@ func (d *acDevice) OnPropsChanged(handler func(props ACProps, changed string)) (
return firstID, nil
}
// recordReadback 记录本次 set 的目标值并调度回读验证(白名单型号)。
// 连续 set 聚合:同属性覆盖目标值,不同属性追加。
func (d *acDevice) recordReadback(name string, want interface{}) {
if d.client == nil || !d.client.NeedReadback(d.info.Model) {
return
}
d.rbMu.Lock()
replaced := false
for i := range d.rbTargets {
if d.rbTargets[i].name == name {
d.rbTargets[i].want = want
replaced = true
break
}
}
if !replaced {
d.rbTargets = append(d.rbTargets, rbTarget{name: name, want: want})
}
d.rbLastSet = time.Now()
if d.rbPending {
d.rbMu.Unlock()
return // 已有回读调度在跑(goroutine 会检查 lastSetAt 继续等待)
}
d.rbPending = true
d.rbMu.Unlock()
if d.client.Logger() != nil {
d.client.Logger().Debugf("[acDevice] readback scheduled: %s (%s) targets=%v\n", d.info.DID, d.info.Model, d.rbTargetNames())
}
go d.readbackLoop()
}
// readbackLoop 防抖 + 重试回读,验证 set 目标值是否被设备执行。
// 静默 readbackIdle 后开始回读;未命中目标则间隔重试,直到命中或超时。
func (d *acDevice) readbackLoop() {
ctx, cancel := context.WithTimeout(context.Background(), readbackTimeout)
defer cancel()
for {
// 防抖:等待静默窗口
select {
case <-time.After(readbackIdle):
case <-ctx.Done():
d.finishReadback("timeout waiting idle")
return
}
d.rbMu.Lock()
idle := time.Since(d.rbLastSet) >= readbackIdle
targets := make([]rbTarget, len(d.rbTargets))
copy(targets, d.rbTargets)
d.rbMu.Unlock()
if !idle {
continue // 期间又有新的 set,继续等待静默
}
// 回读验证(带重试)
for attempt := 1; ; attempt++ {
if attempt > 1 {
select {
case <-time.After(readbackRetryInterval):
case <-ctx.Done():
d.finishReadback("timeout")
return
}
}
props, err := d.FetchProps(ctx)
if err != nil {
if d.client.Logger() != nil {
d.client.Logger().Warnf("[acDevice] readback %s attempt %d failed: %v\n", d.info.DID, attempt, err)
}
if ctx.Err() != nil {
d.finishReadback("timeout")
return
}
continue // 拉取失败重试
}
missed := d.checkReadbackTargets(props, targets)
if len(missed) == 0 {
d.finishReadbackOK(props, targets)
return
}
if d.client.Logger() != nil {
d.client.Logger().Debugf("[acDevice] readback %s attempt %d not applied yet: missed=%v\n", d.info.DID, attempt, missed)
}
if attempt >= readbackMaxRetries {
d.finishReadback("not applied after retries")
return
}
}
}
}
// checkReadbackTargets 检查回读状态是否达到所有 set 目标值,返回未命中的属性名。
func (d *acDevice) checkReadbackTargets(props ACProps, targets []rbTarget) []string {
var missed []string
for _, t := range targets {
if !rbValueEqual(readbackFieldValue(props, t.name), t.want) {
missed = append(missed, t.name)
}
}
return missed
}
// finishReadbackOK 全部目标命中:逐字段上报业务层(与 MQTT 推送同语义)。
func (d *acDevice) finishReadbackOK(props ACProps, targets []rbTarget) {
d.rbMu.Lock()
handler := d.rbHandler
d.rbMu.Unlock()
if handler == nil {
if d.client.Logger() != nil {
d.client.Logger().Warnf("[acDevice] readback %s done but handler not registered (OnPropsChanged never called)\n", d.info.DID)
}
d.rbMu.Lock()
d.rbPending = false
d.rbMu.Unlock()
return
}
for _, t := range targets {
handler(props, t.name)
}
if d.client.Logger() != nil {
d.client.Logger().Debugf("[acDevice] readback %s success: %v on=%v mode=%s target=%.0f temp=%.0f fan=%s hum=%d\n",
d.info.DID, d.rbTargetNamesLocked(), props.On, props.Mode, props.TargetTemperature,
props.Temperature, props.FanLevel, props.CurrentHumidity)
}
d.rbMu.Lock()
d.rbPending = false
d.rbMu.Unlock()
}
// finishReadback 回读终止(超时/未生效):清理调度并记录日志。
func (d *acDevice) finishReadback(reason string) {
d.rbMu.Lock()
names := d.rbTargetNamesLocked()
d.rbTargets = nil
d.rbPending = false
d.rbMu.Unlock()
if d.client.Logger() != nil {
d.client.Logger().Warnf("[acDevice] readback %s finished: %s targets=%v\n", d.info.DID, reason, names)
}
}
func (d *acDevice) rbTargetNames() []string {
d.rbMu.Lock()
defer d.rbMu.Unlock()
return d.rbTargetNamesLocked()
}
func (d *acDevice) rbTargetNamesLocked() []string {
names := make([]string, 0, len(d.rbTargets))
for _, t := range d.rbTargets {
names = append(names, t.name)
}
return names
}
// readbackFieldValue 从 ACProps 取指定属性的值(与 rbTarget.want 类型一致)。
func readbackFieldValue(props ACProps, name string) interface{} {
switch name {
case "on":
return props.On
case "mode":
return props.Mode
case "target-temperature":
return props.TargetTemperature
case "temperature":
return props.Temperature
case "target-humidity":
return float64(props.TargetHumidity)
case "relative-humidity":
return float64(props.CurrentHumidity)
case "fan-level":
return props.FanLevel
case "swing":
return props.Swing
case "heater":
return props.AuxHeat
case "eco":
return props.Eco
case "electric-power":
return props.ElectricPower
case "power-consumption":
return props.PowerConsumption
}
return nil
}
// rbValueEqual 比较回读值与目标值:数值用 epsilon,bool/string 直接相等。
func rbValueEqual(got, want interface{}) bool {
if got == nil || want == nil {
return got == want
}
if gf, ok := got.(float64); ok {
wf, ok2 := toFloatOk(want)
if !ok2 {
return false
}
diff := gf - wf
if diff < 0 {
diff = -diff
}
return diff < 0.01
}
if gb, ok := got.(bool); ok {
wb, ok2 := want.(bool)
return ok2 && gb == wb
}
if gs, ok := got.(string); ok {
ws, ok2 := want.(string)
return ok2 && gs == ws
}
return got == want
}
// toFloatOk 将 interface{} 转为 float64(bool/string 除外),失败返回 false。
func toFloatOk(v interface{}) (float64, bool) {
switch n := v.(type) {
case float64:
return n, true
case float32:
return float64(n), true
case int:
return float64(n), true
case int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64:
return toFloat(v), true
}
return 0, false
}
func (d *acDevice) propFields() []propField {
r := d.getSpecResolver()
var fields []propField