Cardata: add forced refresh (#24777)

This commit is contained in:
andig 2025-10-30 09:00:33 +01:00 • committed by GitHub
parent 4be717f684
commit 224e1fbdc6
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 93 additions and 60 deletions

View file

@ -52,7 +52,9 @@ params:
required: true
- preset: vehicle-features
- name: streaming
default: true
help:
en: Enable if vehicle sends streaming updates during charging.
de: Aktivieren falls das Fahrzeug Streaming Updates während des Ladevorgangs schickt.
render: |
type: cardata
vin: {{ .vin }}

View file

@ -97,8 +97,8 @@ func (v *API) DeleteContainer(id string) error {
return v.DoJSON(req, &res)
}
func (v *API) GetTelematics(vin, container string) (TelematicData, error) {
var res TelematicData
func (v *API) GetTelematics(vin, container string) (ContainerContents, error) {
var res ContainerContents
uri := fmt.Sprintf(ApiURL+"/customers/vehicles/%s/telematicData?containerId=%s", vin, container)
err := v.GetJSON(uri, &res)
return res, err

View file

@ -2,6 +2,7 @@ package cardata
import (
"context"
"fmt"
"maps"
"slices"
"sync"
@ -23,19 +24,24 @@ type Provider struct {
api *API
ts oauth2.TokenSource
vin string
vin string
container string
initial map[string]TelematicDataPoint
rest map[string]TelematicData
streaming map[string]StreamingData
updated time.Time
cache time.Duration
}
// NewProvider creates a vehicle api provider
func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.TokenSource, clientID, vin string) *Provider {
func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.TokenSource, clientID, vin string, cache time.Duration) *Provider {
v := &Provider{
log: log,
api: api,
ts: ts,
vin: vin,
cache: cache,
rest: make(map[string]TelematicData),
streaming: make(map[string]StreamingData),
}
@ -51,6 +57,7 @@ func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.Toke
for msg := range recvC {
v.mu.Lock()
maps.Copy(v.streaming, msg.Data)
v.updated = time.Now()
v.mu.Unlock()
}
}()
@ -58,48 +65,7 @@ func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.Toke
return v
}
func (v *Provider) any(key string) (any, error) {
v.mu.Lock()
defer v.mu.Unlock()
if a, ok := v.streaming[key]; ok {
return a.Value, nil
}
if v.initial == nil {
// don't try as long as there's no token
if _, err := v.ts.Token(); err != nil {
return nil, api.ErrNotAvailable
}
defer func() {
if v.initial == nil {
v.initial = make(map[string]TelematicDataPoint)
}
}()
container, err := v.ensureContainer()
if err != nil {
v.log.ERROR.Printf("get container: %v", err)
return nil, api.ErrNotAvailable
}
if res, err := v.api.GetTelematics(v.vin, container); err == nil {
v.initial = res.TelematicData
} else {
v.log.ERROR.Printf("get telematics: %v", err)
return nil, api.ErrNotAvailable
}
}
if el, ok := v.initial[key]; ok {
return el.Value, nil
}
return nil, api.ErrNotAvailable
}
func (v *Provider) ensureContainer() (string, error) {
func (v *Provider) findOrCreateContainer() (string, error) {
containers, err := v.api.GetContainers()
if err != nil {
return "", err
@ -119,6 +85,62 @@ func (v *Provider) ensureContainer() (string, error) {
return res.ContainerId, err
}
func (v *Provider) setupContainer() error {
container, err := v.findOrCreateContainer()
if err != nil {
return fmt.Errorf("get container: %v", err)
}
v.container = container
return nil
}
func (v *Provider) updateContainerData() error {
res, err := v.api.GetTelematics(v.vin, v.container)
if err != nil {
return fmt.Errorf("get telematics: %v", err)
}
v.rest = res.TelematicData
v.streaming = make(map[string]StreamingData) // reset streaming
return nil
}
func (v *Provider) any(key string) (any, error) {
v.mu.Lock()
defer v.mu.Unlock()
_, tokenErr := v.ts.Token()
switch {
case tokenErr == nil && v.updated.IsZero():
// this will only happen once
if err := v.setupContainer(); err != nil {
v.log.WARN.Println(err)
}
fallthrough
case tokenErr == nil && time.Since(v.updated) > v.cache && v.container != "":
if err := v.updateContainerData(); err != nil {
v.log.WARN.Println(err)
}
v.updated = time.Now()
}
if a, ok := v.streaming[key]; ok {
return a.Value, nil
}
if el, ok := v.rest[key]; ok {
return el.Value, nil
}
return nil, api.ErrNotAvailable
}
func (v *Provider) String(key string) (string, error) {
res, err := v.any(key)
if err != nil {

View file

@ -3,6 +3,7 @@ package cardata
import (
"context"
"testing"
"time"
"github.com/evcc-io/evcc/util"
"github.com/stretchr/testify/require"
@ -15,10 +16,13 @@ func TestCardataStreaming(t *testing.T) {
p := NewProvider(ctx, util.NewLogger("foo"), nil, oauth2.StaticTokenSource(&oauth2.Token{
AccessToken: "at",
}), "client", "vin")
}), "client", "vin", 0)
// prevent container panic
p.updated = time.Now()
keySoc := "vehicle.drivetrain.batteryManagement.header"
p.initial = map[string]TelematicDataPoint{
p.rest = map[string]TelematicData{
keySoc: {Value: "42"},
}

View file

@ -21,14 +21,14 @@ type CreateContainer struct {
TechnicalDescriptors []string `json:"technicalDescriptors"`
}
type TelematicDataPoint struct {
Timestamp time.Time
Unit string
Value string
type ContainerContents struct {
TelematicData map[string]TelematicData
}
type TelematicData struct {
TelematicData map[string]TelematicDataPoint
Timestamp time.Time
Unit string
Value string
}
type StreamingMessage struct {

View file

@ -4,6 +4,7 @@ import (
"context"
"encoding/json"
"errors"
"slices"
"time"
"github.com/evcc-io/evcc/api"
@ -25,12 +26,10 @@ func init() {
// NewCardataFromConfig creates a new BMW/Mini CarData vehicle
func NewCardataFromConfig(ctx context.Context, other map[string]interface{}) (api.Vehicle, error) {
cc := struct {
var cc struct {
embed `mapstructure:",squash"`
ClientID, VIN string
Cache time.Duration
}{
Cache: 30 * time.Minute, // 50 requests per day
Cache time.Duration // 50 requests per day
}
if err := util.DecodeOther(other, &cc); err != nil {
@ -45,6 +44,12 @@ func NewCardataFromConfig(ctx context.Context, other map[string]interface{}) (ap
return nil, api.ErrMissingCredentials
}
if cc.Cache == 0 {
// for non-streaming use 15m, access controlled by loadpoint
isStreaming := slices.Contains(cc.embed.Features(), api.Streaming)
cc.Cache = map[bool]time.Duration{false: 15 * time.Minute, true: 30 * time.Minute}[isStreaming]
}
v := &Cardata{
embed: &cc.embed,
}
@ -77,7 +82,7 @@ func NewCardataFromConfig(ctx context.Context, other map[string]interface{}) (ap
api := cardata.NewAPI(log, ts)
v.Provider = cardata.NewProvider(ctx, log, api, ts, cc.ClientID, cc.VIN)
v.Provider = cardata.NewProvider(ctx, log, api, ts, cc.ClientID, cc.VIN, cc.Cache)
return v, nil
}