diff --git a/bridge/bridge.go b/bridge/bridge.go index 78da418..3fc2b90 100644 --- a/bridge/bridge.go +++ b/bridge/bridge.go @@ -45,13 +45,18 @@ type OccupancyProps struct { // Pool 全局路由:id → Bridge。 type Pool struct { - mu sync.RWMutex - items map[int64]*Bridge + 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,14 +120,15 @@ 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, - DataDir: dataDir, - UUID: uuid, - UID: idStr, - HomeIDs: homeIDs, + AuthInfo: auth, + DataDir: dataDir, + 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 diff --git a/xiaomi/client.go b/xiaomi/client.go index 2a5f305..cbb83f2 100644 --- a/xiaomi/client.go +++ b/xiaomi/client.go @@ -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) } diff --git a/xiaomi/config.go b/xiaomi/config.go index 2844a8b..a5dbc95 100644 --- a/xiaomi/config.go +++ b/xiaomi/config.go @@ -18,4 +18,10 @@ type Config struct { RedirectURL string // OAuth 授权回调 URL HomeIDs []int64 // 需要管理的家庭 ID 列表,nil/空=全部 + + // ReadbackModels 需要"set 成功后回读状态"的设备型号列表。 + // 用于设备固件不广播云端指令的场景(如部分 BLE/Mesh 设备): + // 控制成功后 SDK 防抖延迟回读真实状态,并通过 OnPropsChanged 送达业务层。 + // 留空 = 所有设备都不做回读(默认)。 + ReadbackModels []string } diff --git a/xiaomi/devices/air_conditioner.go b/xiaomi/devices/air_conditioner.go index 97f32b8..b4b1cb7 100644 --- a/xiaomi/devices/air_conditioner.go +++ b/xiaomi/devices/air_conditioner.go @@ -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