Mqtt: handle messages asynchronously (#7687)

This commit is contained in:
andig 2023-04-27 14:17:29 +02:00 • committed by GitHub
parent 9324f0c518
commit 77a8a323b3
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
2 changed files with 14 additions and 16 deletions

View file

@ -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)
}
}

View file

@ -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{}) {