diff --git a/xiaomi/devices/fan.go b/xiaomi/devices/fan.go index 78c178a..c1a5281 100644 --- a/xiaomi/devices/fan.go +++ b/xiaomi/devices/fan.go @@ -2,6 +2,10 @@ package devices import ( "context" + "sync" + "time" + + "github.com/sirupsen/logrus" "git.zeroonesoft.cn/golib/xiaomihome/miot" "git.zeroonesoft.cn/golib/xiaomihome/xiaomi" @@ -59,6 +63,14 @@ type fanDevice struct { fanLevelValues []int // 可用风速档位值列表 propsCache FanProps + + // set 后回读(验证设置是否成功):记录目标值,防抖 + 重试比较。 + // 用于固件不广播云端指令的 BLE/Mesh 设备(如 ecosnu.airfresh.eksn1)。 + rbMu sync.Mutex + rbPending bool // 是否有回读调度在跑 + rbLastSet time.Time // 最近一次 set 时间 + rbTargets []rbTarget // 本次待验证的目标(聚合连续 set) + rbHandler func(props FanProps, changed string) // OnPropsChanged 注册的 handler } // NewFan creates a Fan control for the given device. @@ -119,14 +131,22 @@ func (d *fanDevice) 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 *fanDevice) 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 *fanDevice) IsOn(ctx context.Context) (bool, error) { @@ -151,7 +171,11 @@ func (d *fanDevice) SetFanLevel(ctx context.Context, level int) error { if !d.propFanLevel.valid { return xiaomi.ErrNotSupported } - return d.SetProp(ctx, d.propFanLevel.siid, d.propFanLevel.piid, level) + if err := d.SetProp(ctx, d.propFanLevel.siid, d.propFanLevel.piid, level); err != nil { + return err + } + d.recordReadback("fan-level", level) + return nil } func (d *fanDevice) GetFanLevel(ctx context.Context) (int, error) { @@ -173,7 +197,11 @@ func (d *fanDevice) SetMode(ctx context.Context, mode int) error { if !d.propMode.valid { return xiaomi.ErrNotSupported } - return d.SetProp(ctx, d.propMode.siid, d.propMode.piid, mode) + if err := d.SetProp(ctx, d.propMode.siid, d.propMode.piid, mode); err != nil { + return err + } + d.recordReadback("mode", mode) + return nil } func (d *fanDevice) GetMode(ctx context.Context) (int, error) { @@ -191,7 +219,11 @@ func (d *fanDevice) SetOscillation(ctx context.Context, on bool) error { if !d.propOscillation.valid { return xiaomi.ErrNotSupported } - return d.SetProp(ctx, d.propOscillation.siid, d.propOscillation.piid, on) + if err := d.SetProp(ctx, d.propOscillation.siid, d.propOscillation.piid, on); err != nil { + return err + } + d.recordReadback("oscillation", on) + return nil } func (d *fanDevice) IsOscillating(ctx context.Context) (bool, error) { @@ -216,7 +248,11 @@ func (d *fanDevice) SetNaturalWind(ctx context.Context, on bool) error { if on { v = 1 } - return d.SetProp(ctx, d.propNaturalWind.siid, d.propNaturalWind.piid, v) + if err := d.SetProp(ctx, d.propNaturalWind.siid, d.propNaturalWind.piid, v); err != nil { + return err + } + d.recordReadback("natural-wind", on) + return nil } func (d *fanDevice) IsNaturalWind(ctx context.Context) (bool, error) { @@ -261,6 +297,10 @@ func (d *fanDevice) OnPropsChanged(handler func(props FanProps, changed string)) if len(fields) == 0 { return "", xiaomi.ErrNotSupported } + // 保存 handler 供 set 后回读使用 + d.rbMu.Lock() + d.rbHandler = handler + d.rbMu.Unlock() return d.SubProp(0, 0, func(did string, prop *xiaomi.PropertyValue) { for _, f := range fields { if f.siid == prop.SIID && f.piid == prop.PIID { @@ -314,6 +354,163 @@ func (d *fanDevice) FetchProps(ctx context.Context) (FanProps, error) { return result, nil } +// recordReadback 记录本次 set 的目标值并调度回读验证(白名单型号)。 +// 连续 set 聚合:同属性覆盖目标值,不同属性追加。与 acDevice 同机制。 +func (d *fanDevice) 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() + logrus.Debugf("[fanDevice] readback scheduled: %s (%s) targets=%v\n", d.info.DID, d.info.Model, d.rbTargetNames()) + go d.readbackLoop() +} + +// readbackLoop 防抖 + 重试回读,验证 set 目标值是否被设备执行。 +// 静默 readbackIdle 后开始回读;未命中目标则间隔重试,直到命中或超时。 +func (d *fanDevice) 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 { + logrus.Warnf("[fanDevice] 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 + } + logrus.Debugf("[fanDevice] 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 *fanDevice) checkReadbackTargets(props FanProps, targets []rbTarget) []string { + var missed []string + for _, t := range targets { + if !rbValueEqual(fanReadbackFieldValue(props, t.name), t.want) { + missed = append(missed, t.name) + } + } + return missed +} + +// finishReadbackOK 全部目标命中:逐字段上报业务层(与 MQTT 推送同语义)。 +func (d *fanDevice) finishReadbackOK(props FanProps, targets []rbTarget) { + d.rbMu.Lock() + handler := d.rbHandler + d.rbMu.Unlock() + if handler == nil { + logrus.Warnf("[fanDevice] 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) + } + logrus.Debugf("[fanDevice] readback %s success: on=%v fanLevel=%d mode=%d swing=%v naturalWind=%v\n", + d.info.DID, props.On, props.FanLevel, props.Mode, props.Oscillation, props.NaturalWind) + d.rbMu.Lock() + d.rbPending = false + d.rbMu.Unlock() +} + +// finishReadback 回读终止(超时/未生效):清理调度并记录日志。 +func (d *fanDevice) finishReadback(reason string) { + d.rbMu.Lock() + names := d.rbTargetNamesLocked() + d.rbTargets = nil + d.rbPending = false + d.rbMu.Unlock() + logrus.Warnf("[fanDevice] readback %s finished: %s targets=%v\n", d.info.DID, reason, names) +} + +func (d *fanDevice) rbTargetNames() []string { + d.rbMu.Lock() + defer d.rbMu.Unlock() + return d.rbTargetNamesLocked() +} + +func (d *fanDevice) rbTargetNamesLocked() []string { + names := make([]string, 0, len(d.rbTargets)) + for _, t := range d.rbTargets { + names = append(names, t.name) + } + return names +} + +// fanReadbackFieldValue 从 FanProps 取指定属性的值(与 rbTarget.want 类型一致)。 +func fanReadbackFieldValue(props FanProps, name string) interface{} { + switch name { + case "on": + return props.On + case "fan-level": + return props.FanLevel + case "mode": + return props.Mode + case "oscillation": + return props.Oscillation + case "natural-wind": + return props.NaturalWind + } + return nil +} + // rangeValues 从 SPEC value-range 生成离散值列表,如 [1,8,1] → [1,2,3,4,5,6,7,8] func rangeValues(r *miot.MIoTSpecValueRange) []int { if r == nil {