Tariffs: reduce published data volume by x10 (#24375)
Some checks failed
Some checks failed
This commit is contained in:
parent
62e58b0c7c
commit
4718e5f90f
9 changed files with 141 additions and 51 deletions
|
|
@ -31,6 +31,7 @@ const initialState: State = {
|
|||
offline: false,
|
||||
loadpoints: [],
|
||||
vehicles: {},
|
||||
forecast: {},
|
||||
};
|
||||
|
||||
const state = reactive(initialState);
|
||||
|
|
|
|||
|
|
@ -52,7 +52,7 @@ export interface State {
|
|||
offline: boolean;
|
||||
startup?: boolean;
|
||||
loadpoints: Loadpoint[];
|
||||
forecast?: Forecast;
|
||||
forecast: Forecast;
|
||||
currency?: CURRENCY;
|
||||
fatal?: FatalError[];
|
||||
authProviders?: AuthProviders;
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ import (
|
|||
"github.com/evcc-io/evcc/core/keys"
|
||||
"github.com/evcc-io/evcc/server/db/settings"
|
||||
"github.com/evcc-io/evcc/tariff"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/jinzhu/now"
|
||||
"github.com/samber/lo"
|
||||
)
|
||||
|
|
@ -117,7 +118,7 @@ func (site *Site) publishTariffs(greenShareHome float64, greenShareLoadpoints fl
|
|||
fc.Solar = lo.ToPtr(site.solarDetails(solar))
|
||||
}
|
||||
|
||||
site.publish(keys.Forecast, fc)
|
||||
site.publish(keys.Forecast, util.NewSharder(keys.Forecast, fc))
|
||||
}
|
||||
|
||||
func (site *Site) solarDetails(solar api.Rates) solarDetails {
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ package server
|
|||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
|
|
@ -110,33 +111,59 @@ func (h *SocketHub) deleteSubscriber(s *socketSubscriber) {
|
|||
}
|
||||
|
||||
func (h *SocketHub) welcome(subscriber *socketSubscriber, params []util.Param) {
|
||||
var msg strings.Builder
|
||||
msg.WriteString("{")
|
||||
msg := make(map[string]json.RawMessage, len(params))
|
||||
|
||||
for _, p := range params {
|
||||
if msg.Len() > 1 {
|
||||
msg.WriteString(",")
|
||||
k := p.Key
|
||||
if p.Loadpoint != nil {
|
||||
k = "loadpoints." + p.UniqueID()
|
||||
}
|
||||
msg.WriteString(kv(p))
|
||||
|
||||
msg[k] = json.RawMessage(socketEncode(p.Val))
|
||||
}
|
||||
msg.WriteString("}")
|
||||
|
||||
b, _ := json.Marshal(msg)
|
||||
|
||||
// should not block
|
||||
subscriber.send <- []byte(msg.String())
|
||||
subscriber.send <- b
|
||||
}
|
||||
|
||||
func (h *SocketHub) broadcast(p util.Param) {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
|
||||
if len(h.subscribers) > 0 {
|
||||
msg := "{" + kv(p) + "}"
|
||||
if len(h.subscribers) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
for s := range h.subscribers {
|
||||
select {
|
||||
case s.send <- []byte(msg):
|
||||
default:
|
||||
s.closeSlow()
|
||||
}
|
||||
msg := make(map[string]json.RawMessage)
|
||||
|
||||
k := p.Key
|
||||
if p.Loadpoint != nil {
|
||||
k = "loadpoints." + p.UniqueID()
|
||||
}
|
||||
|
||||
// Sharder splits data into chunks
|
||||
if sp, ok := (p.Val).(util.Sharder); ok {
|
||||
shards := sp.Shards()
|
||||
if len(shards) == 0 {
|
||||
return // nothing changed, skip broadcast
|
||||
}
|
||||
|
||||
for _, shard := range shards {
|
||||
msg[k+"."+shard.Key] = json.RawMessage(socketEncode(shard.Value))
|
||||
}
|
||||
} else {
|
||||
msg[k] = json.RawMessage(socketEncode(p.Val))
|
||||
}
|
||||
|
||||
b, _ := json.Marshal(msg)
|
||||
|
||||
for s := range h.subscribers {
|
||||
select {
|
||||
case s.send <- b:
|
||||
default:
|
||||
s.closeSlow()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,7 +6,6 @@ import (
|
|||
"reflect"
|
||||
"strings"
|
||||
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/evcc-io/evcc/util/encode"
|
||||
)
|
||||
|
||||
|
|
@ -31,36 +30,22 @@ func encodeSliceAsString(v any) (string, error) {
|
|||
return fmt.Sprintf("[%s]", strings.Join(res, ",")), nil
|
||||
}
|
||||
|
||||
func kv(p util.Param) string {
|
||||
func socketEncode(pval any) string {
|
||||
var (
|
||||
val string
|
||||
err error
|
||||
)
|
||||
|
||||
// unwrap slices
|
||||
if p.Val != nil && reflect.TypeOf(p.Val).Kind() == reflect.Slice {
|
||||
val, err = encodeSliceAsString(p.Val)
|
||||
if rv := reflect.ValueOf(pval); pval != nil && rv.Kind() == reflect.Slice && !rv.IsNil() {
|
||||
val, err = encodeSliceAsString(pval)
|
||||
} else {
|
||||
val, err = encodeAsString(p.Val)
|
||||
val, err = encodeAsString(pval)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
if p.Key == "" && val == "" {
|
||||
log.ERROR.Printf("invalid key/val for %+v, please report to https://github.com/evcc-io/evcc/issues/6439", p)
|
||||
return "\"foo\":\"bar\""
|
||||
}
|
||||
|
||||
var msg strings.Builder
|
||||
msg.WriteString("\"")
|
||||
if p.Loadpoint != nil {
|
||||
msg.WriteString(fmt.Sprintf("loadpoints.%d.", *p.Loadpoint))
|
||||
}
|
||||
msg.WriteString(p.Key)
|
||||
msg.WriteString("\":")
|
||||
msg.WriteString(val)
|
||||
|
||||
return msg.String()
|
||||
return val
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
interval: 0.25s
|
||||
interval: 0.1s
|
||||
|
||||
site:
|
||||
title: Plan
|
||||
|
|
|
|||
|
|
@ -4,7 +4,6 @@ import (
|
|||
"maps"
|
||||
"slices"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/evcc-io/evcc/util/encode"
|
||||
|
|
@ -19,15 +18,11 @@ type Param struct {
|
|||
|
||||
// UniqueID returns unique identifier for parameter Loadpoint/Key combination
|
||||
func (p Param) UniqueID() string {
|
||||
var b strings.Builder
|
||||
|
||||
if p.Loadpoint != nil {
|
||||
b.WriteString(strconv.Itoa(*p.Loadpoint) + ".")
|
||||
return strconv.Itoa(*p.Loadpoint) + "." + p.Key
|
||||
}
|
||||
|
||||
b.WriteString(p.Key)
|
||||
|
||||
return b.String()
|
||||
return p.Key
|
||||
}
|
||||
|
||||
// ParamCache is a data store
|
||||
|
|
|
|||
81
util/param_shard.go
Normal file
81
util/param_shard.go
Normal file
|
|
@ -0,0 +1,81 @@
|
|||
package util
|
||||
|
||||
import (
|
||||
"crypto/sha256"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/fatih/structs"
|
||||
)
|
||||
|
||||
// Sharder splits data into chunks, omitting unmodified chunks
|
||||
type Sharder interface {
|
||||
Shards() []Shard
|
||||
}
|
||||
|
||||
type Shard struct {
|
||||
Key string
|
||||
Value any
|
||||
}
|
||||
|
||||
type sharderImpl struct {
|
||||
prefix string
|
||||
struc any
|
||||
}
|
||||
|
||||
// shared shard cache
|
||||
var (
|
||||
shardCache = make(map[string][32]byte)
|
||||
shardMu sync.Mutex
|
||||
)
|
||||
|
||||
func (s *sharderImpl) MarshalJSON() ([]byte, error) {
|
||||
return json.Marshal(s.struc)
|
||||
}
|
||||
|
||||
func (s *sharderImpl) Shards() []Shard {
|
||||
ff := structs.Fields(s.struc)
|
||||
res := make([]Shard, 0, len(ff))
|
||||
|
||||
shardMu.Lock()
|
||||
defer shardMu.Unlock()
|
||||
|
||||
for _, f := range ff {
|
||||
key := f.Name()
|
||||
if t := f.Tag("json"); t != "" {
|
||||
if n := strings.Split(t, ",")[0]; n != "" {
|
||||
key = n
|
||||
}
|
||||
}
|
||||
|
||||
// Use JSON for stable hashing (fmt.Append includes pointer addresses)
|
||||
b, err := json.Marshal(f.Value())
|
||||
if err != nil {
|
||||
// Fallback to fmt.Append if JSON fails
|
||||
b = fmt.Append(nil, f.Value())
|
||||
}
|
||||
|
||||
hash := sha256.Sum256(b)
|
||||
if cached, ok := shardCache[s.prefix+key]; ok && hash == cached {
|
||||
continue
|
||||
}
|
||||
shardCache[s.prefix+key] = hash
|
||||
|
||||
res = append(res, Shard{
|
||||
Key: key,
|
||||
Value: f.Value(),
|
||||
})
|
||||
}
|
||||
|
||||
return res
|
||||
}
|
||||
|
||||
var _ Sharder = (*sharderImpl)(nil)
|
||||
|
||||
// NewSharder creates a Sharder that splits structs into sub-structs for space-efficient socket publishing
|
||||
// Passing anything else than a struct will panic
|
||||
func NewSharder(prefix string, struc any) Sharder {
|
||||
return &sharderImpl{prefix, struc}
|
||||
}
|
||||
12
util/tee.go
12
util/tee.go
|
|
@ -1,7 +1,6 @@
|
|||
package util
|
||||
|
||||
import (
|
||||
"reflect"
|
||||
"sync"
|
||||
|
||||
"github.com/evcc-io/evcc/api"
|
||||
|
|
@ -37,11 +36,12 @@ func (t *Tee) add(out chan<- Param) {
|
|||
// Run starts parameter distribution
|
||||
func (t *Tee) Run(in <-chan Param) {
|
||||
for msg := range in {
|
||||
if val := reflect.ValueOf(msg.Val); val.Kind() == reflect.Ptr {
|
||||
if ptr := reflect.Indirect(val); ptr.IsValid() {
|
||||
msg.Val = ptr.Addr().Elem().Interface()
|
||||
}
|
||||
}
|
||||
// TODO MUST NOT PUBLISH POINTERS (WHO'S VALUES ARE LATER MODIFIED)
|
||||
// if val := reflect.ValueOf(msg.Val); val.Kind() == reflect.Ptr {
|
||||
// if ptr := reflect.Indirect(val); ptr.IsValid() {
|
||||
// fmt.Println("DANGER pointer value:", msg.Key)
|
||||
// }
|
||||
// }
|
||||
|
||||
if val, ok := (msg.Val).(api.Redactor); ok {
|
||||
msg.Val = val.Redacted()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue