From 5bdf309d8ad0753e09a0d56235160b686120d06b Mon Sep 17 00:00:00 2001 From: andig Date: Tue, 3 Oct 2023 11:19:30 +0200 Subject: [PATCH] Mqtt: limit number of inflight token waits (#10131) --- provider/mqtt/client.go | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/provider/mqtt/client.go b/provider/mqtt/client.go index f0042a83e..658a34905 100644 --- a/provider/mqtt/client.go +++ b/provider/mqtt/client.go @@ -6,6 +6,7 @@ import ( "math/rand" "strings" "sync" + "sync/atomic" paho "github.com/eclipse/paho.mqtt.golang" "github.com/evcc-io/evcc/api" @@ -38,6 +39,7 @@ type Client struct { Client paho.Client broker string Qos byte + inflight uint32 listener map[string][]func(string) } @@ -168,6 +170,14 @@ func (m *Client) listen(topic string) { // WaitForToken synchronously waits until token operation completed func (m *Client) WaitForToken(action, topic string, token paho.Token) { + if inflight := atomic.LoadUint32(&m.inflight); inflight > 64 { + return + } + + // track inflight token waits + atomic.AddUint32(&m.inflight, 1) + defer atomic.AddUint32(&m.inflight, ^uint32(0)) + err := api.ErrTimeout if token.WaitTimeout(request.Timeout) { err = token.Error()