From 93eb81a2ebdb50dcd213e5aeff2a6904e4c3b63a Mon Sep 17 00:00:00 2001 From: andig Date: Thu, 4 May 2023 08:24:49 +0200 Subject: [PATCH] Revert "Mqtt: handle messages asynchronously (#7687)" This reverts commit 77a8a323b3d79c6cf37aff59df7af1889a5b6f0b. --- provider/mqtt/client.go | 28 +++++++++++++++------------- server/mqtt.go | 2 +- 2 files changed, 16 insertions(+), 14 deletions(-) diff --git a/provider/mqtt/client.go b/provider/mqtt/client.go index 4373822be..afdcaf656 100644 --- a/provider/mqtt/client.go +++ b/provider/mqtt/client.go @@ -4,12 +4,12 @@ import ( "crypto/tls" "fmt" "math/rand" + "os" "strings" "sync" "time" paho "github.com/eclipse/paho.mqtt.golang" - "github.com/evcc-io/evcc/api" "github.com/evcc-io/evcc/util" ) @@ -129,7 +129,7 @@ func (m *Client) Publish(topic string, retained bool, payload interface{}) error if token.WaitTimeout(publishTimeout) { return token.Error() } - return api.ErrTimeout + return os.ErrDeadlineExceeded } // Listen validates uniqueness and registers and attaches listener @@ -146,7 +146,7 @@ func (m *Client) ListenSetter(topic string, callback func(string)) { m.Listen(topic, func(payload string) { callback(payload) if err := m.Publish(topic, true, ""); err != nil { - m.log.ERROR.Printf("clear: %s: %v", topic, err) + m.log.ERROR.Printf("clear: %v", err) } }) } @@ -158,22 +158,24 @@ func (m *Client) listen(topic string) { m.log.TRACE.Printf("recv %s: '%v'", topic, payload) if len(payload) > 0 { m.mux.Lock() - for _, cb := range m.listener[topic] { - go cb(payload) - } + callbacks := m.listener[topic] m.mux.Unlock() + + for _, cb := range callbacks { + cb(payload) + } } }) - m.WaitForToken("subscribe", topic, token) + m.WaitForToken(token) } // WaitForToken synchronously waits until token operation completed -func (m *Client) WaitForToken(action, topic string, token paho.Token) { - err := api.ErrTimeout +func (m *Client) WaitForToken(token paho.Token) { if token.WaitTimeout(publishTimeout) { - err = token.Error() - } - if err != nil { - m.log.ERROR.Printf("%s: %s: %v", action, topic, err) + if token.Error() != nil { + m.log.ERROR.Printf("error: %s", token.Error()) + } + } else { + m.log.DEBUG.Println("timeout") } } diff --git a/server/mqtt.go b/server/mqtt.go index 346e5efb3..9fc8c1850 100644 --- a/server/mqtt.go +++ b/server/mqtt.go @@ -62,7 +62,7 @@ func (m *MQTT) encode(v interface{}) string { func (m *MQTT) publishSingleValue(topic string, retained bool, payload interface{}) { token := m.Handler.Client.Publish(topic, m.Handler.Qos, retained, m.encode(payload)) - go m.Handler.WaitForToken("send", topic, token) + go m.Handler.WaitForToken(token) } func (m *MQTT) publish(topic string, retained bool, payload interface{}) {