Add admin api for deleting invalid metrics (#32666)
This commit is contained in:
parent
5d1f1193bb
commit
5ab5cfc2ef
7 changed files with 589 additions and 12 deletions
104
core/metrics/db_delete_test.go
Normal file
104
core/metrics/db_delete_test.go
Normal file
|
|
@ -0,0 +1,104 @@
|
|||
package metrics
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/evcc-io/evcc/server/db"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestDeleteEnergy(t *testing.T) {
|
||||
require.NoError(t, db.NewInstance("sqlite", ":memory:"))
|
||||
require.NoError(t, SetupSchema())
|
||||
|
||||
ePv := entity{Id: 2, Name: "pv1", Group: PV}
|
||||
require.NoError(t, db.Instance.Create(&ePv).Error)
|
||||
eFc := entity{Id: 3, Name: Forecast, Group: Forecast}
|
||||
require.NoError(t, db.Instance.Create(&eFc).Error)
|
||||
|
||||
loc := time.Now().Location()
|
||||
base := time.Date(2026, 4, 15, 16, 0, 0, 0, loc)
|
||||
for i := range 4 {
|
||||
ts := base.Add(time.Duration(i) * 15 * time.Minute)
|
||||
require.NoError(t, persist(ePv, ts, 1, 0, nil, false))
|
||||
require.NoError(t, persist(eFc, ts, 2, 0, nil, false))
|
||||
}
|
||||
|
||||
count := func() int64 {
|
||||
var n int64
|
||||
require.NoError(t, db.Instance.Model(new(meter)).Count(&n).Error)
|
||||
return n
|
||||
}
|
||||
require.Equal(t, int64(8), count())
|
||||
|
||||
// both bounds are required
|
||||
_, err := DeleteEnergy(time.Time{}, base, EnergyFilter{})
|
||||
require.Error(t, err)
|
||||
|
||||
// filter narrows to the matching entity, range is half-open
|
||||
rows, err := DeleteEnergy(base.UTC(), base.Add(30*time.Minute).UTC(), EnergyFilter{Group: Forecast})
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(2), rows)
|
||||
require.Equal(t, int64(6), count())
|
||||
|
||||
// pv untouched
|
||||
res, err := QueryEnergy(base.Add(-time.Hour).UTC(), base.Add(time.Hour).UTC(), "hour", false, EnergyFilter{Group: PV})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, res, 1)
|
||||
require.Equal(t, 4.0, res[0].Data[0].Energy)
|
||||
|
||||
// empty filter deletes across entities
|
||||
rows, err = DeleteEnergy(base.UTC(), base.Add(time.Hour).UTC(), EnergyFilter{})
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(6), rows)
|
||||
require.Zero(t, count())
|
||||
}
|
||||
|
||||
func TestDeleteTariffs(t *testing.T) {
|
||||
require.NoError(t, db.NewInstance("sqlite", ":memory:"))
|
||||
require.NoError(t, db.Instance.AutoMigrate(new(tariffValue)))
|
||||
|
||||
base := time.Date(2026, 4, 15, 16, 0, 0, 0, time.UTC)
|
||||
grid, co2 := 0.3, 250.0
|
||||
for i := range 4 {
|
||||
require.NoError(t, PersistTariffs(base.Add(time.Duration(i)*15*time.Minute), &grid, nil, &co2, nil))
|
||||
}
|
||||
|
||||
count := func() int64 {
|
||||
var n int64
|
||||
require.NoError(t, db.Instance.Model(new(tariffValue)).Count(&n).Error)
|
||||
return n
|
||||
}
|
||||
|
||||
_, err := DeleteTariffs(base, time.Time{}, "")
|
||||
require.Error(t, err)
|
||||
|
||||
// unknown usage never reaches the column interpolation
|
||||
_, err = DeleteTariffs(base, base.Add(time.Hour), "unknown")
|
||||
require.ErrorIs(t, err, ErrInvalidUsage)
|
||||
require.Equal(t, int64(4), count())
|
||||
|
||||
// clearing one usage keeps the row as long as another value remains
|
||||
rows, err := DeleteTariffs(base, base.Add(30*time.Minute), "grid")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(2), rows)
|
||||
require.Equal(t, int64(4), count())
|
||||
|
||||
var res tariffValue
|
||||
require.NoError(t, db.Instance.Where("ts = ?", base.Unix()).First(&res).Error)
|
||||
require.Nil(t, res.Grid)
|
||||
require.InDelta(t, 250, *res.Co2, 0.001)
|
||||
|
||||
// clearing the last usage drops the now empty rows
|
||||
rows, err = DeleteTariffs(base, base.Add(30*time.Minute), "co2")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(2), rows)
|
||||
require.Equal(t, int64(2), count())
|
||||
|
||||
// no usage drops the whole row
|
||||
rows, err = DeleteTariffs(base, base.Add(time.Hour), "")
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(2), rows)
|
||||
require.Zero(t, count())
|
||||
}
|
||||
|
|
@ -10,6 +10,7 @@ import (
|
|||
|
||||
"github.com/evcc-io/evcc/server/db"
|
||||
"github.com/evcc-io/evcc/util/export"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Slot represents an aggregated energy time slot
|
||||
|
|
@ -62,6 +63,47 @@ type EnergyFilter struct {
|
|||
Title string
|
||||
}
|
||||
|
||||
// entityQuery returns a subquery selecting the ids of the matching entities,
|
||||
// nil for an empty filter.
|
||||
func entityQuery(f EnergyFilter) *gorm.DB {
|
||||
if f.Group == "" && f.Name == "" && f.Title == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
tx := db.Instance.Model(new(entity)).Select("id")
|
||||
if f.Group != "" {
|
||||
tx = tx.Where(`"group" = ?`, f.Group)
|
||||
}
|
||||
if f.Name != "" {
|
||||
tx = tx.Where("name = ?", f.Name)
|
||||
}
|
||||
if f.Title != "" {
|
||||
tx = tx.Where("title = ?", f.Title)
|
||||
}
|
||||
|
||||
return tx
|
||||
}
|
||||
|
||||
// DeleteEnergy removes the slots in [from,to), narrowed to the matching
|
||||
// entities. Both bounds are required, a full wipe is /api/db/reset.
|
||||
func DeleteEnergy(from, to time.Time, filter ...EnergyFilter) (int64, error) {
|
||||
if from.IsZero() || to.IsZero() {
|
||||
return 0, errors.New("missing from/to")
|
||||
}
|
||||
|
||||
tx := db.Instance.Where("ts >= ? AND ts < ?", from.Unix(), to.Unix())
|
||||
|
||||
if len(filter) > 0 {
|
||||
if sub := entityQuery(filter[0]); sub != nil {
|
||||
tx = tx.Where("meter IN (?)", sub)
|
||||
}
|
||||
}
|
||||
|
||||
res := tx.Delete(new(meter))
|
||||
|
||||
return res.RowsAffected, res.Error
|
||||
}
|
||||
|
||||
// QueryEnergy returns aggregated energy data, per title or per group.
|
||||
func QueryEnergy(from, to time.Time, aggregate string, grouped bool, filter ...EnergyFilter) ([]Series, error) {
|
||||
addDuration := aggregateDurations[aggregate]
|
||||
|
|
@ -114,15 +156,8 @@ func QueryEnergy(from, to time.Time, aggregate string, grouped bool, filter ...E
|
|||
}
|
||||
|
||||
if len(filter) > 0 {
|
||||
f := filter[0]
|
||||
if f.Group != "" {
|
||||
tx = tx.Where(`e."group" = ?`, f.Group)
|
||||
}
|
||||
if f.Name != "" {
|
||||
tx = tx.Where("e.name = ?", f.Name)
|
||||
}
|
||||
if f.Title != "" {
|
||||
tx = tx.Where("e.title = ?", f.Title)
|
||||
if sub := entityQuery(filter[0]); sub != nil {
|
||||
tx = tx.Where("m.meter IN (?)", sub)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,6 +1,10 @@
|
|||
package metrics
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"slices"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/evcc-io/evcc/server/db"
|
||||
|
|
@ -26,6 +30,50 @@ func init() {
|
|||
})
|
||||
}
|
||||
|
||||
// ErrInvalidUsage is returned for an unknown tariff usage
|
||||
var ErrInvalidUsage = errors.New("invalid usage")
|
||||
|
||||
// tariffUsages are the deletable usages, named after their table column
|
||||
var tariffUsages = []string{"grid", "feedin", "co2", "temperature"}
|
||||
|
||||
// DeleteTariffs removes the persisted values in [from,to). An empty usage drops
|
||||
// the entire row, otherwise only that usage is cleared. Both bounds are
|
||||
// required, a full wipe is /api/db/reset. The count is the number of affected
|
||||
// rows; rows dropped by the cleanup are a subset of the cleared ones.
|
||||
func DeleteTariffs(from, to time.Time, usage string) (int64, error) {
|
||||
if from.IsZero() || to.IsZero() {
|
||||
return 0, errors.New("missing from/to")
|
||||
}
|
||||
|
||||
inRange := func() *gorm.DB {
|
||||
return db.Instance.Where("ts >= ? AND ts < ?", from.Unix(), to.Unix())
|
||||
}
|
||||
|
||||
if usage == "" {
|
||||
res := inRange().Delete(new(tariffValue))
|
||||
return res.RowsAffected, res.Error
|
||||
}
|
||||
|
||||
// guards the column interpolated below
|
||||
if !slices.Contains(tariffUsages, usage) {
|
||||
return 0, fmt.Errorf("%w: %s (valid: %s)", ErrInvalidUsage, usage, strings.Join(tariffUsages, ", "))
|
||||
}
|
||||
|
||||
res := inRange().Model(new(tariffValue)).
|
||||
Where(usage+" IS NOT NULL").
|
||||
Update(usage, gorm.Expr("NULL"))
|
||||
if res.Error != nil {
|
||||
return 0, res.Error
|
||||
}
|
||||
|
||||
// drop the rows that no longer hold any value
|
||||
err := inRange().
|
||||
Where("grid IS NULL AND feedin IS NULL AND co2 IS NULL AND temperature IS NULL").
|
||||
Delete(new(tariffValue)).Error
|
||||
|
||||
return res.RowsAffected, err
|
||||
}
|
||||
|
||||
// PersistTariffs stores the tariff values at the given 15min boundary, nil values omitted
|
||||
func PersistTariffs(ts time.Time, grid, feedin, co2, temperature *float64) error {
|
||||
if grid == nil && feedin == nil && co2 == nil && temperature == nil {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue