diff --git a/core/metrics/db.go b/core/metrics/db.go index b6049651f..e5031a54d 100644 --- a/core/metrics/db.go +++ b/core/metrics/db.go @@ -2,6 +2,8 @@ package metrics import ( "errors" + "slices" + "strings" "time" "github.com/evcc-io/evcc/server/db" @@ -10,11 +12,11 @@ import ( ) type meter struct { - Meter int `json:"meter" gorm:"column:meter;uniqueIndex:meters_meter_ts"` - Timestamp time.Time `json:"ts" gorm:"column:ts;uniqueIndex:meters_meter_ts"` // start of 15min slot - Entity entity `json:"-" gorm:"foreignkey:Meter;references:Id"` - Import float64 `json:"import" gorm:"column:import"` - Export float64 `json:"export" gorm:"column:export"` + Meter int `json:"meter" gorm:"column:meter;uniqueIndex:meters_meter_ts"` + Timestamp int64 `json:"ts" gorm:"column:ts;uniqueIndex:meters_meter_ts"` // start of 15min slot + Entity entity `json:"-" gorm:"foreignkey:Meter;references:Id"` + Import float64 `json:"import" gorm:"column:import"` + Export float64 `json:"export" gorm:"column:export"` } type entity struct { @@ -23,118 +25,6 @@ type entity struct { Name string `gorm:"column:name;uniqueIndex:entities_group_name"` } -var ErrIncomplete = errors.New("meter profile incomplete") - -// Slot represents an aggregated energy time slot -type Slot struct { - Start time.Time `json:"start"` - End time.Time `json:"end"` - Import float64 `json:"import"` - Export float64 `json:"export"` -} - -// Series represents a named series of energy slots -type Series struct { - Name string `json:"name,omitempty"` - Group string `json:"group"` - Data []Slot `json:"data"` -} - -var aggregateFormats = map[string]string{ - "15m": "%Y-%m-%d %H:%M", - "hour": "%Y-%m-%d %H:00", - "day": "%Y-%m-%d", - "month": "%Y-%m", -} - -var aggregateGoFormats = map[string]string{ - "15m": "2006-01-02 15:04", - "hour": "2006-01-02 15:00", - "day": "2006-01-02", - "month": "2006-01", -} - -var aggregateDurations = map[string]func(time.Time) time.Time{ - "15m": func(t time.Time) time.Time { return t.Add(15 * time.Minute) }, - "hour": func(t time.Time) time.Time { return t.Add(time.Hour) }, - "day": func(t time.Time) time.Time { return t.AddDate(0, 0, 1) }, - "month": func(t time.Time) time.Time { return t.AddDate(0, 1, 0) }, -} - -// QueryImportEnergy returns aggregated energy data, per entity or per group. -func QueryImportEnergy(from, to time.Time, aggregate string, grouped bool) ([]Series, error) { - format, ok := aggregateFormats[aggregate] - if !ok { - return nil, errors.New("invalid aggregate value") - } - - addDuration := aggregateDurations[aggregate] - - // match timezone of stored timestamps for correct SQLite comparison - from = from.Local() - to = to.Local() - tz := from.Format("-07:00") - - groupCols := `e.name, e."group", bucket` - if grouped { - groupCols = `e."group", bucket` - } - - type row struct { - Name string - Group string - Bucket string - Import float64 - Export float64 - } - - tx := db.Instance.Table("meters m"). - Select(`e.name, e."group", strftime(?, m.ts, ?) AS bucket, - COALESCE(SUM(m."import"), 0) AS import, COALESCE(SUM(m.export), 0) AS export`, format, tz). - Joins("JOIN entities e ON m.meter = e.id"). - Group(groupCols). - Order(groupCols) - - if !from.IsZero() { - tx = tx.Where("m.ts >= ?", from) - } - if !to.IsZero() { - tx = tx.Where("m.ts < ?", to) - } - - var rows []row - if err := tx.Scan(&rows).Error; err != nil { - return nil, err - } - - var res []Series - for _, r := range rows { - start, err := time.ParseInLocation(aggregateGoFormats[aggregate], r.Bucket, time.Now().Location()) - if err != nil { - return nil, err - } - - name := r.Name - if grouped { - name = "" - } - - if n := len(res); n == 0 || res[n-1].Name != name || res[n-1].Group != r.Group { - res = append(res, Series{Name: name, Group: r.Group}) - } - - s := &res[len(res)-1] - s.Data = append(s.Data, Slot{ - Start: start, - End: addDuration(start), - Import: r.Import, - Export: r.Export, - }) - } - - return res, nil -} - func init() { db.Register(func(_ *gorm.DB) error { return SetupSchema() @@ -202,6 +92,31 @@ func SetupSchema() error { return err } + // meter: ts migration + if m.HasTable(new(meter)) { + types, err := m.ColumnTypes(new(meter)) + if err != nil { + return err + } + tsIdx := slices.IndexFunc(types, func(typ gorm.ColumnType) bool { + return typ.Name() == "ts" + }) + if tsIdx == -1 { + return errors.New("missing meters.ts") + } + + if tsTyp, _ := types[tsIdx].ColumnType(); !strings.EqualFold(tsTyp, "INTEGER") { + db, err := db.Instance.DB() + if err != nil { + return err + } + + if _, err := db.Exec(`UPDATE meters SET ts = unixepoch(ts)`); err != nil { + return err + } + } + } + return db.Instance.AutoMigrate(new(meter)) } @@ -209,60 +124,8 @@ func SetupSchema() error { func persist(entity entity, ts time.Time, imp, exp float64) error { return db.Instance.Create(&meter{ Meter: entity.Id, - Timestamp: ts.Truncate(tariff.SlotDuration), + Timestamp: ts.Truncate(tariff.SlotDuration).Unix(), Import: imp, Export: exp, }).Error } - -// importProfile returns a 15min average meter profile in Wh. The profile -// is sorted by timestamp starting at 00:00. It is guaranteed to contain 96 15min values. -func importProfile(entity entity, from time.Time) (*[96]float64, error) { - db, err := db.Instance.DB() - if err != nil { - return nil, err - } - - tz := from.Format("-07:00") - - rows, err := db.Query(`SELECT min(ts) AS ts, avg(import) AS import - FROM meters - WHERE meter = ? AND ts >= ? - GROUP BY strftime("%H:%M", ts, '`+tz+`') - ORDER BY strftime("%H:%M", ts, '`+tz+`') ASC`, entity.Id, from, - ) - if err != nil { - return nil, err - } - defer rows.Close() - - var prev time.Time - res := make([]float64, 0, 96) - - for rows.Next() { - var ts SqlTime - var val float64 - - if err := rows.Scan(&ts, &val); err != nil { - return nil, err - } - - // interpolate single missing value, maybe due to regular restarts? - if time.Time(ts).Sub(prev) == 2*tariff.SlotDuration { - res = append(res, (val+res[len(res)-1])/2) - } - prev = time.Time(ts) - - res = append(res, val) - } - - if err := rows.Err(); err != nil { - return nil, err - } - - if len(res) != 96 { - return nil, ErrIncomplete - } - - return (*[96]float64)(res), nil -} diff --git a/core/metrics/db_history.go b/core/metrics/db_history.go new file mode 100644 index 000000000..883338f99 --- /dev/null +++ b/core/metrics/db_history.go @@ -0,0 +1,113 @@ +package metrics + +import ( + "errors" + "time" + + "github.com/evcc-io/evcc/server/db" +) + +// Slot represents an aggregated energy time slot +type Slot struct { + Start time.Time `json:"start"` + End time.Time `json:"end"` + Import float64 `json:"import"` + Export float64 `json:"export"` +} + +// Series represents a named series of energy slots +type Series struct { + Name string `json:"name,omitempty"` + Group string `json:"group"` + Data []Slot `json:"data"` +} + +var aggregateFormats = map[string]string{ + "15m": "%Y-%m-%d %H:%M", + "hour": "%Y-%m-%d %H:00", + "day": "%Y-%m-%d", + "month": "%Y-%m", +} + +var aggregateGoFormats = map[string]string{ + "15m": "2006-01-02 15:04", + "hour": "2006-01-02 15:00", + "day": "2006-01-02", + "month": "2006-01", +} + +var aggregateDurations = map[string]func(time.Time) time.Time{ + "15m": func(t time.Time) time.Time { return t.Add(15 * time.Minute) }, + "hour": func(t time.Time) time.Time { return t.Add(time.Hour) }, + "day": func(t time.Time) time.Time { return t.AddDate(0, 0, 1) }, + "month": func(t time.Time) time.Time { return t.AddDate(0, 1, 0) }, +} + +// QueryImportEnergy returns aggregated energy data, per entity or per group. +func QueryImportEnergy(from, to time.Time, aggregate string, grouped bool) ([]Series, error) { + format, ok := aggregateFormats[aggregate] + if !ok { + return nil, errors.New("invalid aggregate value") + } + + addDuration := aggregateDurations[aggregate] + + groupCols := `e.name, e."group", bucket` + if grouped { + groupCols = `e."group", bucket` + } + + type row struct { + Name string + Group string + Bucket string + Import float64 + Export float64 + } + + tx := db.Instance.Table("meters m"). + Select(`e.name, e."group", strftime(?, m.ts, 'unixepoch', 'localtime') AS bucket, + COALESCE(SUM(m."import"), 0) AS import, COALESCE(SUM(m.export), 0) AS export`, format). + Joins("JOIN entities e ON m.meter = e.id"). + Group(groupCols). + Order(groupCols) + + if !from.IsZero() { + tx = tx.Where("m.ts >= ?", from.Unix()) + } + if !to.IsZero() { + tx = tx.Where("m.ts < ?", to.Unix()) + } + + var rows []row + if err := tx.Scan(&rows).Error; err != nil { + return nil, err + } + + var res []Series + for _, r := range rows { + start, err := time.ParseInLocation(aggregateGoFormats[aggregate], r.Bucket, from.Location()) + if err != nil { + return nil, err + } + + name := r.Name + if grouped { + name = "" + } + + if n := len(res); n == 0 || res[n-1].Name != name || res[n-1].Group != r.Group { + res = append(res, Series{Name: name, Group: r.Group}) + } + + s := &res[len(res)-1] + s.Data = append(s.Data, Slot{ + Start: start, + End: addDuration(start), + Import: r.Import, + Export: r.Export, + }) + } + + return res, nil +} diff --git a/core/metrics/db_profile.go b/core/metrics/db_profile.go new file mode 100644 index 000000000..217663f85 --- /dev/null +++ b/core/metrics/db_profile.go @@ -0,0 +1,62 @@ +package metrics + +import ( + "errors" + "time" + + "github.com/evcc-io/evcc/server/db" + "github.com/evcc-io/evcc/tariff" +) + +var ErrIncomplete = errors.New("meter profile incomplete") + +// importProfile returns a 15min average meter profile in Wh. The profile +// is sorted by timestamp starting at 00:00. It is guaranteed to contain 96 15min values. +func importProfile(entity entity, from time.Time) (*[96]float64, error) { + db, err := db.Instance.DB() + if err != nil { + return nil, err + } + + rows, err := db.Query(`SELECT min(ts) AS ts, avg(import) AS import + FROM meters + WHERE meter = ? AND ts >= ? + GROUP BY strftime("%H:%M", ts, 'unixepoch', 'localtime') + ORDER BY strftime("%H:%M", ts, 'unixepoch', 'localtime') ASC`, + entity.Id, from.Unix(), + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var prev time.Time + res := make([]float64, 0, 96) + + for rows.Next() { + var ts SqlTime + var val float64 + + if err := rows.Scan(&ts, &val); err != nil { + return nil, err + } + + // interpolate single missing value, maybe due to regular restarts? + if time.Time(ts).Sub(prev) == 2*tariff.SlotDuration { + res = append(res, (val+res[len(res)-1])/2) + } + prev = time.Time(ts) + + res = append(res, val) + } + + if err := rows.Err(); err != nil { + return nil, err + } + + if len(res) != 96 { + return nil, ErrIncomplete + } + + return (*[96]float64)(res), nil +} diff --git a/core/metrics/db_test.go b/core/metrics/db_test.go index dcc409a40..9b752711f 100644 --- a/core/metrics/db_test.go +++ b/core/metrics/db_test.go @@ -1,6 +1,7 @@ package metrics import ( + "slices" "testing" "time" @@ -24,26 +25,22 @@ func TestSqliteTimestamp(t *testing.T) { db, err := db.Instance.DB() require.NoError(t, err) - var ( - ts SqlTime - val float64 - ) + var ts SqlTime for _, sql := range []string{ - `SELECT ts, import FROM meters`, - `SELECT min(ts), import FROM meters`, - `SELECT unixepoch(ts), import FROM meters`, - `SELECT unixepoch(min(ts)), import FROM meters`, - `SELECT min(ts) AS ts, avg(import) AS import - FROM meters + `SELECT ts FROM meters`, + `SELECT min(ts) FROM meters`, + // `SELECT unixepoch(ts) FROM meters`, + // `SELECT unixepoch(min(ts)) FROM meters`, + `SELECT min(ts) AS ts FROM meters GROUP BY strftime("%H:%M", ts) ORDER BY ts`, } { - require.NoError(t, db.QueryRow(sql).Scan(&ts, &val)) - require.True(t, clock.Now().Equal(time.Time(ts)), "expected %v, got %v", clock.Now().Local(), time.Time(ts).Local()) + require.NoError(t, db.QueryRow(sql).Scan(&ts)) + require.True(t, clock.Now().Equal(time.Time(ts)), "expected %v, got %v (%s)", clock.Now().Local(), time.Time(ts).Local(), sql) } - require.NoError(t, db.QueryRow(`SELECT ts, import FROM meters WHERE ts >= ?`, clock.Now()).Scan(&ts, &val)) + require.NoError(t, db.QueryRow(`SELECT ts FROM meters WHERE ts >= ?`, clock.Now().Unix()).Scan(&ts)) require.True(t, clock.Now().Equal(time.Time(ts)), "expected %v, got %v", clock.Now().Local(), time.Time(ts).Local()) } @@ -138,6 +135,11 @@ func TestUpdateProfile(t *testing.T) { clock.Add(15 * time.Minute) } + // validate records written + var count int64 + require.NoError(t, db.Instance.Model(new(meter)).Count(&count).Error) + require.Equal(t, int64(24*2*4), count) + { from := clock.Now().Local().AddDate(0, 0, -2).Add(12 * time.Hour) // 12:00 of day 0 @@ -170,3 +172,39 @@ func TestUpdateProfile(t *testing.T) { require.Equal(t, expected, *prof, "full profile: expected %v, got %v", expected, *prof) } } + +func TestTimeMigration(t *testing.T) { + require.NoError(t, db.NewInstance("sqlite", ":memory:")) + mig := db.Instance.Migrator() + + require.NoError(t, db.Instance.AutoMigrate(new(entity))) + + type v1 struct { + Meter int `json:"meter" gorm:"column:meter;uniqueIndex:idx_meter_ts"` + Timestamp time.Time `json:"ts" gorm:"column:ts;uniqueIndex:idx_meter_ts"` + Entity entity `json:"-" gorm:"foreignkey:Meter;references:Id"` + } + + require.NoError(t, db.Instance.AutoMigrate(new(v1))) + { + tables, err := mig.GetTables() + require.NoError(t, err) + require.True(t, slices.Contains(tables, "v1")) + } + + require.NoError(t, mig.RenameTable("v1", "v2")) + + type v2 struct { + Meter int `json:"meter" gorm:"column:meter;uniqueIndex:idx_meter_ts"` + Timestamp int64 `json:"ts" gorm:"column:ts;uniqueIndex:idx_meter_ts"` + Entity entity `json:"-" gorm:"foreignkey:Meter;references:Id"` + } + + require.NoError(t, db.Instance.AutoMigrate(new(v2))) + { + tables, err := mig.GetTables() + require.NoError(t, err) + require.False(t, slices.Contains(tables, "v1")) + require.True(t, slices.Contains(tables, "v2")) + } +} diff --git a/go.mod b/go.mod index 6bee2bf02..62c5605f6 100644 --- a/go.mod +++ b/go.mod @@ -271,3 +271,5 @@ replace github.com/grid-x/modbus => github.com/evcc-io/modbus v0.0.0-20250501165 replace github.com/lorenzodonini/ocpp-go => github.com/evcc-io/ocpp-go v0.0.0-20251212212612-b7f92ee0443b replace go.yaml.in/yaml/v4 => go.yaml.in/yaml/v4 v4.0.0-rc.3 + +replace github.com/glebarez/sqlite => github.com/evcc-io/sqlite v0.0.0-20260421123006-d66e0643f9bb diff --git a/go.sum b/go.sum index 5203e7ee6..57f06df95 100644 --- a/go.sum +++ b/go.sum @@ -206,6 +206,8 @@ github.com/evcc-io/optimizer v0.0.0-20260411145738-bf13a64d411c h1:KUJqCEBzEXepn github.com/evcc-io/optimizer v0.0.0-20260411145738-bf13a64d411c/go.mod h1:cHMeq6GfoXQ8LxqY4LH1wn/hrgYEn+ka5RPI1CF37SY= github.com/evcc-io/rct v0.2.0 h1:qjXKI7NKwW5YpEKEb27MwLMBaR4ne2Kd2cPV2PYTKn4= github.com/evcc-io/rct v0.2.0/go.mod h1:n6MTBU36QOadGlxxiADu86VaT5l0eYdbmFBHF14AD/s= +github.com/evcc-io/sqlite v0.0.0-20260421123006-d66e0643f9bb h1:k0oq1Bcuidy30jREK6drmSVk0i5M9PENZULzBfk9Amw= +github.com/evcc-io/sqlite v0.0.0-20260421123006-d66e0643f9bb/go.mod h1:jTnU6XZ6pKY74rm4VrCSxgXmH0x3p3PcA00uAzr0OtQ= github.com/evcc-io/tesla-proxy-client v0.0.0-20260324063928-151fe10796ae h1:0ihqBMBjkekSlwD3ZYQIiYite6ZAln9ze60evOE6FkQ= github.com/evcc-io/tesla-proxy-client v0.0.0-20260324063928-151fe10796ae/go.mod h1:L8AR3E7/MbbvCgnrtzsT6YsgQJtN5gk/LumR1gDBeJM= github.com/fatih/camelcase v1.0.0/go.mod h1:yN2Sb0lFhZJUdVvtELVWefmrXpuZESvPmqwoZc+/fpc= @@ -230,8 +232,6 @@ github.com/getkin/kin-openapi v0.135.0/go.mod h1:6dd5FJl6RdX4usBtFBaQhk9q62Yb2J0 github.com/ghodss/yaml v1.0.0/go.mod h1:4dBDuWmgqj2HViK6kFavaiC9ZROes6MMH2rRYeMEF04= github.com/glebarez/go-sqlite v1.22.0 h1:uAcMJhaA6r3LHMTFgP0SifzgXg46yJkgxqyuyec+ruQ= github.com/glebarez/go-sqlite v1.22.0/go.mod h1:PlBIdHe0+aUEFn+r2/uthrWq4FxbzugL0L8Li6yQJbc= -github.com/glebarez/sqlite v1.11.0 h1:wSG0irqzP6VurnMEpFGer5Li19RpIRi2qvQz++w0GMw= -github.com/glebarez/sqlite v1.11.0/go.mod h1:h8/o8j5wiAsqSPoWELDUdJXhjAhsVliSn7bWZjOhrgQ= github.com/go-http-utils/etag v0.0.0-20161124023236-513ea8f21eb1 h1:zga7zaRE8HCbWjcXMDlfvmQtH0/kMVLo7cQ48dy6kWg= github.com/go-http-utils/etag v0.0.0-20161124023236-513ea8f21eb1/go.mod h1:PumS+5d59wmAGsZo6IfRpVNaJUq+6xjC4Utt/k8GO6Q= github.com/go-http-utils/fresh v0.0.0-20161124030543-7231e26a4b27 h1:O6yi4xa9b2DMosGsXzlMe2E9qXgXCVkRLCoRX+5amxI= diff --git a/tests/energy-history.sql b/tests/energy-history.sql index b0aeb0f26..9c0de65d9 100644 --- a/tests/energy-history.sql +++ b/tests/energy-history.sql @@ -11,7 +11,7 @@ CREATE UNIQUE INDEX `entities_group_name` ON `entities`(`group`, `name`); CREATE TABLE `meters` ( `meter` integer, - `ts` datetime, + `ts` integer, `import` real, `export` real ); @@ -26,17 +26,17 @@ INSERT INTO `entities` VALUES (2, 'grid', 'grid'); -- 2026-03-25: 00:00, 00:15 (2 slots) -- home (id=1): import=0.1, export=0 per slot -INSERT INTO `meters` VALUES (1, '2026-03-24 22:00:00+01:00', 0.1, 0); -INSERT INTO `meters` VALUES (1, '2026-03-24 22:15:00+01:00', 0.1, 0); -INSERT INTO `meters` VALUES (1, '2026-03-24 22:30:00+01:00', 0.1, 0); -INSERT INTO `meters` VALUES (1, '2026-03-24 22:45:00+01:00', 0.1, 0); -INSERT INTO `meters` VALUES (1, '2026-03-25 00:00:00+01:00', 0.1, 0); -INSERT INTO `meters` VALUES (1, '2026-03-25 00:15:00+01:00', 0.1, 0); +INSERT INTO `meters` VALUES (1, 1774386000, 0.1, 0); +INSERT INTO `meters` VALUES (1, 1774386900, 0.1, 0); +INSERT INTO `meters` VALUES (1, 1774387800, 0.1, 0); +INSERT INTO `meters` VALUES (1, 1774388700, 0.1, 0); +INSERT INTO `meters` VALUES (1, 1774393200, 0.1, 0); +INSERT INTO `meters` VALUES (1, 1774394100, 0.1, 0); -- grid (id=2): import=0.5, export=0.1 per slot -INSERT INTO `meters` VALUES (2, '2026-03-24 22:00:00+01:00', 0.5, 0.1); -INSERT INTO `meters` VALUES (2, '2026-03-24 22:15:00+01:00', 0.5, 0.1); -INSERT INTO `meters` VALUES (2, '2026-03-24 22:30:00+01:00', 0.5, 0.1); -INSERT INTO `meters` VALUES (2, '2026-03-24 22:45:00+01:00', 0.5, 0.1); -INSERT INTO `meters` VALUES (2, '2026-03-25 00:00:00+01:00', 0.5, 0.1); -INSERT INTO `meters` VALUES (2, '2026-03-25 00:15:00+01:00', 0.5, 0.1); +INSERT INTO `meters` VALUES (2, 1774386000, 0.5, 0.1); +INSERT INTO `meters` VALUES (2, 1774386900, 0.5, 0.1); +INSERT INTO `meters` VALUES (2, 1774387800, 0.5, 0.1); +INSERT INTO `meters` VALUES (2, 1774388700, 0.5, 0.1); +INSERT INTO `meters` VALUES (2, 1774393200, 0.5, 0.1); +INSERT INTO `meters` VALUES (2, 1774394100, 0.5, 0.1);