From dc143a6dc7dd75fc1e01524781ae9e746dc2c960 Mon Sep 17 00:00:00 2001 From: andig Date: Mon, 10 Aug 2026 21:54:37 +0200 Subject: [PATCH] MQTT: publish one forecast message per key (BC) (#32716) --- core/site_tariffs.go | 32 +++++++++++++++++++++++++------- server/mqtt_test.go | 22 ++++++++++++++++++++++ 2 files changed, 47 insertions(+), 7 deletions(-) diff --git a/core/site_tariffs.go b/core/site_tariffs.go index e9790a020..f26a7ffd5 100644 --- a/core/site_tariffs.go +++ b/core/site_tariffs.go @@ -1,6 +1,7 @@ package core import ( + "encoding/json" "math" "time" @@ -13,6 +14,17 @@ import ( "github.com/samber/lo" ) +// forecastSeries and solarDetails implement BytesMarshaler so MQTT publishes one +// json message per forecast key instead of decomposing every slot into its own +// topic (several thousand messages per update). +type forecastSeries [][]float64 + +var _ api.BytesMarshaler = (*forecastSeries)(nil) + +func (s forecastSeries) MarshalBytes() ([]byte, error) { + return json.Marshal(s) +} + type solarDetails struct { Scale float64 `json:"scale"` // scale factor yield/forecasted today, 1 if unscaled Today dailyDetails `json:"today,omitempty"` // tomorrow @@ -21,6 +33,12 @@ type solarDetails struct { Timeseries timeseries `json:"timeseries,omitempty"` // timeseries of forecasted energy } +var _ api.BytesMarshaler = (*solarDetails)(nil) + +func (d solarDetails) MarshalBytes() ([]byte, error) { + return json.Marshal(d) +} + type dailyDetails struct { Yield float64 `json:"energy"` Complete bool `json:"complete"` @@ -29,7 +47,7 @@ type dailyDetails struct { // forecastRates publishes rates as [start, end, value] with the timestamps in // unix seconds. The forecast is the largest payload evcc sends and RFC3339 // timestamps are two thirds of it. -func forecastRates(rr api.Rates) [][]float64 { +func forecastRates(rr api.Rates) forecastSeries { // keep nil for empty rates: shards are published without omitempty if len(rr) == 0 { return nil @@ -116,12 +134,12 @@ func (site *Site) publishTariffs(greenShareHome float64, greenShareLoadpoints fl } fc := struct { - Co2 [][]float64 `json:"co2,omitempty"` - FeedIn [][]float64 `json:"feedin,omitempty"` - Grid [][]float64 `json:"grid,omitempty"` - Planner [][]float64 `json:"planner,omitempty"` - Solar *solarDetails `json:"solar,omitempty"` - Temperature [][]float64 `json:"temperature,omitempty"` + Co2 forecastSeries `json:"co2,omitempty"` + FeedIn forecastSeries `json:"feedin,omitempty"` + Grid forecastSeries `json:"grid,omitempty"` + Planner forecastSeries `json:"planner,omitempty"` + Solar *solarDetails `json:"solar,omitempty"` + Temperature forecastSeries `json:"temperature,omitempty"` }{ Co2: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageCo2))), FeedIn: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageFeedIn))), diff --git a/server/mqtt_test.go b/server/mqtt_test.go index 87b05d1e0..115800433 100644 --- a/server/mqtt_test.go +++ b/server/mqtt_test.go @@ -1,6 +1,7 @@ package server import ( + "encoding/json" "math" "slices" "strconv" @@ -8,6 +9,7 @@ import ( "time" "github.com/evcc-io/evcc/core/types" + "github.com/evcc-io/evcc/util" "github.com/samber/lo" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/suite" @@ -97,6 +99,26 @@ func (suite *mqttSuite) TestSlice() { suite.Equal([]string{"2", "10", "20"}, suite.payloads, "payloads") } +type jsonSeries [][]float64 + +func (s jsonSeries) MarshalBytes() ([]byte, error) { + return json.Marshal(s) +} + +// a sharded struct publishes one message per key, BytesMarshaler fields as json +func (suite *mqttSuite) TestSharderBytesMarshaler() { + p := struct { + Foo jsonSeries `json:"foo,omitempty"` + Bar jsonSeries `json:"bar,omitempty"` + }{ + Foo: jsonSeries{{1, 2, 3}, {4, 5, 6}}, + } + + suite.publish("test", false, util.NewSharder("test", p)) + suite.Equal([]string{"test/foo", "test/bar"}, suite.topics, "topics") + suite.Equal([]string{"[[1,2,3],[4,5,6]]", ""}, suite.payloads, "payloads") +} + func (suite *mqttSuite) TestNilInterface() { var ptr *time.Time suite.publish("test", false, ptr)