DriveSomethingGreater: single store per account (#30430)
This commit is contained in:
parent
1933be48f4
commit
5bbe1bc897
5 changed files with 127 additions and 55 deletions
|
|
@ -10,6 +10,7 @@ import (
|
|||
"net/http/cookiejar"
|
||||
"net/url"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/evcc-io/evcc/api"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
|
|
@ -46,13 +47,37 @@ type API struct {
|
|||
user, password string
|
||||
}
|
||||
|
||||
// NewAPI creates an EU Data Act client and performs the initial login
|
||||
// apiKey identifies a portal account. All vehicles of the same brand and user
|
||||
// share one authenticated client.
|
||||
type apiKey struct {
|
||||
brand brand
|
||||
user string
|
||||
}
|
||||
|
||||
var (
|
||||
apiMu sync.Mutex
|
||||
apiReg = make(map[apiKey]*API)
|
||||
)
|
||||
|
||||
// NewAPI returns the EU Data Act client for the given brand and user, performing
|
||||
// the initial login on first use. Subsequent calls for the same brand and user
|
||||
// return the already authenticated client so that several vehicles of one
|
||||
// account share a single portal session instead of competing for it.
|
||||
func NewAPI(log *util.Logger, brandName, user, password string) (*API, error) {
|
||||
b, ok := resolveBrand(brandName)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("unknown brand: %s", brandName)
|
||||
}
|
||||
|
||||
key := apiKey{brand: b, user: user}
|
||||
|
||||
apiMu.Lock()
|
||||
defer apiMu.Unlock()
|
||||
|
||||
if v, ok := apiReg[key]; ok {
|
||||
return v, nil
|
||||
}
|
||||
|
||||
v := &API{
|
||||
Helper: request.NewHelper(log),
|
||||
log: log,
|
||||
|
|
@ -65,6 +90,8 @@ func NewAPI(log *util.Logger, brandName, user, password string) (*API, error) {
|
|||
return nil, fmt.Errorf("login failed: %w", err)
|
||||
}
|
||||
|
||||
apiReg[key] = v
|
||||
|
||||
return v, nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -55,9 +55,10 @@ func TestParseDataset(t *testing.T) {
|
|||
},
|
||||
}
|
||||
|
||||
data, err := parseDataset(zipJSON(t, doc))
|
||||
vin, data, err := parseDataset(zipJSON(t, doc))
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, "WVWZZZ123", vin, "dataset vin must be returned for filtering")
|
||||
assert.Equal(t, "80", data[FieldSoc].Value, "newest timestamp must win")
|
||||
assert.Equal(t, "12345", data[FieldOdometer].Value)
|
||||
assert.Equal(t, "210", data[FieldRange].Value)
|
||||
|
|
|
|||
|
|
@ -34,18 +34,18 @@ type Provider struct {
|
|||
// NewProvider creates a vehicle api provider
|
||||
func NewProvider(api *API, vin string, cache time.Duration) *Provider {
|
||||
v := &Provider{}
|
||||
s := newStore(api, vin)
|
||||
s := sharedStore(api)
|
||||
|
||||
var cached util.Cacheable[map[string]point]
|
||||
cached = util.ResettableCached(func() (map[string]point, error) {
|
||||
ts, err := s.update()
|
||||
ts, err := s.update(vin)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !ts.IsZero() {
|
||||
time.AfterFunc(resetDelay(ts, time.Now()), cached.Reset)
|
||||
}
|
||||
return s.snapshot(), nil
|
||||
return s.snapshot(vin), nil
|
||||
}, cache)
|
||||
|
||||
v.statusG = cached.Get
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ package eudataact
|
|||
import (
|
||||
"maps"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
|
|
@ -14,47 +15,80 @@ import (
|
|||
// on the first poll of a session to seed the merged map.
|
||||
const maxBackfill = 8
|
||||
|
||||
// store holds the merged dataset state for a single vehicle across a session.
|
||||
// The portal only ever appends datasets. Each dataset is downloaded at most once;
|
||||
// its data points are merged into data, keeping the newest value per field. after
|
||||
// is the delivery time of the newest dataset already merged, so later polls only
|
||||
// fetch datasets delivered after it and nothing is downloaded twice.
|
||||
// store holds the merged dataset state for all vehicles of a single portal
|
||||
// account consisting of username and brand
|
||||
type store struct {
|
||||
mu sync.Mutex
|
||||
api *API
|
||||
vin string
|
||||
mu sync.Mutex // guards vehicles
|
||||
api *API
|
||||
vehicles map[string]*vehicleState
|
||||
}
|
||||
|
||||
// vehicleState holds the merged dataset state for a single vehicle
|
||||
type vehicleState struct {
|
||||
mu sync.Mutex // guards the fields below
|
||||
identifier string
|
||||
data map[string]point
|
||||
after time.Time
|
||||
}
|
||||
|
||||
// newStore creates an empty store for the given vehicle
|
||||
func newStore(api *API, vin string) *store {
|
||||
return &store{
|
||||
api: api,
|
||||
vin: vin,
|
||||
data: make(map[string]point),
|
||||
var (
|
||||
storeMu sync.Mutex
|
||||
storeReg = make(map[*API]*store)
|
||||
)
|
||||
|
||||
// sharedStore returns the store shared by all vehicles of the given account,
|
||||
// creating it on first use. Since NewAPI already returns one client per username
|
||||
// and brand, keying on the client yields a single store per account.
|
||||
func sharedStore(api *API) *store {
|
||||
storeMu.Lock()
|
||||
defer storeMu.Unlock()
|
||||
|
||||
if s, ok := storeReg[api]; ok {
|
||||
return s
|
||||
}
|
||||
|
||||
s := &store{
|
||||
api: api,
|
||||
vehicles: make(map[string]*vehicleState),
|
||||
}
|
||||
storeReg[api] = s
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
// update downloads any datasets delivered after the newest one already merged
|
||||
// and merges them into the map oldest to newest. It returns the newest dataset's
|
||||
// delivery time (used to schedule the next poll); the merged data is read with
|
||||
// snapshot. On the first poll only the latest maxBackfill content datasets are
|
||||
// downloaded.
|
||||
func (s *store) update() (time.Time, error) {
|
||||
// state returns the per-vehicle state for vin, creating it on first use
|
||||
func (s *store) state(vin string) *vehicleState {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
if s.identifier == "" {
|
||||
id, err := s.api.identifier(s.vin)
|
||||
v := s.vehicles[vin]
|
||||
if v == nil {
|
||||
v = &vehicleState{data: make(map[string]point)}
|
||||
s.vehicles[vin] = v
|
||||
}
|
||||
|
||||
return v
|
||||
}
|
||||
|
||||
// update downloads any datasets for vin delivered after the newest one already
|
||||
// merged and merges them into the vehicle's map oldest to newest. It returns the
|
||||
// newest dataset's delivery time (used to schedule the next poll).
|
||||
// On first poll latest maxBackfill content datasets are downloaded.
|
||||
func (s *store) update(vin string) (time.Time, error) {
|
||||
v := s.state(vin)
|
||||
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
if v.identifier == "" {
|
||||
id, err := s.api.identifier(vin)
|
||||
if err != nil {
|
||||
return time.Time{}, err
|
||||
}
|
||||
s.identifier = id
|
||||
v.identifier = id
|
||||
}
|
||||
|
||||
list, err := s.api.datasets(s.vin, s.identifier)
|
||||
list, err := s.api.datasets(vin, v.identifier)
|
||||
if err != nil {
|
||||
return time.Time{}, err
|
||||
}
|
||||
|
|
@ -74,48 +108,56 @@ func (s *store) update() (time.Time, error) {
|
|||
// on the first poll the backfilled datasets are logged once as the final
|
||||
// merged map below; afterwards each newly received dataset is logged as it
|
||||
// arrives
|
||||
initial := s.after.IsZero()
|
||||
initial := v.after.IsZero()
|
||||
|
||||
for _, d := range pending(content, s.after) {
|
||||
b, err := s.api.download(s.vin, s.identifier, d.Name)
|
||||
for _, d := range pending(content, v.after) {
|
||||
b, err := s.api.download(vin, v.identifier, d.Name)
|
||||
if err != nil {
|
||||
return newest, err
|
||||
}
|
||||
|
||||
data, err := parseDataset(b)
|
||||
dvin, data, err := parseDataset(b)
|
||||
if err != nil {
|
||||
return newest, err
|
||||
}
|
||||
|
||||
merge(s.data, data)
|
||||
|
||||
// advance the high-water mark so this dataset is never downloaded again
|
||||
if d.CreatedOn.After(s.after) {
|
||||
s.after = d.CreatedOn
|
||||
// advance the high-water mark so this dataset is never downloaded again,
|
||||
// even when it is dropped below
|
||||
if d.CreatedOn.After(v.after) {
|
||||
v.after = d.CreatedOn
|
||||
}
|
||||
|
||||
// only merge points that belong to the requested vehicle
|
||||
if dvin != "" && !strings.EqualFold(dvin, vin) {
|
||||
continue
|
||||
}
|
||||
|
||||
merge(v.data, data)
|
||||
|
||||
if !initial {
|
||||
logData(s.api.log, data)
|
||||
}
|
||||
}
|
||||
|
||||
if len(s.data) == 0 {
|
||||
if len(v.data) == 0 {
|
||||
return time.Time{}, api.ErrNotAvailable
|
||||
}
|
||||
|
||||
if initial {
|
||||
logData(s.api.log, s.data)
|
||||
logData(s.api.log, v.data)
|
||||
}
|
||||
|
||||
return newest, nil
|
||||
}
|
||||
|
||||
// snapshot returns a copy of the merged data
|
||||
func (s *store) snapshot() map[string]point {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
// snapshot returns a copy of the merged data for vin
|
||||
func (s *store) snapshot(vin string) map[string]point {
|
||||
v := s.state(vin)
|
||||
|
||||
return maps.Clone(s.data)
|
||||
v.mu.Lock()
|
||||
defer v.mu.Unlock()
|
||||
|
||||
return maps.Clone(v.data)
|
||||
}
|
||||
|
||||
// logData logs every field of data at DEBUG level, sorted by field name, with
|
||||
|
|
|
|||
|
|
@ -135,12 +135,14 @@ func contentDatasets(list []dataset) ([]dataset, error) {
|
|||
}
|
||||
|
||||
// parseDataset extracts the inner JSON document from the dataset zip archive and
|
||||
// decodes it into a map of data points keyed by the dotted data field name. On
|
||||
// duplicate field names the entry with the newest timestamp wins.
|
||||
func parseDataset(b []byte) (map[string]point, error) {
|
||||
// decodes it into the dataset's VIN and a map of data points keyed by the dotted
|
||||
// data field name. On duplicate field names the entry with the newest timestamp
|
||||
// wins. The VIN is returned so the caller can drop datasets that do not belong
|
||||
// to the requested vehicle.
|
||||
func parseDataset(b []byte) (string, map[string]point, error) {
|
||||
zr, err := zip.NewReader(bytes.NewReader(b), int64(len(b)))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return "", nil, err
|
||||
}
|
||||
|
||||
var file *zip.File
|
||||
|
|
@ -151,23 +153,23 @@ func parseDataset(b []byte) (map[string]point, error) {
|
|||
}
|
||||
}
|
||||
if file == nil {
|
||||
return nil, errors.New("no json document in dataset")
|
||||
return "", nil, errors.New("no json document in dataset")
|
||||
}
|
||||
|
||||
rc, err := file.Open()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return "", nil, err
|
||||
}
|
||||
defer rc.Close()
|
||||
|
||||
raw, err := io.ReadAll(rc)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return "", nil, err
|
||||
}
|
||||
|
||||
var ds datasetFile
|
||||
if err := json.Unmarshal(raw, &ds); err != nil {
|
||||
return nil, err
|
||||
return "", nil, err
|
||||
}
|
||||
|
||||
res := make(map[string]point, len(ds.Data))
|
||||
|
|
@ -178,7 +180,7 @@ func parseDataset(b []byte) (map[string]point, error) {
|
|||
|
||||
ts, err := time.Parse(time.RFC3339, p.TimestampUtc)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return "", nil, err
|
||||
}
|
||||
|
||||
if cur, ok := res[p.DataFieldName]; ok && cur.Timestamp.After(ts) {
|
||||
|
|
@ -188,5 +190,5 @@ func parseDataset(b []byte) (map[string]point, error) {
|
|||
res[p.DataFieldName] = point{Value: p.Value, Timestamp: ts}
|
||||
}
|
||||
|
||||
return res, nil
|
||||
return ds.VIN, res, nil
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue