From 5e6f1e7356d807785859f6015f7cbd81078f6dbc Mon Sep 17 00:00:00 2001 From: Michael Geers Date: Sat, 18 Jul 2026 18:14:23 +0200 Subject: [PATCH] Metrics: recover downtime energy via persisted meter readings (#31795) --- core/metrics/accumulator.go | 26 +++++++ core/metrics/collector.go | 40 ++++++++++- core/metrics/collector_test.go | 116 ++++++++++++++++++++++++++++++- core/metrics/db.go | 18 +++-- core/metrics/db_entities_test.go | 4 +- core/metrics/db_history_test.go | 4 +- core/metrics/db_profile.go | 2 +- core/metrics/db_test.go | 28 ++++---- core/metrics/stats_test.go | 14 ++-- 9 files changed, 215 insertions(+), 37 deletions(-) diff --git a/core/metrics/accumulator.go b/core/metrics/accumulator.go index b19644893..943af91b2 100644 --- a/core/metrics/accumulator.go +++ b/core/metrics/accumulator.go @@ -18,6 +18,32 @@ type Accumulator struct { SocTemp *float64 `json:"socTemp,omitempty"` } +// AccumulatorState is the resumable meter-reading checkpoint of an Accumulator. +type AccumulatorState struct { + EnergyMeter *float64 // kWh, last absolute reading + ReturnEnergyMeter *float64 // kWh, last absolute reading +} + +// Snapshot returns the current meter readings for persistence. +func (m *Accumulator) Snapshot() AccumulatorState { + return AccumulatorState{EnergyMeter: m.energyMeter, ReturnEnergyMeter: m.returnEnergyMeter} +} + +// Restore seeds the meter readings so the first delta covers the downtime. +func (m *Accumulator) Restore(s AccumulatorState) { + m.energyMeter = s.EnergyMeter + m.returnEnergyMeter = s.ReturnEnergyMeter +} + +// CompleteFor reports whether the state can seed a collector of the given group. +// Bidirectional groups need both readings for a complete restore. +func (s AccumulatorState) CompleteFor(group string) bool { + if group == Battery || group == Grid { + return s.EnergyMeter != nil && s.ReturnEnergyMeter != nil + } + return s.EnergyMeter != nil || s.ReturnEnergyMeter != nil +} + // setSocTemp keeps the first reading per slot. func (m *Accumulator) setSocTemp(value float64) { if m.SocTemp == nil { diff --git a/core/metrics/collector.go b/core/metrics/collector.go index 3c86aff79..4ce5c26e3 100644 --- a/core/metrics/collector.go +++ b/core/metrics/collector.go @@ -24,6 +24,8 @@ type Collector struct { entity entity accu *Accumulator started time.Time + restored bool // meter readings seeded from db + lastSlot time.Time // last persisted slot at restore, for contiguity check statsCache EnergyStats } @@ -38,6 +40,17 @@ func NewCollector(group, name, title string, opt ...func(*Accumulator)) (*Collec accu: NewAccumulator(opt...), } + // seed saved readings so the first delta covers the downtime + if state := (AccumulatorState{entity.EnergyMeter, entity.ReturnEnergyMeter}); state.CompleteFor(group) { + c.accu.Restore(state) + c.restored = true + // last persisted slot distinguishes a contiguous restart (energy stays + // time-correct) from one that skipped whole slots (inflated catchup) + var lastTs int64 + db.Instance.Model(new(meter)).Where("meter = ?", entity.Id).Select("COALESCE(max(ts), 0)").Scan(&lastTs) + c.lastSlot = time.Unix(lastTs, 0) + } + return c, nil } @@ -94,6 +107,11 @@ func (c *Collector) process(fun func()) error { switch { case c.started.IsZero(): + if c.restored { + // seeded readings make the mid-slot start complete energy-wise + c.started = slotStart + return nil + } // keep started un-truncated so a mid-slot start stays distinguishable c.started = now @@ -102,11 +120,16 @@ func (c *Collector) process(fun func()) error { // preceding slot boundary - false for the mid-slot first slot and // for a slot reached after a data gap if c.started.Equal(slotStart.Add(-tariff.SlotDuration)) { - if err := c.persist(); err != nil { + // a restore that skipped whole slots dumps the downtime energy into + // this single slot, inflating it - a contiguous restart keeps the + // slot's meter delta time-correct, so only the former is recovered + recovered := c.restored && !c.started.Equal(c.lastSlot.Add(tariff.SlotDuration)) + if err := c.persist(recovered); err != nil { return err } } + c.restored = false // only the first slot inherits recovery energy c.started = slotStart default: @@ -119,8 +142,19 @@ func (c *Collector) process(fun func()) error { return nil } -func (c *Collector) persist() error { - return persist(c.entity, c.started, c.accu.Energy, c.accu.ReturnEnergy, c.accu.SocTemp) +func (c *Collector) persist(recovered bool) error { + if err := persist(c.entity, c.started, c.accu.Energy, c.accu.ReturnEnergy, c.accu.SocTemp, recovered); err != nil { + return err + } + + // checkpoint meter readings for downtime recovery (scoped write, keeps identity intact) + s := c.accu.Snapshot() + c.entity.EnergyMeter = s.EnergyMeter + c.entity.ReturnEnergyMeter = s.ReturnEnergyMeter + return db.Instance.Model(&c.entity).UpdateColumns(map[string]any{ + "energy_meter": s.EnergyMeter, + "return_energy_meter": s.ReturnEnergyMeter, + }).Error } // SetSocTemp records the slot-start soc (temperature when isTemp). diff --git a/core/metrics/collector_test.go b/core/metrics/collector_test.go index 58f371660..5754d2eb4 100644 --- a/core/metrics/collector_test.go +++ b/core/metrics/collector_test.go @@ -266,6 +266,120 @@ func TestCollectorSkipsPartialFirstSlot(t *testing.T) { require.Equal(t, int64(15*60), m.Timestamp, "persisted slot should start at 00:15") } +// TestCollectorRecoversDowntimeViaMeterReadings verifies that saved readings +// seed a new collector so the first slot after restart contains downtime energy. +func TestCollectorRecoversDowntimeViaMeterReadings(t *testing.T) { + clk := clock.NewMock() // 1970-01-01 00:00:00 UTC, on a slot boundary + + require.NoError(t, db.NewInstance("sqlite", ":memory:")) + require.NoError(t, SetupSchema()) + + col, err := NewCollector(Grid, "restore", "", WithClock(clk)) + require.NoError(t, err) + require.False(t, col.restored, "fresh entity must not restore") + + // seed meters at slot start, advance within slot + require.NoError(t, col.AddEnergy(new(100.0), new(200.0), 0)) + clk.Add(15 * time.Minute) // 00:15 + require.NoError(t, col.AddEnergy(new(100.5), new(200.2), 0)) + + // cross into 00:30: slot 00:15 persisted, readings saved on entity + clk.Add(15 * time.Minute) // 00:30 + require.NoError(t, col.AddEnergy(new(101.0), new(200.4), 0)) + + var e entity + require.NoError(t, db.Instance.First(&e, col.entity.Id).Error) + require.Equal(t, 101.0, *e.EnergyMeter) + require.Equal(t, 200.4, *e.ReturnEnergyMeter) + + var count int64 + require.NoError(t, db.Instance.Model(new(meter)).Where("meter = ?", col.entity.Id).Count(&count).Error) + require.EqualValues(t, 2, count) + + // restart after 1h downtime, joining slot 01:30 mid-way + clk.Add(65 * time.Minute) // 01:35 + col2, err := NewCollector(Grid, "restore", "", WithClock(clk)) + require.NoError(t, err) + require.True(t, col2.restored) + + // first reading yields the delta across the downtime + require.NoError(t, col2.AddEnergy(new(103.0), new(201.4), 0)) + require.InDelta(t, 2.0, col2.accu.Energy, 1e-10) + require.InDelta(t, 1.0, col2.accu.ReturnEnergy, 1e-10) + + // cross into 01:45: catchup slot 01:30 persisted despite mid-slot start, + // downtime delta plus the 0.5 kWh accrued since restart + clk.Add(10 * time.Minute) // 01:45 + require.NoError(t, col2.AddEnergy(new(103.5), new(201.4), 0)) + + var m meter + require.NoError(t, db.Instance.Where("meter = ? AND ts = ?", col2.entity.Id, 90*60).First(&m).Error) + require.InDelta(t, 2.5, m.Energy, 1e-10) + require.InDelta(t, 1.0, m.ReturnEnergy, 1e-10) + require.True(t, m.Recovered, "catchup slot must be flagged recovered") + + // the recovered slot is excluded from the household profile + require.False(t, col2.restored, "recovery flag cleared after first slot") + var recovered int64 + require.NoError(t, db.Instance.Model(new(meter)).Where("meter = ? AND recovered", col2.entity.Id).Count(&recovered).Error) + require.EqualValues(t, 1, recovered, "only the catchup slot is recovered") + + // readings advanced with the persisted slot + require.NoError(t, db.Instance.First(&e, col2.entity.Id).Error) + require.Equal(t, 103.5, *e.EnergyMeter) + require.Equal(t, 201.4, *e.ReturnEnergyMeter) +} + +// TestCollectorRecoveryWithinCurrentSlot verifies that a restart that stays +// within the slot after the last persisted one keeps its meter delta +// time-correct: the catchup slot is persisted normally and not flagged +// recovered, so it stays in the household profile. +func TestCollectorRecoveryWithinCurrentSlot(t *testing.T) { + clk := clock.NewMock() // 1970-01-01 00:00:00 UTC, on a slot boundary + + require.NoError(t, db.NewInstance("sqlite", ":memory:")) + require.NoError(t, SetupSchema()) + + col, err := NewCollector(Grid, "within", "", WithClock(clk)) + require.NoError(t, err) + + // seed meters, advance a slot, cross into 00:30 to persist slot 00:15 + require.NoError(t, col.AddEnergy(new(100.0), new(200.0), 0)) + clk.Add(15 * time.Minute) // 00:15 + require.NoError(t, col.AddEnergy(new(100.5), new(200.2), 0)) + clk.Add(15 * time.Minute) // 00:30 + require.NoError(t, col.AddEnergy(new(101.0), new(200.4), 0)) + + // restart 10min into slot 00:30 - the slot right after the last persisted + clk.Add(10 * time.Minute) // 00:40 + col2, err := NewCollector(Grid, "within", "", WithClock(clk)) + require.NoError(t, err) + require.True(t, col2.restored) + require.EqualValues(t, 15*60, col2.lastSlot.Unix(), "last persisted slot is 00:15") + + // recovery happens during the current slot: the meter delta is applied but + // no boundary was skipped, so nothing is persisted yet + require.NoError(t, col2.AddEnergy(new(103.0), new(201.4), 0)) + require.InDelta(t, 2.0, col2.accu.Energy, 1e-10) + var count int64 + require.NoError(t, db.Instance.Model(new(meter)).Where("meter = ?", col2.entity.Id).Count(&count).Error) + require.EqualValues(t, 2, count, "current slot not persisted mid-slot") + + // cross into 00:45: slot 00:30 persisted with the full slot delta, not recovered + clk.Add(5 * time.Minute) // 00:45 + require.NoError(t, col2.AddEnergy(new(103.5), new(201.4), 0)) + + var m meter + require.NoError(t, db.Instance.Where("meter = ? AND ts = ?", col2.entity.Id, 30*60).First(&m).Error) + require.InDelta(t, 2.5, m.Energy, 1e-10) + require.False(t, m.Recovered, "contiguous restart with meter totals is not recovered") + + // no slot is excluded from the household profile + var recovered int64 + require.NoError(t, db.Instance.Model(new(meter)).Where("meter = ? AND recovered", col2.entity.Id).Count(&recovered).Error) + require.Zero(t, recovered) +} + // TestCreateEntityRefreshesTitle verifies that a second call to createEntity // with a non-empty title fills in (or updates) the title on an existing row, // and that passing an empty title never clears a previously stored value. @@ -375,7 +489,7 @@ func TestCreateEntityReconcilesExtToConsumer(t *testing.T) { // ext meter with a persisted history slot ext, err := createEntity(Meter, "db:5", "Fridge") require.NoError(t, err) - require.NoError(t, persist(ext, time.Unix(15*60, 0), 0.3, 0, nil)) + require.NoError(t, persist(ext, time.Unix(15*60, 0), 0.3, 0, nil, false)) // reconfigured as consumer: same row relabeled, history intact con, err := createEntity(Consumer, "db:5", "Fridge") diff --git a/core/metrics/db.go b/core/metrics/db.go index fd6a2a130..47998ec6e 100644 --- a/core/metrics/db.go +++ b/core/metrics/db.go @@ -17,15 +17,18 @@ type meter struct { Entity entity `json:"-" gorm:"foreignkey:Meter;references:Id"` Energy float64 `json:"energy" gorm:"column:energy"` ReturnEnergy float64 `json:"returnEnergy" gorm:"column:return_energy"` - SocTemp *float64 `json:"socTemp,omitempty" gorm:"column:soc_temp"` // at start of slot + SocTemp *float64 `json:"socTemp,omitempty" gorm:"column:soc_temp"` // at start of slot + Recovered bool `json:"recovered,omitempty" gorm:"column:recovered"` // downtime catchup slot, excluded from profile } type entity struct { - Id int `gorm:"column:id;primarykey"` - Group string `gorm:"column:group;uniqueIndex:entities_group_name"` - Name string `gorm:"column:name;uniqueIndex:entities_group_name"` - Title string `gorm:"column:title"` - IsTemp bool `gorm:"column:is_temp"` // soc_temp holds temperature, not soc + Id int `gorm:"column:id;primarykey"` + Group string `gorm:"column:group;uniqueIndex:entities_group_name"` + Name string `gorm:"column:name;uniqueIndex:entities_group_name"` + Title string `gorm:"column:title"` + IsTemp bool `gorm:"column:is_temp"` // soc_temp holds temperature, not soc + EnergyMeter *float64 `gorm:"column:energy_meter"` // kWh, at last persisted slot + ReturnEnergyMeter *float64 `gorm:"column:return_energy_meter"` // kWh, at last persisted slot } func init() { @@ -134,7 +137,7 @@ func SetupSchema() error { var OnPersist func(slot time.Time) // persist stores a completed 15min slot -func persist(entity entity, ts time.Time, energy, returnEnergy float64, socTemp *float64) error { +func persist(entity entity, ts time.Time, energy, returnEnergy float64, socTemp *float64, recovered bool) error { slot := ts.Truncate(tariff.SlotDuration) if err := db.Instance.Create(&meter{ Meter: entity.Id, @@ -142,6 +145,7 @@ func persist(entity entity, ts time.Time, energy, returnEnergy float64, socTemp Energy: energy, ReturnEnergy: returnEnergy, SocTemp: socTemp, + Recovered: recovered, }).Error; err != nil { return err } diff --git a/core/metrics/db_entities_test.go b/core/metrics/db_entities_test.go index 008edf266..673e7d214 100644 --- a/core/metrics/db_entities_test.go +++ b/core/metrics/db_entities_test.go @@ -19,8 +19,8 @@ func TestListEntities(t *testing.T) { require.NoError(t, db.Instance.Create(&pv).Error) base := time.Date(2026, 4, 15, 16, 0, 0, 0, time.Now().Location()) - require.NoError(t, persist(grid, base, 1, 0, nil)) - require.NoError(t, persist(grid, base.Add(time.Hour), 2, 0, nil)) + require.NoError(t, persist(grid, base, 1, 0, nil, false)) + require.NoError(t, persist(grid, base.Add(time.Hour), 2, 0, nil, false)) entities, err := ListEntities() require.NoError(t, err) diff --git a/core/metrics/db_history_test.go b/core/metrics/db_history_test.go index e771cb852..d2daf11a9 100644 --- a/core/metrics/db_history_test.go +++ b/core/metrics/db_history_test.go @@ -175,8 +175,8 @@ func TestQueryEnergySoc(t *testing.T) { require.NoError(t, db.Instance.Create(&e).Error) base := time.Date(2026, 4, 15, 16, 0, 0, 0, time.Now().Location()) - require.NoError(t, persist(e, base, 1, 0, new(80.0))) - require.NoError(t, persist(e, base.Add(15*time.Minute), 1, 0, new(70.0))) + require.NoError(t, persist(e, base, 1, 0, new(80.0), false)) + require.NoError(t, persist(e, base.Add(15*time.Minute), 1, 0, new(70.0), false)) from := base.Add(-time.Hour).UTC() to := base.Add(time.Hour).UTC() diff --git a/core/metrics/db_profile.go b/core/metrics/db_profile.go index 664f5a60a..2986be8a2 100644 --- a/core/metrics/db_profile.go +++ b/core/metrics/db_profile.go @@ -21,7 +21,7 @@ func energyProfile(entity entity, from time.Time) (*[96]float64, error) { // COALESCE guards against legacy rows with NULL energy rows, err := db.Query(`SELECT min(ts) AS ts, COALESCE(avg(energy), 0) AS energy FROM meters - WHERE meter = ? AND ts >= ? + WHERE meter = ? AND ts >= ? AND COALESCE(recovered, 0) = 0 GROUP BY strftime("%H:%M", ts, 'unixepoch', 'localtime') ORDER BY strftime("%H:%M", ts, 'unixepoch', 'localtime') ASC`, entity.Id, from.Unix(), diff --git a/core/metrics/db_test.go b/core/metrics/db_test.go index d61c2125b..5c66f1bd3 100644 --- a/core/metrics/db_test.go +++ b/core/metrics/db_test.go @@ -20,7 +20,7 @@ func TestSqliteTimestamp(t *testing.T) { entity := entity{Name: "foo"} require.NoError(t, db.Instance.FirstOrCreate(&entity).Error) - persist(entity, clock.Now(), 0, 0, nil) + persist(entity, clock.Now(), 0, 0, nil, false) db, err := db.Instance.DB() require.NoError(t, err) @@ -55,8 +55,8 @@ func TestQueryEnergyUTCFilter(t *testing.T) { loc := time.Now().Location() base := time.Date(2026, 4, 15, 16, 0, 0, 0, loc) - require.NoError(t, persist(e, base, 0, 1, nil)) - require.NoError(t, persist(e, base.Add(time.Hour), 0, 2, nil)) + require.NoError(t, persist(e, base, 0, 1, nil, false)) + require.NoError(t, persist(e, base.Add(time.Hour), 0, 2, nil, false)) // query with UTC times spanning both slots from := base.Add(-time.Hour).UTC() @@ -83,10 +83,10 @@ func TestQueryEnergyGrouped(t *testing.T) { loc := time.Now().Location() base := time.Date(2026, 4, 15, 16, 0, 0, 0, loc) - require.NoError(t, persist(e1, base, 1, 0, nil)) - require.NoError(t, persist(e2, base, 2, 0, nil)) - require.NoError(t, persist(e1, base.Add(time.Hour), 3, 0, nil)) - require.NoError(t, persist(e2, base.Add(time.Hour), 4, 0, nil)) + require.NoError(t, persist(e1, base, 1, 0, nil, false)) + require.NoError(t, persist(e2, base, 2, 0, nil, false)) + require.NoError(t, persist(e1, base.Add(time.Hour), 3, 0, nil, false)) + require.NoError(t, persist(e2, base.Add(time.Hour), 4, 0, nil, false)) from := base.Add(-time.Hour).UTC() to := base.Add(3 * time.Hour).UTC() @@ -125,9 +125,9 @@ func TestQueryEnergyMultipleSeries(t *testing.T) { // 2 hourly slots per entity for i := range 2 { ts := base.Add(time.Duration(i) * time.Hour) - require.NoError(t, persist(eGrid, ts, float64(1+i), 0, nil)) - require.NoError(t, persist(ePv1, ts, 0, float64(10+i), nil)) - require.NoError(t, persist(ePv2, ts, 0, float64(20+i), nil)) + require.NoError(t, persist(eGrid, ts, float64(1+i), 0, nil, false)) + require.NoError(t, persist(ePv1, ts, 0, float64(10+i), nil, false)) + require.NoError(t, persist(ePv2, ts, 0, float64(20+i), nil, false)) } from := base.Add(-time.Hour).UTC() @@ -187,9 +187,9 @@ func TestQueryEnergyFilter(t *testing.T) { base := time.Date(2026, 4, 15, 16, 0, 0, 0, loc) for i := range 2 { ts := base.Add(time.Duration(i) * time.Hour) - require.NoError(t, persist(eGrid, ts, float64(1+i), 0, nil)) - require.NoError(t, persist(ePv, ts, 0, float64(10+i), nil)) - require.NoError(t, persist(eBat, ts, float64(5+i), 0, nil)) + require.NoError(t, persist(eGrid, ts, float64(1+i), 0, nil, false)) + require.NoError(t, persist(ePv, ts, 0, float64(10+i), nil, false)) + require.NoError(t, persist(eBat, ts, float64(5+i), 0, nil, false)) } from := base.Add(-time.Hour).UTC() @@ -237,7 +237,7 @@ func TestUpdateProfile(t *testing.T) { // day 1: 0 ... 95 // day 2: 96 ... 181 for i := range 4 * 2 * 24 { - persist(entity, clock.Now(), float64(i), float64(i), nil) + persist(entity, clock.Now(), float64(i), float64(i), nil, false) clock.Add(15 * time.Minute) } diff --git a/core/metrics/stats_test.go b/core/metrics/stats_test.go index 97848797b..85653ff7e 100644 --- a/core/metrics/stats_test.go +++ b/core/metrics/stats_test.go @@ -25,13 +25,13 @@ func TestCollectorEnergyStats(t *testing.T) { require.NoError(t, err) e := col.entity - require.NoError(t, persist(e, slotStart.AddDate(0, 0, -8), 100, 0, nil)) // outside 7d - require.NoError(t, persist(e, slotStart.Add(-24*time.Hour-15*time.Minute), 5, 0, nil)) // 7d only - require.NoError(t, persist(e, slotStart.Add(-24*time.Hour), 7, 0, nil)) // 24h window start - require.NoError(t, persist(e, slotStart.Add(-12*time.Hour), 3, 0, nil)) // yesterday, within 24h - require.NoError(t, persist(e, midnight, 1, 0, nil)) // first slot today - require.NoError(t, persist(e, slotStart.Add(-15*time.Minute), 2, 0, nil)) // last completed slot - require.NoError(t, persist(e, slotStart, 4, 0, nil)) // current slot, excluded + require.NoError(t, persist(e, slotStart.AddDate(0, 0, -8), 100, 0, nil, false)) // outside 7d + require.NoError(t, persist(e, slotStart.Add(-24*time.Hour-15*time.Minute), 5, 0, nil, false)) // 7d only + require.NoError(t, persist(e, slotStart.Add(-24*time.Hour), 7, 0, nil, false)) // 24h window start + require.NoError(t, persist(e, slotStart.Add(-12*time.Hour), 3, 0, nil, false)) // yesterday, within 24h + require.NoError(t, persist(e, midnight, 1, 0, nil, false)) // first slot today + require.NoError(t, persist(e, slotStart.Add(-15*time.Minute), 2, 0, nil, false)) // last completed slot + require.NoError(t, persist(e, slotStart, 4, 0, nil, false)) // current slot, excluded stats, err := col.EnergyStats() require.NoError(t, err)