MQTT: publish one forecast message per key (BC) (#32716)
This commit is contained in:
parent
fff5502a36
commit
dc143a6dc7
2 changed files with 47 additions and 7 deletions
|
|
@ -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))),
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue