From 6b053bf77126c3d62fa289c45ad0dc6b032dca76 Mon Sep 17 00:00:00 2001 From: andig Date: Fri, 15 May 2026 11:50:23 +0200 Subject: [PATCH] PV: track energy metrics and apply forecast scaling to optimizer (#29784) --- core/keys/site.go | 2 - core/metrics/collector.go | 1 + core/metrics/collector_test.go | 26 +++++++++++ core/site.go | 79 ++++++++++++---------------------- core/site_optimizer.go | 10 ++--- core/site_tariffs.go | 68 ++++++++++++++++++----------- 6 files changed, 101 insertions(+), 85 deletions(-) diff --git a/core/keys/site.go b/core/keys/site.go index 347693b43..8ea653b98 100644 --- a/core/keys/site.go +++ b/core/keys/site.go @@ -20,8 +20,6 @@ const ( SmartCostType = "smartCostType" Statistics = "statistics" Forecast = "forecast" - SolarAccYield = "solarAccYield" - SolarAccForecast = "solarAccForecast" TariffCo2 = "tariffCo2" TariffCo2Home = "tariffCo2Home" TariffCo2Loadpoints = "tariffCo2Loadpoints" diff --git a/core/metrics/collector.go b/core/metrics/collector.go index 2a18ca9ac..50d285aa8 100644 --- a/core/metrics/collector.go +++ b/core/metrics/collector.go @@ -14,6 +14,7 @@ const ( PV = "pv" Home = "home" // meter and group (virtual measurement) Loadpoint = "loadpoint" + Forecast = "forecast" ) type Collector struct { diff --git a/core/metrics/collector_test.go b/core/metrics/collector_test.go index 2eb3a315b..6263af368 100644 --- a/core/metrics/collector_test.go +++ b/core/metrics/collector_test.go @@ -87,6 +87,32 @@ func TestCollectorAddEnergyWithImportMeterAndExport(t *testing.T) { require.InDelta(t, 600.0*3/60/1e3, col.accu.Exported(), 1e-10) // 0.03 kWh } +func TestCollectorAddEnergyWithExportMeterAndImport(t *testing.T) { + clock := clock.NewMock() + + require.NoError(t, db.NewInstance("sqlite", ":memory:")) + require.NoError(t, SetupSchema()) + + col, err := NewCollector("baz2", "baz2", WithClock(clock)) + require.NoError(t, err) + + // seed export meter + clock.Add(3 * time.Minute) + require.NoError(t, col.AddEnergy(nil, new(1000.0), 0)) + + // negative power: export via meter delta, no import + clock.Add(3 * time.Minute) + require.NoError(t, col.AddEnergy(nil, new(1000.3), -500)) + require.InDelta(t, 0.3, col.accu.Exported(), 1e-10) + require.Equal(t, 0.0, col.accu.Imported()) + + // positive power: export via meter (no change), import via power integration + clock.Add(3 * time.Minute) + require.NoError(t, col.AddEnergy(nil, new(1000.3), 600)) + require.InDelta(t, 0.3, col.accu.Exported(), 1e-10) + require.InDelta(t, 600.0*3/60/1e3, col.accu.Imported(), 1e-10) // 0.03 kWh +} + func TestCollectorAddEnergyWithBothMeters(t *testing.T) { clock := clock.NewMock() diff --git a/core/site.go b/core/site.go index 8f3c39d14..984fdc4b6 100644 --- a/core/site.go +++ b/core/site.go @@ -84,9 +84,10 @@ type Site struct { coordinator *coordinator.Coordinator // Vehicles prioritizer *prioritizer.Prioritizer // Power budgets stats *Stats // Stats - fcstEnergy *metrics.Accumulator - pvEnergy map[string]*metrics.Accumulator + // metrics + fcstEnergy *metrics.Collector + pvEnergy map[string]*metrics.Collector homeEnergy, gridEnergy *metrics.Collector batteryEnergy map[string]*metrics.Collector // per-battery, keyed by meter ref @@ -206,10 +207,21 @@ func (site *Site) Boot(log *util.Logger, loadpoints []*Loadpoint, tariffs *tarif } site.pvMeters = append(site.pvMeters, dev) - // accumulator - site.pvEnergy[ref] = metrics.NewAccumulator() + // energy collector (for history persistence and forecast scaling) + me, err := metrics.NewCollector(metrics.PV, ref) + if err != nil { + return err + } + site.pvEnergy[ref] = me } + // solar forecast collector (mirrors PV history shape, used for scale lookup) + fc, err := metrics.NewCollector(metrics.Forecast, metrics.Forecast) + if err != nil { + return err + } + site.fcstEnergy = fc + // multiple batteries for _, ref := range site.Meters.BatteryMetersRef { dev, err := config.Meters().ByName(ref) @@ -260,9 +272,8 @@ func NewSite() *Site { site := &Site{ log: util.NewLogger("site"), Voltage: 230, // V - pvEnergy: make(map[string]*metrics.Accumulator), + pvEnergy: make(map[string]*metrics.Collector), batteryEnergy: make(map[string]*metrics.Collector), - fcstEnergy: metrics.NewAccumulator(), } return site @@ -329,36 +340,10 @@ func (site *Site) restoreSettings() error { } } - // restore accumulated energy - pvEnergy := make(map[string]metrics.Accumulator) - fcstEnergy, err := settings.Float(keys.SolarAccForecast) - - if err == nil && settings.Json(keys.SolarAccYield, &pvEnergy) == nil { - var nok bool - for _, name := range site.Meters.PVMetersRef { - if fcst, ok := pvEnergy[name]; ok { - site.pvEnergy[name].Import = fcst.Import - } else { - nok = true - site.log.WARN.Printf("accumulated solar yield: cannot restore %s", name) - } - } - - if !nok { - site.fcstEnergy.Import = fcstEnergy - site.log.DEBUG.Printf("accumulated solar yield: restored %.3fkWh forecasted, %+v produced", fcstEnergy, pvEnergy) - } else { - // reset metrics - site.log.WARN.Printf("accumulated solar yield: metrics reset") - - settings.Delete(keys.SolarAccForecast) - settings.Delete(keys.SolarAccYield) - - for _, pe := range site.pvEnergy { - pe.Import = 0 - } - } - } + // drop legacy accumulator-based forecast settings (now stored via metrics collector) + settings.Delete("solarAccForecast") + settings.Delete("solarAccYield") + settings.Delete("solarAccDay") return nil } @@ -593,27 +578,17 @@ func (site *Site) updatePvMeters() { site.publish(keys.PvEnergy, totalEnergy) site.publish(keys.Pv, mm) - // update solar yield + // persist per-meter PV energy slots (used for history and forecast scaling) for i, dev := range site.pvMeters { - // use stored devices, not ui-updated instances! - name := dev.Config().Name + c := site.pvEnergy[dev.Config().Name] - prev := site.pvEnergy[name].Imported() + var importEnergy *float64 if mm[i].Energy > 0 { - site.log.DEBUG.Printf("!! solar production: accumulate set %s %.3fkWh meter total (was: %s)", name, mm[i].Energy, site.pvEnergy[name]) - site.pvEnergy[name].SetImportMeterTotal(mm[i].Energy) - } else { - site.log.DEBUG.Printf("!! solar production: accumulate add %s %.3fW power (was: %s)", name, mm[i].Energy, site.pvEnergy[name]) - site.pvEnergy[name].AddPower(mm[i].Power) + importEnergy = &mm[i].Energy } - site.log.DEBUG.Printf("!! solar production: accumulate moved %s from %.3f to %.3f", name, prev, site.pvEnergy[name].Imported()) - } - // store - if err := settings.SetJson(keys.SolarAccYield, site.pvEnergy); err != nil { - site.log.ERROR.Println("accumulated solar production:", err) - for k, v := range site.pvEnergy { - site.log.ERROR.Printf("!! %s: %+v", k, v) + if err := c.AddEnergy(importEnergy, nil, mm[i].Power); err != nil { + site.log.ERROR.Printf("persist pv %d energy: %v", i+1, err) } } } diff --git a/core/site_optimizer.go b/core/site_optimizer.go index 85036ca20..4f5be3021 100644 --- a/core/site_optimizer.go +++ b/core/site_optimizer.go @@ -157,7 +157,7 @@ func (site *Site) optimizerUpdate(battery []types.Measurement) error { return err } - ft = prorate(scaleAndPrune(solarEnergy, 1, minLen), firstSlotDuration) + ft = prorate(scaleAndPrune(solarEnergy, site.solarScale(), minLen), firstSlotDuration) } req := optimizer.OptimizationInput{ @@ -171,8 +171,8 @@ func (site *Site) optimizerUpdate(battery []types.Measurement) error { Dt: dt, Gt: prorate(gt, firstSlotDuration), Ft: ft, - PN: scaleAndPrune(grid, 1e3, minLen), - PE: scaleAndPrune(feedIn, 1e3, minLen), + PN: scaleAndPrune(grid, 0.001, minLen), + PE: scaleAndPrune(feedIn, 0.001, minLen), }, } @@ -635,11 +635,11 @@ func asTimestamps(dt []int, now time.Time) []time.Time { return res } -func scaleAndPrune(rates api.Rates, div float64, maxLen int) []float32 { +func scaleAndPrune(rates api.Rates, scale float64, maxLen int) []float32 { res := make([]float32, 0, maxLen) for _, slot := range rates { - res = append(res, float32(slot.Value/div)) + res = append(res, float32(slot.Value*scale)) if len(res) >= maxLen { break } diff --git a/core/site_tariffs.go b/core/site_tariffs.go index 5c95ec4be..da3684894 100644 --- a/core/site_tariffs.go +++ b/core/site_tariffs.go @@ -1,19 +1,15 @@ package core import ( - "maps" "math" - "slices" "time" "github.com/evcc-io/evcc/api" "github.com/evcc-io/evcc/core/keys" "github.com/evcc-io/evcc/core/metrics" - "github.com/evcc-io/evcc/server/db/settings" "github.com/evcc-io/evcc/tariff" "github.com/evcc-io/evcc/util" "github.com/jinzhu/now" - "github.com/samber/lo" ) type solarDetails struct { @@ -150,34 +146,54 @@ func (site *Site) solarDetails(solar api.Rates) solarDetails { Complete: !last.Before(eot.AddDate(0, 0, 1)), } - // accumulate forecasted energy since last update - fcstUpdated := site.fcstEnergy.Updated() - energy := solarEnergy(solar, fcstUpdated, time.Now()) / 1e3 - site.log.DEBUG.Printf("solar forecast: accumulated %.3fWh from %v to %v", - energy, fcstUpdated.Truncate(time.Second), time.Now().Truncate(time.Second), - ) - - site.fcstEnergy.AddImportEnergy(energy) - settings.SetFloat(keys.SolarAccForecast, site.fcstEnergy.Imported()) - - produced := lo.SumBy(slices.Collect(maps.Values(site.pvEnergy)), func(v *metrics.Accumulator) float64 { - return v.Imported() - }) - site.log.DEBUG.Printf("solar forecast: produced %.3f", produced) - - if fcst := site.fcstEnergy.Imported(); fcst > 0 { - scale := produced / fcst - site.log.DEBUG.Printf("solar forecast: accumulated %.3fkWh, produced %.3fkWh, scale %.3f", fcst, produced, scale) - - const minEnergy = 0.5 // kWh - if produced+fcst > minEnergy { - res.Scale = new(scale) + if r, err := solar.At(time.Now()); err == nil { + if err := site.fcstEnergy.AddEnergy(nil, nil, r.Value); err != nil { + site.log.ERROR.Printf("solar forecast collector: %v", err) } } + if scale := site.solarScale(); scale != 1 { + res.Scale = &scale + } + return res } +// solarScale returns the ratio of produced solar energy to forecasted solar +// energy for the current day, queried from the metrics database. Used to +// adjust forecasts when PV is consistently under-/over-producing relative +// to the forecast. Returns 1.0 when not enough data is available to make +// the ratio meaningful. +func (site *Site) solarScale() float64 { + series, err := metrics.QueryImportEnergy(now.BeginningOfDay(), time.Now(), "day", true) + if err != nil { + site.log.ERROR.Printf("solar forecast scale: %v", err) + return 1 + } + + var pv, fcst float64 + for _, s := range series { + if len(s.Data) == 0 { + continue + } + switch s.Group { + case metrics.PV: + pv = s.Data[0].Import + case metrics.Forecast: + fcst = s.Data[0].Import + } + } + + const minEnergy = 0.5 // kWh + if fcst <= 0 || pv+fcst <= minEnergy { + return 1 + } + + scale := pv / fcst + site.log.DEBUG.Printf("solar forecast: produced %.3fkWh, forecasted %.3fkWh, scale %.3f", pv, fcst, scale) + return scale +} + func (site *Site) isDynamicTariff(usage api.TariffUsage) bool { tariff := site.GetTariff(usage) return tariff != nil && tariff.Type() != api.TariffTypePriceStatic