From 22263b5c2e70c4f6b06cbdf2f73e39a093d90ee5 Mon Sep 17 00:00:00 2001 From: andig Date: Sun, 5 Oct 2025 22:15:30 +0200 Subject: [PATCH] Cardata: support multiple client ids and vins (#24142) --- vehicle/bmw/cardata/mqtt.go | 131 ++++++++++++++++++++++++++++++++ vehicle/bmw/cardata/provider.go | 80 ++----------------- vehicle/cardata.go | 2 +- 3 files changed, 138 insertions(+), 75 deletions(-) create mode 100644 vehicle/bmw/cardata/mqtt.go diff --git a/vehicle/bmw/cardata/mqtt.go b/vehicle/bmw/cardata/mqtt.go new file mode 100644 index 000000000..e042dbf61 --- /dev/null +++ b/vehicle/bmw/cardata/mqtt.go @@ -0,0 +1,131 @@ +package cardata + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "sync" + "time" + + "github.com/cenkalti/backoff/v4" + mqtt "github.com/eclipse/paho.mqtt.golang" + "github.com/evcc-io/evcc/util" + "golang.org/x/oauth2" +) + +type MqttConnector struct { + mu sync.RWMutex + log *util.Logger + subscriptions map[string]chan StreamingMessage +} + +var ( + mqttMu sync.Mutex + mqttConnections = make(map[string]*MqttConnector) +) + +func NewMqttConnector(ctx context.Context, log *util.Logger, clientID string, ts oauth2.TokenSource) *MqttConnector { + mqttMu.Lock() + defer mqttMu.Unlock() + + if conn, ok := mqttConnections[clientID]; ok { + return conn + } + + v := &MqttConnector{ + log: log, + subscriptions: make(map[string]chan StreamingMessage), + } + + go v.run(ctx, ts) + + mqttConnections[clientID] = v + + return v +} + +func (v *MqttConnector) Subscribe(vin string) <-chan StreamingMessage { + v.mu.Lock() + defer v.mu.Unlock() + + ch := make(chan StreamingMessage, 1) + v.subscriptions[vin] = ch + + return ch +} + +func (v *MqttConnector) run(ctx context.Context, ts oauth2.TokenSource) { + bo := backoff.NewExponentialBackOff(backoff.WithMaxInterval(time.Minute)) + + for ctx.Err() == nil { + time.Sleep(bo.NextBackOff()) + + token, err := ts.Token() + if err != nil { + if !tokenError(err) { + v.log.ERROR.Println(err) + } + + continue + } + + bo.Reset() + + if err := v.runMqtt(ctx, token); err != nil { + v.log.ERROR.Println(err) + } + } +} + +func (v *MqttConnector) runMqtt(ctx context.Context, token *oauth2.Token) error { + gcid := TokenExtra(token, "gcid") + idToken := TokenExtra(token, "id_token") + + paho := mqtt.NewClient( + mqtt.NewClientOptions(). + AddBroker(StreamingURL). + SetAutoReconnect(true). + SetUsername(gcid). + SetPassword(idToken)) + + timeout := 30 * time.Second + if t := paho.Connect(); !t.WaitTimeout(timeout) { + return errors.New("connect timeout") + } else if err := t.Error(); err != nil { + return fmt.Errorf("connect: %w", err) + } + defer paho.Disconnect(0) + + topic := fmt.Sprintf("%s/#", gcid) + + if t := paho.Subscribe(topic, 0, v.handler); !t.WaitTimeout(timeout) { + return errors.New("subcribe timeout") + } else if err := t.Error(); err != nil { + return fmt.Errorf("subscribe: %w", err) + } + + ctx, cancel := context.WithDeadline(ctx, token.Expiry) + defer cancel() + + <-ctx.Done() + + return nil +} + +func (v *MqttConnector) handler(c mqtt.Client, m mqtt.Message) { + var res StreamingMessage + if err := json.Unmarshal(m.Payload(), &res); err != nil { + v.log.ERROR.Println(m.Topic(), string(m.Payload()), err) + return + } + + v.log.TRACE.Println("recv: " + string(m.Payload())) + + v.mu.RLock() + defer v.mu.RUnlock() + + if ch, ok := v.subscriptions[res.Vin]; ok { + ch <- res + } +} diff --git a/vehicle/bmw/cardata/provider.go b/vehicle/bmw/cardata/provider.go index 32108d93d..9c6c7fb50 100644 --- a/vehicle/bmw/cardata/provider.go +++ b/vehicle/bmw/cardata/provider.go @@ -2,16 +2,11 @@ package cardata import ( "context" - "encoding/json" - "errors" - "fmt" "maps" "slices" "sync" "time" - "github.com/cenkalti/backoff/v4" - mqtt "github.com/eclipse/paho.mqtt.golang" "github.com/evcc-io/evcc/api" "github.com/evcc-io/evcc/util" "github.com/spf13/cast" @@ -34,7 +29,7 @@ type Provider struct { } // NewProvider creates a vehicle api provider -func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.TokenSource, vin string) *Provider { +func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.TokenSource, clientID, vin string) *Provider { v := &Provider{ log: log, api: api, @@ -44,81 +39,18 @@ func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.Toke } go func() { - bo := backoff.NewExponentialBackOff(backoff.WithMaxInterval(time.Minute)) + mqtt := NewMqttConnector(ctx, log, clientID, ts) - for ctx.Err() == nil { - time.Sleep(bo.NextBackOff()) - - token, err := ts.Token() - if err != nil { - if !tokenError(err) { - v.log.ERROR.Println(err) - } - - continue - } - - bo.Reset() - - if err := v.runMqtt(ctx, vin, token); err != nil { - v.log.ERROR.Println(err) - } + for msg := range mqtt.Subscribe(vin) { + v.mu.Lock() + maps.Copy(v.streaming, msg.Data) + v.mu.Unlock() } }() return v } -func (v *Provider) runMqtt(ctx context.Context, vin string, token *oauth2.Token) error { - gcid := TokenExtra(token, "gcid") - idToken := TokenExtra(token, "id_token") - - paho := mqtt.NewClient( - mqtt.NewClientOptions(). - AddBroker(StreamingURL). - SetAutoReconnect(true). - SetUsername(gcid). - SetPassword(idToken)) - - timeout := 30 * time.Second - if t := paho.Connect(); !t.WaitTimeout(timeout) { - return errors.New("connect timeout") - } else if err := t.Error(); err != nil { - return fmt.Errorf("connect: %w", err) - } - defer paho.Disconnect(0) - - topic := fmt.Sprintf("%s/%s", gcid, vin) - - if t := paho.Subscribe(topic, 0, v.handler); !t.WaitTimeout(timeout) { - return errors.New("subcribe timeout") - } else if err := t.Error(); err != nil { - return fmt.Errorf("subscribe: %w", err) - } - - ctx, cancel := context.WithDeadline(ctx, token.Expiry) - defer cancel() - - <-ctx.Done() - - return nil -} - -func (v *Provider) handler(c mqtt.Client, m mqtt.Message) { - var res StreamingMessage - if err := json.Unmarshal(m.Payload(), &res); err != nil { - v.log.ERROR.Println(m.Topic(), string(m.Payload()), err) - return - } - - v.log.TRACE.Println("recv: " + string(m.Payload())) - - v.mu.Lock() - defer v.mu.Unlock() - - maps.Copy(v.streaming, res.Data) -} - func (v *Provider) any(key string) (any, error) { v.mu.Lock() defer v.mu.Unlock() diff --git a/vehicle/cardata.go b/vehicle/cardata.go index 10c736bcb..dca542409 100644 --- a/vehicle/cardata.go +++ b/vehicle/cardata.go @@ -87,7 +87,7 @@ func NewCardataFromConfig(ctx context.Context, other map[string]interface{}) (ap return nil, err } - v.Provider = cardata.NewProvider(ctx, log, api, ts, vehicle) + v.Provider = cardata.NewProvider(ctx, log, api, ts, cc.ClientID, vehicle) return v, nil }