From 93c182e486f32180d9befe06841f3c5c4a7d46c4 Mon Sep 17 00:00:00 2001 From: andig Date: Mon, 10 Aug 2026 20:42:46 +0200 Subject: [PATCH] MQTT: publish the forecast as a single JSON message (#32710) --- core/site_tariffs.go | 28 ++++++++++++++++++++-------- server/mqtt.go | 20 ++++++++++---------- server/mqtt_test.go | 18 ++++++++++++++++++ 3 files changed, 48 insertions(+), 18 deletions(-) diff --git a/core/site_tariffs.go b/core/site_tariffs.go index e9790a020..ebfa767c0 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,24 @@ import ( "github.com/samber/lo" ) +// forecastDetails is the published forecast payload. It implements BytesMarshaler +// so MQTT sends a single JSON message instead of decomposing every slot into an +// individual topic (several thousand messages per update). +type forecastDetails 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"` +} + +var _ api.BytesMarshaler = (*forecastDetails)(nil) + +func (fc forecastDetails) MarshalBytes() ([]byte, error) { + return json.Marshal(fc) +} + type solarDetails struct { Scale float64 `json:"scale"` // scale factor yield/forecasted today, 1 if unscaled Today dailyDetails `json:"today,omitempty"` // tomorrow @@ -115,14 +134,7 @@ func (site *Site) publishTariffs(greenShareHome float64, greenShareLoadpoints fl site.publish(keys.TariffCo2Loadpoints, v) } - 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"` - }{ + fc := forecastDetails{ Co2: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageCo2))), FeedIn: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageFeedIn))), Planner: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsagePlanner))), diff --git a/server/mqtt.go b/server/mqtt.go index cfc8400f6..a4412d495 100644 --- a/server/mqtt.go +++ b/server/mqtt.go @@ -84,6 +84,16 @@ func mqttTagAttribute(attr string, f reflect.StructField) bool { } func (m *MQTT) publishComplex(topic string, retained bool, payload any) { + // unwrap first so the wrapped value is still checked for the marshalers below + if mm, ok := payload.(api.StructMarshaler); ok { + d, err := mm.MarshalStruct() + if err != nil { + m.log.ERROR.Printf("marshal struct: %v", err) + return + } + payload = d + } + if _, ok := payload.(fmt.Stringer); ok || payload == nil { m.publishSingleValue(topic, retained, payload) return @@ -98,16 +108,6 @@ func (m *MQTT) publishComplex(topic string, retained bool, payload any) { return } - if mm, ok := payload.(api.StructMarshaler); ok { - if d, err := mm.MarshalStruct(); err != nil { - m.log.ERROR.Printf("marshal struct: %v", err) - return - } else { - payload = d - // fallthrough - } - } - switch typ := reflect.TypeOf(payload); typ.Kind() { case reflect.Slice: // publish count diff --git a/server/mqtt_test.go b/server/mqtt_test.go index 87b05d1e0..b2a2655cf 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,22 @@ func (suite *mqttSuite) TestSlice() { suite.Equal([]string{"2", "10", "20"}, suite.payloads, "payloads") } +type jsonPayload struct { + Foo [][]float64 `json:"foo,omitempty"` +} + +func (p jsonPayload) MarshalBytes() ([]byte, error) { + return json.Marshal(p) +} + +// a BytesMarshaler wrapped in a Sharder must still publish as a single message +func (suite *mqttSuite) TestShardedBytesMarshaler() { + p := jsonPayload{Foo: [][]float64{{1, 2, 3}, {4, 5, 6}}} + suite.publish("test", false, util.NewSharder("test", p)) + suite.Equal([]string{"test"}, suite.topics, "topics") + suite.Equal([]string{`{"foo":[[1,2,3],[4,5,6]]}`}, suite.payloads, "payloads") +} + func (suite *mqttSuite) TestNilInterface() { var ptr *time.Time suite.publish("test", false, ptr)