diff --git a/charger/openwb.go b/charger/openwb.go index 9685fb99c..406898354 100644 --- a/charger/openwb.go +++ b/charger/openwb.go @@ -60,39 +60,18 @@ func NewOpenWB(log *util.Logger, mqttconf mqtt.Config, id int, topic string, p1p } // timeout handler - timer := provider.NewMqtt(log, client, + to := provider.NewTimeoutHandler(provider.NewMqtt(log, client, fmt.Sprintf("%s/system/%s", topic, openwb.TimestampTopic), 1, timeout, - ).IntGetter() + ).StringGetter()) - // getters boolG := func(topic string) func() (bool, error) { g := provider.NewMqtt(log, client, topic, 1, 0).BoolGetter() - return func() (val bool, err error) { - if val, err = g(); err == nil { - _, err = timer() - } - return val, err - } + return to.BoolGetter(g) } - // intG := func(topic string) func() (int64, error) { - // g := provider.NewMqtt(log, client, topic, 1, 0).IntGetter() - // return func() (val int64, err error) { - // if val, err = g(); err == nil { - // _, err = timer() - // } - // return val, err - // } - // } - floatG := func(topic string) func() (float64, error) { g := provider.NewMqtt(log, client, topic, 1, 0).FloatGetter() - return func() (val float64, err error) { - if val, err = g(); err == nil { - _, err = timer() - } - return val, err - } + return to.FloatGetter(g) } // check if loadpoint configured diff --git a/charger/warp.go b/charger/warp.go index 131a9e3ed..8720bd3c4 100644 --- a/charger/warp.go +++ b/charger/warp.go @@ -92,18 +92,13 @@ func NewWarp(mqttconf mqtt.Config, topic string, timeout time.Duration) (*Warp, } // timeout handler - timer := provider.NewMqtt(log, client, + to := provider.NewTimeoutHandler(provider.NewMqtt(log, client, fmt.Sprintf("%s/evse/state", topic), 1, timeout, - ).StringGetter() + ).StringGetter()) stringG := func(topic string) func() (string, error) { g := provider.NewMqtt(log, client, topic, 1, 0).StringGetter() - return func() (val string, err error) { - if val, err = g(); err == nil { - _, err = timer() - } - return val, err - } + return to.StringGetter(g) } wb.enabledG = stringG(fmt.Sprintf("%s/evse/auto_start_charging", topic)) diff --git a/meter/openwb.go b/meter/openwb.go index cfce77501..b13f6c9c3 100644 --- a/meter/openwb.go +++ b/meter/openwb.go @@ -41,33 +41,18 @@ func NewOpenWBFromConfig(other map[string]interface{}) (api.Meter, error) { } // timeout handler - timer := provider.NewMqtt(log, client, + to := provider.NewTimeoutHandler(provider.NewMqtt(log, client, fmt.Sprintf("%s/system/%s", cc.Topic, openwb.TimestampTopic), 1, cc.Timeout, - ).IntGetter() + ).StringGetter()) - // getters boolG := func(topic string) func() (bool, error) { g := provider.NewMqtt(log, client, topic, 1, 0).BoolGetter() - return func() (val bool, err error) { - if val, err = g(); err == nil { - _, err = timer() - } - return val, err - } + return to.BoolGetter(g) } - floatG := func(topic string, scaler ...float64) func() (float64, error) { - scale := 1.0 - if len(scaler) == 1 { - scale = scaler[0] - } + floatG := func(topic string) func() (float64, error) { g := provider.NewMqtt(log, client, topic, 1, 0).FloatGetter() - return func() (val float64, err error) { - if val, err = g(); err == nil { - _, err = timer() - } - return scale * val, err - } + return to.FloatGetter(g) } var power func() (float64, error) @@ -87,7 +72,7 @@ func NewOpenWBFromConfig(other map[string]interface{}) (api.Meter, error) { currents = collectCurrentProviders(curr) case "pv": - configuredG := boolG(fmt.Sprintf("%s/pv/%s", cc.Topic, openwb.PvConfigured)) + configuredG := boolG(fmt.Sprintf("%s/pv/1/%s", cc.Topic, openwb.PvConfigured)) // first pv configured, err := configuredG() if err != nil { return nil, err @@ -110,7 +95,7 @@ func NewOpenWBFromConfig(other map[string]interface{}) (api.Meter, error) { return nil, errors.New("battery not available") } - power = floatG(fmt.Sprintf("%s/housebattery/%s", cc.Topic, openwb.PowerTopic), -1) + power = floatG(fmt.Sprintf("%s/housebattery/%s", cc.Topic, openwb.PowerTopic)) soc = floatG(fmt.Sprintf("%s/housebattery/%s", cc.Topic, openwb.SoCTopic)) default: diff --git a/provider/mqtt_timeout.go b/provider/mqtt_timeout.go new file mode 100644 index 000000000..1a3d71036 --- /dev/null +++ b/provider/mqtt_timeout.go @@ -0,0 +1,37 @@ +package provider + +// TimeoutHandler is a wrapper for a Getter that times out after a given duration +type TimeoutHandler struct { + ticker func() (string, error) +} + +func NewTimeoutHandler(ticker func() (string, error)) *TimeoutHandler { + return &TimeoutHandler{ticker} +} + +func (h *TimeoutHandler) BoolGetter(g func() (bool, error)) func() (bool, error) { + return func() (val bool, err error) { + if val, err = g(); err == nil { + _, err = h.ticker() + } + return val, err + } +} + +func (h *TimeoutHandler) FloatGetter(g func() (float64, error)) func() (float64, error) { + return func() (val float64, err error) { + if val, err = g(); err == nil { + _, err = h.ticker() + } + return val, err + } +} + +func (h *TimeoutHandler) StringGetter(g func() (string, error)) func() (string, error) { + return func() (val string, err error) { + if val, err = g(); err == nil { + _, err = h.ticker() + } + return val, err + } +}