Revert "MQTT: publish the forecast as a single JSON message (#32710)"

This reverts commit 93c182e486.
This commit is contained in:
andig 2026-08-10 21:27:08 +02:00
parent 93c182e486
commit fff5502a36
3 changed files with 18 additions and 48 deletions

View file

@ -1,7 +1,6 @@
package core package core
import ( import (
"encoding/json"
"math" "math"
"time" "time"
@ -14,24 +13,6 @@ import (
"github.com/samber/lo" "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 { type solarDetails struct {
Scale float64 `json:"scale"` // scale factor yield/forecasted today, 1 if unscaled Scale float64 `json:"scale"` // scale factor yield/forecasted today, 1 if unscaled
Today dailyDetails `json:"today,omitempty"` // tomorrow Today dailyDetails `json:"today,omitempty"` // tomorrow
@ -134,7 +115,14 @@ func (site *Site) publishTariffs(greenShareHome float64, greenShareLoadpoints fl
site.publish(keys.TariffCo2Loadpoints, v) site.publish(keys.TariffCo2Loadpoints, v)
} }
fc := forecastDetails{ 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: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageCo2))), Co2: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageCo2))),
FeedIn: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageFeedIn))), FeedIn: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsageFeedIn))),
Planner: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsagePlanner))), Planner: forecastRates(tariff.Rates(site.GetTariff(api.TariffUsagePlanner))),

View file

@ -84,16 +84,6 @@ func mqttTagAttribute(attr string, f reflect.StructField) bool {
} }
func (m *MQTT) publishComplex(topic string, retained bool, payload any) { 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 { if _, ok := payload.(fmt.Stringer); ok || payload == nil {
m.publishSingleValue(topic, retained, payload) m.publishSingleValue(topic, retained, payload)
return return
@ -108,6 +98,16 @@ func (m *MQTT) publishComplex(topic string, retained bool, payload any) {
return 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() { switch typ := reflect.TypeOf(payload); typ.Kind() {
case reflect.Slice: case reflect.Slice:
// publish count // publish count

View file

@ -1,7 +1,6 @@
package server package server
import ( import (
"encoding/json"
"math" "math"
"slices" "slices"
"strconv" "strconv"
@ -9,7 +8,6 @@ import (
"time" "time"
"github.com/evcc-io/evcc/core/types" "github.com/evcc-io/evcc/core/types"
"github.com/evcc-io/evcc/util"
"github.com/samber/lo" "github.com/samber/lo"
"github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert"
"github.com/stretchr/testify/suite" "github.com/stretchr/testify/suite"
@ -99,22 +97,6 @@ func (suite *mqttSuite) TestSlice() {
suite.Equal([]string{"2", "10", "20"}, suite.payloads, "payloads") 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() { func (suite *mqttSuite) TestNilInterface() {
var ptr *time.Time var ptr *time.Time
suite.publish("test", false, ptr) suite.publish("test", false, ptr)