From d21d15e9d73180da8ac6e0b61fcb90b5745d3cf3 Mon Sep 17 00:00:00 2001 From: andig Date: Tue, 9 Mar 2021 14:58:43 +0100 Subject: [PATCH] Fix mqtt logging and return values (#739) --- api/error.go | 10 ++++++++++ provider/mqtt.go | 21 ++++++++++++--------- provider/mqtt/client.go | 13 ++++++++----- provider/mqtt/registry.go | 6 +++--- 4 files changed, 33 insertions(+), 17 deletions(-) diff --git a/api/error.go b/api/error.go index 17180a634..39a3bf7e5 100644 --- a/api/error.go +++ b/api/error.go @@ -4,3 +4,13 @@ import "errors" // ErrNotAvailable indicates that a feature is not available var ErrNotAvailable = errors.New("not available") + +// ErrTimeout is the error returned when a timeout happened. +// Modeled after context.DeadlineError +var ErrTimeout error = errTimeoutError{} + +type errTimeoutError struct{} + +func (errTimeoutError) Error() string { return "timeout" } +func (errTimeoutError) Timeout() bool { return true } +func (errTimeoutError) Temporary() bool { return true } diff --git a/provider/mqtt.go b/provider/mqtt.go index 0894e2e95..ddbe22ef3 100644 --- a/provider/mqtt.go +++ b/provider/mqtt.go @@ -67,7 +67,6 @@ func NewMqtt(log *util.Logger, client *mqtt.Client, topic string, payload string // FloatGetter creates handler for float64 from MQTT topic that returns cached value func (m *Mqtt) FloatGetter() func() (float64, error) { h := &msgHandler{ - log: m.log, topic: m.topic, scale: m.scale, mux: util.NewWaiter(m.timeout, func() { m.log.TRACE.Printf("%s wait for initial value", m.topic) }), @@ -80,7 +79,6 @@ func (m *Mqtt) FloatGetter() func() (float64, error) { // IntGetter creates handler for int64 from MQTT topic that returns cached value func (m *Mqtt) IntGetter() func() (int64, error) { h := &msgHandler{ - log: m.log, topic: m.topic, scale: float64(m.scale), mux: util.NewWaiter(m.timeout, func() { m.log.TRACE.Printf("%s wait for initial value", m.topic) }), @@ -93,7 +91,6 @@ func (m *Mqtt) IntGetter() func() (int64, error) { // StringGetter creates handler for string from MQTT topic that returns cached value func (m *Mqtt) StringGetter() func() (string, error) { h := &msgHandler{ - log: m.log, topic: m.topic, mux: util.NewWaiter(m.timeout, func() { m.log.TRACE.Printf("%s wait for initial value", m.topic) }), } @@ -105,7 +102,6 @@ func (m *Mqtt) StringGetter() func() (string, error) { // BoolGetter creates handler for string from MQTT topic that returns cached value func (m *Mqtt) BoolGetter() func() (bool, error) { h := &msgHandler{ - log: m.log, topic: m.topic, mux: util.NewWaiter(m.timeout, func() { m.log.TRACE.Printf("%s wait for initial value", m.topic) }), } @@ -122,7 +118,6 @@ func (m *Mqtt) IntSetter(param string) func(int64) error { return err } - m.log.TRACE.Printf("send %s: '%s'", m.topic, payload) return m.client.Publish(m.topic, false, payload) } } @@ -135,13 +130,23 @@ func (m *Mqtt) BoolSetter(param string) func(bool) error { return err } - m.log.TRACE.Printf("send %s: '%s'", m.topic, payload) + return m.client.Publish(m.topic, false, payload) + } +} + +// StringSetter invokes script with parameter replaced by string value +func (m *Mqtt) StringSetter(param string) func(string) error { + return func(v string) error { + payload, err := setFormattedValue(m.payload, param, v) + if err != nil { + return err + } + return m.client.Publish(m.topic, false, payload) } } type msgHandler struct { - log *util.Logger mux *util.Waiter scale float64 topic string @@ -149,8 +154,6 @@ type msgHandler struct { } func (h *msgHandler) receive(payload string) { - h.log.TRACE.Printf("recv %s: '%s'", h.topic, payload) - h.mux.Lock() defer h.mux.Unlock() diff --git a/provider/mqtt/client.go b/provider/mqtt/client.go index f82920ce3..bb63c6324 100644 --- a/provider/mqtt/client.go +++ b/provider/mqtt/client.go @@ -6,6 +6,7 @@ import ( "sync" "time" + "github.com/andig/evcc/api" "github.com/andig/evcc/util" mqtt "github.com/eclipse/paho.mqtt.golang" ) @@ -92,13 +93,14 @@ func (m *Client) ConnectionHandler(client mqtt.Client) { } } -// Publish synchronously pulishes payload using client qos +// Publish synchronously publishes payload using client qos func (m *Client) Publish(topic string, retained bool, payload interface{}) error { + m.log.TRACE.Printf("send %s: '%v'", topic, payload) token := m.Client.Publish(topic, m.Qos, retained, payload) if token.WaitTimeout(publishTimeout) { return token.Error() } - return nil + return api.ErrTimeout } // Listen validates uniqueness and registers and attaches listener @@ -113,14 +115,15 @@ func (m *Client) Listen(topic string, callback func(string)) { // listen attaches listener to topic func (m *Client) listen(topic string) { token := m.Client.Subscribe(topic, m.Qos, func(c mqtt.Client, msg mqtt.Message) { - s := string(msg.Payload()) - if len(s) > 0 { + payload := string(msg.Payload()) + m.log.TRACE.Printf("recv %s: '%v'", topic, payload) + if len(payload) > 0 { m.mux.Lock() callbacks := m.listener[topic] m.mux.Unlock() for _, cb := range callbacks { - cb(s) + cb(payload) } } }) diff --git a/provider/mqtt/registry.go b/provider/mqtt/registry.go index 893ed06c6..7960165bb 100644 --- a/provider/mqtt/registry.go +++ b/provider/mqtt/registry.go @@ -47,12 +47,12 @@ func RegisteredClientOrDefault(log *util.Logger, cc Config) (*Client, error) { var err error client := Instance - if client == nil && cc.Broker == "" { - return nil, errors.New("missing mqtt broker configuration") + if cc.Broker != "" { + client, err = RegisteredClient(log, cc.Broker, cc.User, cc.Password, ClientID(), 1) } if client == nil { - client, err = RegisteredClient(log, cc.Broker, cc.User, cc.Password, ClientID(), 1) + err = errors.New("missing mqtt broker configuration") } return client, err