From f0cabc6850c7d63dfbd5eb499b69d388563b194c Mon Sep 17 00:00:00 2001 From: andig Date: Sun, 31 May 2026 16:12:42 +0200 Subject: [PATCH] Vehicle: refresh drivesomethinggreater on portal schedule (#30368) --- vehicle/vw/eudataact/api.go | 44 ++----- vehicle/vw/eudataact/eudataact_test.go | 162 +++++++++++++++++++++--- vehicle/vw/eudataact/provider.go | 58 +++++++-- vehicle/vw/eudataact/store.go | 166 +++++++++++++++++++++++++ vehicle/vw/eudataact/types.go | 93 ++++++++++---- 5 files changed, 433 insertions(+), 90 deletions(-) create mode 100644 vehicle/vw/eudataact/store.go diff --git a/vehicle/vw/eudataact/api.go b/vehicle/vw/eudataact/api.go index 30415adbf..c108d77dd 100644 --- a/vehicle/vw/eudataact/api.go +++ b/vehicle/vw/eudataact/api.go @@ -41,6 +41,7 @@ var portalHost = strings.TrimPrefix(BaseURL, "https://") // the newest dataset and decoding its flat list of data points. type API struct { *request.Helper + log *util.Logger brand brand user, password string } @@ -54,6 +55,7 @@ func NewAPI(log *util.Logger, brandName, user, password string) (*API, error) { v := &API{ Helper: request.NewHelper(log), + log: log, brand: b, user: user, password: password, @@ -261,6 +263,11 @@ func (v *API) datasets(vin, identifier string) ([]dataset, error) { b, err := v.get(uri, map[string]string{"Accept": request.JSONContent, "type": "partial"}) if err != nil { + // the portal answers 404 "No files available for this request" until the + // vehicle has delivered its first dataset + if se, ok := errors.AsType[*request.StatusError](err); ok && se.HasStatus(http.StatusNotFound) { + return nil, nil + } return nil, err } @@ -269,14 +276,14 @@ func (v *API) datasets(vin, identifier string) ([]dataset, error) { return arr, nil } - var wrap struct { + var res struct { Files []dataset `json:"files"` } - if err := json.Unmarshal(b, &wrap); err != nil { + if err := json.Unmarshal(b, &res); err != nil { return nil, err } - return wrap.Files, nil + return res.Files, nil } // download fetches the dataset zip archive @@ -284,34 +291,3 @@ func (v *API) download(vin, identifier, name string) ([]byte, error) { uri := fmt.Sprintf("%s/proxy_api/euda-apim/datadelivery/vehicles/%s/%s/download", BaseURL, vin, identifier) return v.get(uri, map[string]string{"filename": name, "type": "partial"}) } - -// Status downloads the newest dataset for the vehicle and decodes its data points -// into a map keyed by the dotted field name (e.g. "state_of_charge"). -func (v *API) Status(vin string) (map[string]string, error) { - identifier, err := v.identifier(vin) - if err != nil { - return nil, err - } - - list, err := v.datasets(vin, identifier) - if err != nil { - return nil, err - } - - name := newestDataset(list) - if name == "" { - if len(list) > 0 { - // the portal only emits "_no_content_found" placeholders while the - // vehicle is asleep and has not produced a dataset with content yet - return nil, fmt.Errorf("no dataset with content- wake the vehicle via the app: %w", api.ErrNotAvailable) - } - return nil, api.ErrNotAvailable - } - - b, err := v.download(vin, identifier, name) - if err != nil { - return nil, err - } - - return parseDataset(b) -} diff --git a/vehicle/vw/eudataact/eudataact_test.go b/vehicle/vw/eudataact/eudataact_test.go index 668a4949f..a4328a045 100644 --- a/vehicle/vw/eudataact/eudataact_test.go +++ b/vehicle/vw/eudataact/eudataact_test.go @@ -4,7 +4,9 @@ import ( "archive/zip" "bytes" "encoding/json" + "fmt" "testing" + "time" "github.com/evcc-io/evcc/api" "github.com/stretchr/testify/assert" @@ -29,29 +31,38 @@ func zipJSON(t *testing.T, doc datasetFile) []byte { return buf.Bytes() } +// testProvider returns a provider serving the given static data +func testProvider(data map[string]point) *Provider { + return &Provider{ + statusG: func() (map[string]point, error) { + return data, nil + }, + } +} + func TestParseDataset(t *testing.T) { doc := datasetFile{ VIN: "WVWZZZ123", Data: []dataPoint{ - {Key: "ffff", DataFieldName: FieldSoc, Value: "73"}, - {Key: "0001", DataFieldName: FieldSoc, Value: "80"}, // smaller key wins - {Key: "aaaa", DataFieldName: FieldOdometer, Value: "12345"}, - {Key: "bbbb", DataFieldName: FieldRange, Value: "210"}, - {Key: "cccc", DataFieldName: FieldChargingState, Value: "charging"}, - {Key: "dddd", DataFieldName: FieldPlugState, Value: "connected"}, - {Key: "eeee", DataFieldName: FieldTargetSoc, Value: "90"}, - {Key: "0002", DataFieldName: "", Value: "ignored"}, // empty field name skipped + {DataFieldName: FieldSoc, Value: "73", TimestampUtc: "2026-05-31T07:00:00Z"}, + {DataFieldName: FieldSoc, Value: "80", TimestampUtc: "2026-05-31T08:00:00Z"}, // newest timestamp wins + {DataFieldName: FieldOdometer, Value: "12345", TimestampUtc: "2026-05-31T08:00:00Z"}, + {DataFieldName: FieldRange, Value: "210", TimestampUtc: "2026-05-31T08:00:00Z"}, + {DataFieldName: FieldChargingState, Value: "charging", TimestampUtc: "2026-05-31T08:00:00Z"}, + {DataFieldName: FieldPlugState, Value: "connected", TimestampUtc: "2026-05-31T08:00:00Z"}, + {DataFieldName: FieldTargetSoc, Value: "90", TimestampUtc: "2026-05-31T08:00:00Z"}, + {DataFieldName: "", Value: "ignored"}, // empty field name skipped }, } data, err := parseDataset(zipJSON(t, doc)) require.NoError(t, err) - assert.Equal(t, "80", data[FieldSoc], "smallest key must win") - assert.Equal(t, "12345", data[FieldOdometer]) - assert.Equal(t, "210", data[FieldRange]) + 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) - p := &Provider{statusG: func() (map[string]string, error) { return data, nil }} + p := testProvider(data) soc, err := p.Soc() require.NoError(t, err) @@ -86,8 +97,11 @@ func TestStatusPlugStates(t *testing.T) { } for _, tc := range tc { - data := map[string]string{FieldPlugState: tc.plug, FieldChargingState: tc.charge} - p := &Provider{statusG: func() (map[string]string, error) { return data, nil }} + data := map[string]point{ + FieldPlugState: {Value: tc.plug}, + FieldChargingState: {Value: tc.charge}, + } + p := testProvider(data) status, err := p.Status() require.NoError(t, err) @@ -106,13 +120,121 @@ func TestResolveBrand(t *testing.T) { assert.False(t, ok) } -func TestNewestDataset(t *testing.T) { +func TestContentDatasets(t *testing.T) { list := []dataset{ - {Name: "2026-05-31T08-00.zip", CreatedOn: "2026-05-31T08:00:00Z"}, - {Name: "2026-05-31T09-00.zip", CreatedOn: "2026-05-31T09:00:00Z"}, - {Name: "2026-05-31T09-15_no_content_found.zip", CreatedOn: "2026-05-31T09:15:00Z"}, + {Name: "20260531090000_WVWZZZ.zip"}, + {Name: "20260531080000_WVWZZZ.zip"}, + {Name: "20260531091500_WVWZZZ_no_content_found.zip"}, } - assert.Equal(t, "2026-05-31T09-00.zip", newestDataset(list), "newest with content, no-content skipped") - assert.Empty(t, newestDataset([]dataset{{Name: "x_no_content_found.zip"}})) + content, err := contentDatasets(list) + require.NoError(t, err) + require.Len(t, content, 2, "no-content placeholder dropped") + assert.Equal(t, "20260531080000_WVWZZZ.zip", content[0].Name, "oldest first") + assert.Equal(t, "20260531090000_WVWZZZ.zip", content[1].Name, "newest last") + assert.Equal(t, time.Date(2026, 5, 31, 8, 0, 0, 0, time.UTC), content[0].Timestamp, "timestamp parsed") + + // no-content placeholders are skipped without parsing + empty, err := contentDatasets([]dataset{{Name: "x_no_content_found.zip"}}) + require.NoError(t, err) + assert.Empty(t, empty) + + // a content dataset with an unparseable timestamp is an error + _, err = contentDatasets([]dataset{{Name: "no-timestamp.zip"}}) + require.Error(t, err) +} + +func TestPending(t *testing.T) { + content := make([]dataset, 0, 11) + for i := range 10 { + content = append(content, dataset{ + Name: fmt.Sprintf("20260531%02d0000_WVWZZZ.zip", i), + Timestamp: time.Date(2026, 5, 31, i, 0, 0, 0, time.UTC), + }) // hour i, oldest first + } + + // first poll (zero high-water): only the latest maxBackfill + got := pending(content, time.Time{}) + require.Len(t, got, maxBackfill) + assert.Equal(t, content[len(content)-maxBackfill].Name, got[0].Name, "oldest within the backfill window") + assert.Equal(t, content[len(content)-1].Name, got[len(got)-1].Name, "newest") + + // fewer datasets than the backfill cap: all returned + assert.Len(t, pending(content[:3], time.Time{}), 3) + + // high-water at the newest merged dataset: nothing new to download + after := content[len(content)-1].Timestamp + assert.Empty(t, pending(content, after)) + + // a newer dataset arrives + content = append(content, dataset{ + Name: "20260531100000_WVWZZZ.zip", + Timestamp: time.Date(2026, 5, 31, 10, 0, 0, 0, time.UTC), + }) + got = pending(content, after) + require.Len(t, got, 1) + assert.Equal(t, "20260531100000_WVWZZZ.zip", got[0].Name, "only the newer dataset is pending") +} + +func TestMerge(t *testing.T) { + t0 := time.Date(2026, 5, 31, 7, 0, 0, 0, time.UTC) + t1 := time.Date(2026, 5, 31, 8, 0, 0, 0, time.UTC) + + dst := map[string]point{ + FieldSoc: {Value: "70", Timestamp: t0}, + FieldOdometer: {Value: "100", Timestamp: t1}, + } + src := map[string]point{ + FieldSoc: {Value: "80", Timestamp: t1}, // newer -> wins + FieldOdometer: {Value: "90", Timestamp: t0}, // older -> ignored + FieldRange: {Value: "200", Timestamp: t1}, // new field -> added + } + + merge(dst, src) + + assert.Equal(t, "80", dst[FieldSoc].Value, "newer datapoint wins") + assert.Equal(t, "100", dst[FieldOdometer].Value, "older datapoint ignored") + assert.Equal(t, "200", dst[FieldRange].Value, "new field added") +} + +func TestDatasetTime(t *testing.T) { + ref := time.Date(2026, 5, 31, 8, 0, 0, 0, time.UTC) + + tc := []struct { + d dataset + expected time.Time + err bool + }{ + {dataset{Name: "20260531080000_WVWZZZ_no_content_found.zip"}, ref, false}, // real portal format + {dataset{Name: "20260531080000_WVWZZZ.zip"}, ref, false}, // content file + {dataset{CreatedOn: "2026-05-31T08:00:00Z"}, ref, false}, // createdOn fallback + {dataset{CreatedOn: "2026-05-31T08:00:00Z", Name: "no-stamp.zip"}, ref, false}, // name unparseable, createdOn used + {dataset{Name: "no-timestamp.zip"}, time.Time{}, true}, // nothing parseable + } + + for _, tc := range tc { + got, err := tc.d.time() + if tc.err { + assert.Error(t, err, "dataset %+v", tc.d) + continue + } + require.NoError(t, err, "dataset %+v", tc.d) + assert.Equal(t, tc.expected, got, "dataset %+v", tc.d) + } +} + +// TestResetDelay verifies the cache reset is scheduled for when the portal is +// expected to deliver the dataset following the one just read. +func TestResetDelay(t *testing.T) { + now := time.Date(2026, 5, 31, 12, 0, 0, 0, time.UTC) + + // fresh dataset: reset one interval + latency later + assert.Equal(t, portalInterval+portalLatency, resetDelay(now, now)) + + // dataset already 5 min old: reset interval + latency after its timestamp + assert.Equal(t, portalInterval+portalLatency-5*time.Minute, resetDelay(now.Add(-5*time.Minute), now)) + + // next dataset already due: never reset sooner than the latency margin + assert.Equal(t, portalLatency, resetDelay(now.Add(-portalInterval), now)) + assert.Equal(t, portalLatency, resetDelay(now.Add(-time.Hour), now)) } diff --git a/vehicle/vw/eudataact/provider.go b/vehicle/vw/eudataact/provider.go index 29fae5263..53269fcd4 100644 --- a/vehicle/vw/eudataact/provider.go +++ b/vehicle/vw/eudataact/provider.go @@ -9,25 +9,65 @@ import ( "github.com/evcc-io/evcc/util" ) -// Provider implements the vehicle api on top of the EU Data Act dataset +const ( + // portalInterval is the cadence at which the portal delivers a new dataset + portalInterval = 15 * time.Minute + // portalLatency is the margin added to a dataset's timestamp before the + // following dataset is expected to be available for download + portalLatency = 30 * time.Second +) + +// Provider implements the vehicle api on top of the EU Data Act dataset. +// +// The portal is not a live api: it stores a new dataset roughly every +// portalInterval and only ever appends. Rather than re-reading the full history +// on every poll, a store downloads each dataset once and merges its data points +// into a single map, keeping the newest value per field across all datasets. The +// status getter is cached; instead of relying on the cache ttl alone, each read +// schedules a reset for the moment the next dataset is expected (the dataset's +// timestamp plus portalInterval and a latency margin), so the map is updated as +// soon as the portal delivers a new dataset. type Provider struct { - statusG func() (map[string]string, error) + statusG func() (map[string]point, error) } // NewProvider creates a vehicle api provider func NewProvider(api *API, vin string, cache time.Duration) *Provider { - return &Provider{ - statusG: util.Cached(func() (map[string]string, error) { - return api.Status(vin) - }, cache), + v := &Provider{} + s := newStore(api, vin) + + var cached util.Cacheable[map[string]point] + cached = util.ResettableCached(func() (map[string]point, error) { + ts, err := s.update() + if err != nil { + return nil, err + } + if !ts.IsZero() { + time.AfterFunc(resetDelay(ts, time.Now()), cached.Reset) + } + return s.snapshot(), nil + }, cache) + + v.statusG = cached.Get + + return v +} + +// resetDelay returns the delay until the dataset following the one delivered at +// ts is expected to be available. It never returns less than portalLatency so a +// late or repeated dataset does not cause immediate re-polling. +func resetDelay(ts, now time.Time) time.Duration { + if d := ts.Add(portalInterval + portalLatency).Sub(now); d > portalLatency { + return d } + return portalLatency } // lookup returns the first present, non-empty value among the given field names -func lookup(data map[string]string, fields ...string) (string, bool) { +func lookup(data map[string]point, fields ...string) (string, bool) { for _, f := range fields { - if v, ok := data[f]; ok && v != "" { - return v, true + if v, ok := data[f]; ok && v.Value != "" { + return v.Value, true } } return "", false diff --git a/vehicle/vw/eudataact/store.go b/vehicle/vw/eudataact/store.go new file mode 100644 index 000000000..1b4270ee7 --- /dev/null +++ b/vehicle/vw/eudataact/store.go @@ -0,0 +1,166 @@ +package eudataact + +import ( + "maps" + "slices" + "sync" + "time" + + "github.com/evcc-io/evcc/api" + "github.com/evcc-io/evcc/util" +) + +// maxBackfill bounds how many of the most recent content datasets are downloaded +// 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. +type store struct { + mu sync.Mutex + api *API + vin string + 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), + } +} + +// 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) { + s.mu.Lock() + defer s.mu.Unlock() + + if s.identifier == "" { + id, err := s.api.identifier(s.vin) + if err != nil { + return time.Time{}, err + } + s.identifier = id + } + + list, err := s.api.datasets(s.vin, s.identifier) + if err != nil { + return time.Time{}, err + } + + content, err := contentDatasets(list) + if err != nil { + return time.Time{}, err + } + + var newest time.Time + for _, d := range list { + t, err := d.time() + if err != nil { + return time.Time{}, err + } + if t.After(newest) { + newest = t + } + } + + // 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() + + for _, d := range pending(content, s.after) { + b, err := s.api.download(s.vin, s.identifier, d.Name) + if err != nil { + return newest, err + } + + 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.Timestamp.After(s.after) { + s.after = d.Timestamp + } + + if !initial { + logData(s.api.log, data) + } + } + + if len(s.data) == 0 { + return time.Time{}, api.ErrNotAvailable + } + + if initial { + logData(s.api.log, s.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() + + return maps.Clone(s.data) +} + +// logData logs every field of data at DEBUG level, sorted by field name, with +// its value and own timestamp in local time. +func logData(log *util.Logger, data map[string]point) { + for _, k := range slices.Sorted(maps.Keys(data)) { + p := data[k] + log.DEBUG.Printf("recv %s: %s (%s)", k, p.Value, p.Timestamp.Local().Format("2006-01-02 15:04:05")) + } +} + +// pending returns the content datasets that still need downloading, oldest to +// newest. content must be sorted oldest to newest. On the first poll (after is +// zero) only the latest maxBackfill datasets are returned to seed the map; +// afterwards only datasets delivered after the newest one already merged are +// returned, so no dataset is downloaded twice. +func pending(content []dataset, after time.Time) []dataset { + if after.IsZero() { + if len(content) > maxBackfill { + return content[len(content)-maxBackfill:] + } + return content + } + + res := make([]dataset, 0, len(content)) + for _, d := range content { + if d.Timestamp.After(after) { + res = append(res, d) + } + } + + return res +} + +// merge copies the data points from src into dst, keeping the newest value per +// field across datasets. +func merge(dst, src map[string]point) { + for k, p := range src { + if cur, ok := dst[k]; ok && cur.Timestamp.After(p.Timestamp) { + continue + } + dst[k] = p + } +} diff --git a/vehicle/vw/eudataact/types.go b/vehicle/vw/eudataact/types.go index f542889ee..0a0c3e806 100644 --- a/vehicle/vw/eudataact/types.go +++ b/vehicle/vw/eudataact/types.go @@ -5,8 +5,11 @@ import ( "bytes" "encoding/json" "errors" + "fmt" "io" + "slices" "strings" + "time" ) // brand holds the OIDC client id and state suffix for a VW group brand. @@ -70,25 +73,44 @@ func (v Vehicle) Name() string { return "" } -// dataset describes a single delivered dataset file +// dataset describes a single delivered dataset file. Timestamp is the parsed +// delivery time; it is populated by contentDatasets from the file name or the +// createdOn field. type dataset struct { - Name string `json:"name"` - CreatedOn string `json:"createdOn"` + Name string `json:"name"` + CreatedOn string `json:"createdOn"` + Timestamp time.Time `json:"-"` } -// sortKey returns the value used to find the newest dataset -func (d dataset) sortKey() string { - if d.CreatedOn != "" { - return d.CreatedOn +// nameTime parses the compact timestamp the portal prefixes to a dataset file +// name, e.g. 20260531102941_WAUZZZ..._no_content_found.zip. +func nameTime(name string) (time.Time, error) { + prefix, _, _ := strings.Cut(name, "_") + return time.Parse("20060102150405", prefix) +} + +// time parses the delivery time the dataset carries. The portal embeds it in the +// file name and also delivers it as the createdOn field; the file name is +// preferred and createdOn is the fallback. An error is returned when neither +// carries a parseable timestamp. +func (d dataset) time() (time.Time, error) { + if t, err := nameTime(d.Name); err == nil { + return t, nil } - return d.Name + return time.Parse(time.RFC3339, d.CreatedOn) } -// dataPoint is a single decoded telemetry value of a dataset +// dataPoint is a single data point as delivered in the dataset JSON document type dataPoint struct { - Key string `json:"key"` DataFieldName string `json:"dataFieldName"` Value string `json:"value"` + TimestampUtc string `json:"timestampUtc"` +} + +// point is a decoded data point: its value and the time it was recorded +type point struct { + Value string + Timestamp time.Time } // datasetFile is the JSON document contained in a dataset zip archive @@ -105,30 +127,42 @@ const ( FieldRangePrimary = "cruising_range_primary_engine" FieldOdometer = "mileage" FieldChargingState = "charging_state" - FieldPlugState = "charging_plug1_connectionstate" + FieldPlugState = "plug_state" FieldTargetSoc = "settings.target_soc" ) -// newestDataset returns the name of the most recent dataset that actually -// carries content. The portal emits "..._no_content_found.zip" placeholders -// while the vehicle is asleep, which are skipped. -func newestDataset(list []dataset) string { - var best dataset +// contentDatasets returns the datasets that actually carry content, with their +// delivery time parsed into Timestamp and sorted from oldest to newest. The +// portal emits "..._no_content_found.zip" placeholders while the vehicle is +// asleep, which are skipped. An error is returned when a content dataset's +// timestamp cannot be parsed. +func contentDatasets(list []dataset) ([]dataset, error) { + content := make([]dataset, 0, len(list)) for _, d := range list { if strings.HasSuffix(strings.ToLower(d.Name), "_no_content_found.zip") { continue } - if best.Name == "" || d.sortKey() > best.sortKey() { - best = d + + t, err := d.time() + if err != nil { + return nil, fmt.Errorf("dataset %q: %w", d.Name, err) } + d.Timestamp = t + + content = append(content, d) } - return best.Name + + slices.SortStableFunc(content, func(a, b dataset) int { + return a.Timestamp.Compare(b.Timestamp) + }) + + return content, nil } // parseDataset extracts the inner JSON document from the dataset zip archive and -// decodes it into a map keyed by the dotted data field name. On duplicate field -// names the entry with the smallest key uuid wins, matching the adapter. -func parseDataset(b []byte) (map[string]string, error) { +// 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) { zr, err := zip.NewReader(bytes.NewReader(b), int64(len(b))) if err != nil { return nil, err @@ -161,17 +195,22 @@ func parseDataset(b []byte) (map[string]string, error) { return nil, err } - res := make(map[string]string, len(ds.Data)) - keys := make(map[string]string, len(ds.Data)) + res := make(map[string]point, len(ds.Data)) for _, p := range ds.Data { if p.DataFieldName == "" { continue } - if k, ok := keys[p.DataFieldName]; ok && k <= p.Key { + + ts, err := time.Parse(time.RFC3339, p.TimestampUtc) + if err != nil { + return nil, err + } + + if cur, ok := res[p.DataFieldName]; ok && cur.Timestamp.After(ts) { continue } - res[p.DataFieldName] = p.Value - keys[p.DataFieldName] = p.Key + + res[p.DataFieldName] = point{Value: p.Value, Timestamp: ts} } return res, nil