diff --git a/provider/mqtt/client.go b/provider/mqtt/client.go index afdcaf656..4373822be 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 os.ErrDeadlineExceeded + return api.ErrTimeout } // 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: %v", err) + m.log.ERROR.Printf("clear: %s: %v", topic, err) } }) } @@ -158,24 +158,22 @@ func (m *Client) listen(topic string) { 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(payload) + for _, cb := range m.listener[topic] { + go cb(payload) } + m.mux.Unlock() } }) - m.WaitForToken(token) + m.WaitForToken("subscribe", topic, token) } // WaitForToken synchronously waits until token operation completed -func (m *Client) WaitForToken(token paho.Token) { +func (m *Client) WaitForToken(action, topic string, token paho.Token) { + err := api.ErrTimeout if token.WaitTimeout(publishTimeout) { - if token.Error() != nil { - m.log.ERROR.Printf("error: %s", token.Error()) - } - } else { - m.log.DEBUG.Println("timeout") + err = token.Error() + } + if err != nil { + m.log.ERROR.Printf("%s: %s: %v", action, topic, err) } } diff --git a/server/mqtt.go b/server/mqtt.go index 9fc8c1850..346e5efb3 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(token) + go m.Handler.WaitForToken("send", topic, token) } func (m *MQTT) publish(topic string, retained bool, payload interface{}) {