diff --git a/charger/openwb.go b/charger/openwb.go index 62d1db99c..baed6950d 100644 --- a/charger/openwb.go +++ b/charger/openwb.go @@ -74,7 +74,7 @@ func NewOpenWB(log *util.Logger, mqttconf mqtt.Config, id int, topic string, p1p } // check if loadpoint configured - configured, err := to.BoolGetter(mq(openwb.ConfiguredTopic)) + configured, err := provider.NewMqtt(log, client, fmt.Sprintf("%s/lp/%d/%s", topic, id, openwb.ConfiguredTopic), timeout).BoolGetter() if err != nil { return nil, err } diff --git a/charger/warp2.go b/charger/warp2.go index cd405eb09..dd15be577 100644 --- a/charger/warp2.go +++ b/charger/warp2.go @@ -60,19 +60,19 @@ func NewWarp2FromConfig(other map[string]interface{}) (api.Charger, error) { } var currentPower, totalEnergy func() (float64, error) - if wb.hasFeature(cc.Topic, warp.FeatureMeter) { + if wb.hasFeature(cc.Topic, warp.FeatureMeter, cc.Timeout) { currentPower = wb.currentPower totalEnergy = wb.totalEnergy } var currents, voltages func() (float64, float64, float64, error) - if wb.hasFeature(cc.Topic, warp.FeatureMeterPhases) { + if wb.hasFeature(cc.Topic, warp.FeatureMeterPhases, cc.Timeout) { currents = wb.currents voltages = wb.voltages } var identity func() (string, error) - if wb.hasFeature(cc.Topic, warp.FeatureNfc) { + if wb.hasFeature(cc.Topic, warp.FeatureNfc, cc.Timeout) { identity = wb.identify } @@ -160,14 +160,14 @@ func NewWarp2(mqttconf mqtt.Config, topic, emTopic string, timeout time.Duration return wb, nil } -func (wb *Warp2) hasFeature(root, feature string) bool { +func (wb *Warp2) hasFeature(root, feature string, timeout time.Duration) bool { if wb.features != nil { return slices.Contains(wb.features, feature) } topic := fmt.Sprintf("%s/info/features", root) - if dataG, err := provider.NewMqtt(wb.log, wb.client, topic, 0).StringGetter(); err == nil { + if dataG, err := provider.NewMqtt(wb.log, wb.client, topic, timeout).StringGetter(); err == nil { if data, err := dataG(); err == nil { if err := json.Unmarshal([]byte(data), &wb.features); err == nil { return slices.Contains(wb.features, feature) diff --git a/meter/openwb.go b/meter/openwb.go index 3f6234777..11242cb2f 100644 --- a/meter/openwb.go +++ b/meter/openwb.go @@ -26,8 +26,8 @@ func NewOpenWBFromConfig(other map[string]interface{}) (api.Meter, error) { Usage string capacity `mapstructure:",squash"` }{ - Topic: "openWB", - Timeout: 15 * time.Second, + Topic: openwb.RootTopic, + Timeout: openwb.Timeout, } if err := util.DecodeOther(other, &cc); err != nil { @@ -76,7 +76,8 @@ func NewOpenWBFromConfig(other map[string]interface{}) (api.Meter, error) { currents = collectPhaseProviders(curr) case "pv": - configuredG, err := to.BoolGetter(mq("%s/pv/1/%s", cc.Topic, openwb.PvConfigured)) // first pv + // first pv + configuredG, err := provider.NewMqtt(log, client, fmt.Sprintf("%s/pv/1/%s", cc.Topic, openwb.PvConfigured), cc.Timeout).BoolGetter() if err != nil { return nil, err } @@ -99,7 +100,7 @@ func NewOpenWBFromConfig(other map[string]interface{}) (api.Meter, error) { } case "battery": - configuredG, err := to.BoolGetter(mq("%s/housebattery/%s", cc.Topic, openwb.BatteryConfigured)) + configuredG, err := provider.NewMqtt(log, client, fmt.Sprintf("%s/housebattery/%s", cc.Topic, openwb.BatteryConfigured), cc.Timeout).BoolGetter() if err != nil { return nil, err } diff --git a/provider/mqtt/client.go b/provider/mqtt/client.go index 2aeb7568e..e59324643 100644 --- a/provider/mqtt/client.go +++ b/provider/mqtt/client.go @@ -141,8 +141,6 @@ func (m *Client) Listen(topic string, callback func(string)) error { case <-time.After(request.Timeout): return fmt.Errorf("subscribe: %s: %w", topic, api.ErrTimeout) case <-token.Done(): - // TODO depends on https://github.com/eclipse/paho.mqtt.golang/issues/663 - time.Sleep(100 * time.Millisecond) return nil } } diff --git a/util/monitor.go b/util/monitor.go index 5be4ac26b..98693a7cb 100644 --- a/util/monitor.go +++ b/util/monitor.go @@ -64,9 +64,32 @@ func (m *Monitor[T]) GetFunc(get func(T)) error { default: return api.ErrOutdated } - } else if time.Since(m.updated) > m.timeout { + } + + if time.Since(m.updated) > m.timeout { + err := api.ErrOutdated + + // wait once on very first call + if m.updated.IsZero() { + m.mu.RUnlock() + + // mark as waited once + m.mu.Lock() + m.updated.Add(time.Nanosecond) + m.mu.Unlock() + + select { + case <-m.done: + // got value and updated timestamp + err = nil + case <-time.After(m.timeout): + } + + m.mu.RLock() + } + get(m.val) - return api.ErrOutdated + return err } get(m.val)