Revert "Mqtt: handle messages asynchronously (#7687)"
This reverts commit 77a8a323b3.
This commit is contained in:
parent
1df1a6ea81
commit
93eb81a2eb
2 changed files with 16 additions and 14 deletions
|
|
@ -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")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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{}) {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue