chore: refactor param cache (#18377)

This commit is contained in:
andig 2025-01-23 16:16:21 +01:00 • committed by GitHub
parent 07fa1b64f8
commit 9e9c15ed58
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
14 changed files with 128 additions and 146 deletions

View file

@ -153,7 +153,7 @@ func runRoot(cmd *cobra.Command, args []string) {
go tee.Run(valueChan)
// value cache
cache := util.NewCache()
cache := util.NewParamCache()
go cache.Run(pipe.NewDropper(ignoreLogs...).Pipe(tee.Attach()))
// create web server

View file

@ -134,7 +134,7 @@ If you know what you're doing, you can skip the database check with the --ignore
}
func isWritable(filePath string) bool {
file, err := os.OpenFile(filePath, os.O_WRONLY, 0666)
file, err := os.OpenFile(filePath, os.O_WRONLY, 0o666)
if err != nil {
return false
}
@ -709,7 +709,7 @@ func configureEEBus(conf *eebus.Config) error {
}
// setup messaging
func configureMessengers(conf *globalconfig.Messaging, vehicles push.Vehicles, valueChan chan<- util.Param, cache *util.Cache) (chan push.Event, error) {
func configureMessengers(conf *globalconfig.Messaging, vehicles push.Vehicles, valueChan chan<- util.Param, cache *util.ParamCache) (chan push.Event, error) {
// migrate settings
if settings.Exists(keys.Messaging) {
if err := settings.Yaml(keys.Messaging, new(map[string]any), &conf); err != nil {

View file

@ -511,14 +511,14 @@ func TestSetModeAndSocAtDisconnect(t *testing.T) {
}
// cacheExpecter can be used to verify asynchronously written values from cache
func cacheExpecter(t *testing.T, lp *Loadpoint) (*util.Cache, func(key string, val interface{})) {
func cacheExpecter(t *testing.T, lp *Loadpoint) (*util.ParamCache, func(key string, val interface{})) {
t.Helper()
// attach cache for verifying values
paramC := make(chan util.Param)
lp.uiChan = paramC
cache := util.NewCache()
cache := util.NewParamCache()
go cache.Run(paramC)
expect := func(key string, val interface{}) {

View file

@ -9,13 +9,9 @@ import (
"github.com/asaskevich/EventBus"
"github.com/benbjohnson/clock"
"github.com/evcc-io/evcc/api"
"github.com/evcc-io/evcc/util"
)
var (
bus = EventBus.New()
log = util.NewLogger("cache")
)
var bus = EventBus.New()
const (
reset = "reset"

View file

@ -138,9 +138,6 @@ func (p *HTTP) WithPipeline(pipeline *pipeline.Pipeline) *HTTP {
func (p *HTTP) WithAuth(typ, user, password string) (*HTTP, error) {
switch strings.ToLower(typ) {
case "basic":
basicAuth := transport.BasicAuthHeader(user, password)
log.Redact(basicAuth)
p.Client.Transport = transport.BasicAuth(user, password, p.Client.Transport)
case "bearer":
p.Client.Transport = transport.BearerAuth(password, p.Client.Transport)

View file

@ -47,12 +47,12 @@ func NewScriptPluginFromConfig(other map[string]interface{}) (Plugin, error) {
return nil, err
}
p, err := NewScripPlugin(cc.Cmd, cc.Timeout, cc.Scale, cc.Cache)
p, err := NewScriptPlugin(cc.Cmd, cc.Timeout, cc.Scale, cc.Cache)
p.getter = defaultGetters(p, cc.Scale)
if err == nil {
var pipe *pipeline.Pipeline
pipe, err = pipeline.New(log, cc.Settings)
pipe, err = pipeline.New(p.log, cc.Settings)
p.pipeline = pipe
}
@ -61,7 +61,7 @@ func NewScriptPluginFromConfig(other map[string]interface{}) (Plugin, error) {
// NewScriptProvider creates a script plugin.
// Script execution is aborted after given timeout.
func NewScripPlugin(script string, timeout time.Duration, scale float64, cache time.Duration) (*Script, error) {
func NewScriptPlugin(script string, timeout time.Duration, scale float64, cache time.Duration) (*Script, error) {
if strings.TrimSpace(script) == "" {
return nil, errors.New("script is required")
}

View file

@ -30,12 +30,12 @@ type Vehicles interface {
type Hub struct {
definitions map[string]EventTemplateConfig
sender []Messenger
cache *util.Cache
cache *util.ParamCache
vehicles Vehicles
}
// NewHub creates push hub with definitions and receiver
func NewHub(cc map[string]EventTemplateConfig, vv Vehicles, cache *util.Cache) (*Hub, error) {
func NewHub(cc map[string]EventTemplateConfig, vv Vehicles, cache *util.ParamCache) (*Hub, error) {
// instantiate all event templates
for k, v := range cc {
if _, err := template.New("out").Funcs(sprig.FuncMap()).Parse(v.Title); err != nil {

View file

@ -192,7 +192,7 @@ func (s *HTTPd) RegisterSiteHandlers(site site.API, valueChan chan<- util.Param)
}
// RegisterSystemHandler provides system level handlers
func (s *HTTPd) RegisterSystemHandler(valueChan chan<- util.Param, cache *util.Cache, auth auth.Auth, shutdown func()) {
func (s *HTTPd) RegisterSystemHandler(valueChan chan<- util.Param, cache *util.ParamCache, auth auth.Auth, shutdown func()) {
router := s.Server.Handler.(*mux.Router)
// api

View file

@ -197,7 +197,7 @@ func updateSmartCostLimit(site site.API) http.HandlerFunc {
}
// stateHandler returns the combined state
func stateHandler(cache *util.Cache) http.HandlerFunc {
func stateHandler(cache *util.ParamCache) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
res := cache.State(encode.NewEncoder(encode.WithDuration()))
for _, k := range ignoreState {

View file

@ -142,7 +142,7 @@ func (h *SocketHub) broadcast(p util.Param) {
}
// Run starts data and status distribution
func (h *SocketHub) Run(in <-chan util.Param, cache *util.Cache) {
func (h *SocketHub) Run(in <-chan util.Param, cache *util.ParamCache) {
for {
select {
case client := <-h.register:

View file

@ -1,114 +0,0 @@
package util
import (
"fmt"
"maps"
"slices"
"sync"
"github.com/evcc-io/evcc/util/encode"
)
// Cache is a data store
type Cache struct {
mu sync.RWMutex
val map[string]Param
}
// flush is the value type used as parameter for flushing the cache.
// Flushing is implemented by closing the channel. At this time, it is guaranteed
// that the cache has catched up processing all pending messages.
type flush chan struct{}
// Flusher returns a new flush channel
func Flusher() flush {
return make(flush)
}
// NewCache creates cache
func NewCache() *Cache {
return &Cache{
val: make(map[string]Param),
}
}
// Run adds input channel's values to cache
func (c *Cache) Run(in <-chan Param) {
log := NewLogger("cache")
for p := range in {
if flushC, ok := p.Val.(flush); ok {
close(flushC)
continue
}
key := p.Key
if p.Loadpoint != nil {
key = fmt.Sprintf("lp-%d/%s", *p.Loadpoint+1, key)
}
log.TRACE.Printf("%s: %v", key, p.Val)
c.Add(p.UniqueID(), p)
}
}
// State provides a structured copy of the cached values.
// Loadpoints are aggregated as loadpoints array.
// Result values are formatted using encoder.
func (c *Cache) State(enc encode.Encoder) map[string]any {
c.mu.RLock()
defer c.mu.RUnlock()
res := make(map[string]any)
lps := make(map[int]map[string]any)
for _, param := range c.val {
if param.Loadpoint == nil {
res[param.Key] = enc.Encode(param.Val)
} else {
lp, ok := lps[*param.Loadpoint]
if !ok {
lp = make(map[string]any)
lps[*param.Loadpoint] = lp
}
lp[param.Key] = enc.Encode(param.Val)
}
}
// convert map to array
loadpoints := make([]map[string]any, len(lps))
for id, lp := range lps {
loadpoints[id] = lp
}
res["loadpoints"] = loadpoints
return res
}
// All provides a copy of the cached values
func (c *Cache) All() []Param {
c.mu.RLock()
defer c.mu.RUnlock()
return slices.Collect(maps.Values(c.val))
}
// Add entry to cache
func (c *Cache) Add(key string, param Param) {
c.mu.Lock()
defer c.mu.Unlock()
c.val[key] = param
}
// Get entry from cache
func (c *Cache) Get(key string) Param {
c.mu.RLock()
defer c.mu.RUnlock()
if val, ok := c.val[key]; ok {
return val
}
return Param{}
}

View file

@ -1,11 +0,0 @@
package util
import (
"testing"
)
func TestCache(t *testing.T) {
c := NewCache()
c.Add("foo", Param{})
}

View file

@ -1,8 +1,14 @@
package util
import (
"fmt"
"maps"
"slices"
"strconv"
"strings"
"sync"
"github.com/evcc-io/evcc/util/encode"
)
// Param is the broadcast channel data type
@ -24,3 +30,107 @@ func (p Param) UniqueID() string {
return b.String()
}
// ParamCache is a data store
type ParamCache struct {
mu sync.RWMutex
val map[string]Param
}
// flush is the value type used as parameter for flushing the cache.
// Flushing is implemented by closing the channel. At this time, it is guaranteed
// that the cache has catched up processing all pending messages.
type flush chan struct{}
// Flusher returns a new flush channel
func Flusher() flush {
return make(flush)
}
// NewCache creates cache
func NewParamCache() *ParamCache {
return &ParamCache{
val: make(map[string]Param),
}
}
// Run adds input channel's values to cache
func (c *ParamCache) Run(in <-chan Param) {
log := NewLogger("cache")
for p := range in {
if flushC, ok := p.Val.(flush); ok {
close(flushC)
continue
}
key := p.Key
if p.Loadpoint != nil {
key = fmt.Sprintf("lp-%d/%s", *p.Loadpoint+1, key)
}
log.TRACE.Printf("%s: %v", key, p.Val)
c.Add(p.UniqueID(), p)
}
}
// State provides a structured copy of the cached values.
// Loadpoints are aggregated as loadpoints array.
// Result values are formatted using encoder.
func (c *ParamCache) State(enc encode.Encoder) map[string]any {
c.mu.RLock()
defer c.mu.RUnlock()
res := make(map[string]any)
lps := make(map[int]map[string]any)
for _, param := range c.val {
if param.Loadpoint == nil {
res[param.Key] = enc.Encode(param.Val)
} else {
lp, ok := lps[*param.Loadpoint]
if !ok {
lp = make(map[string]any)
lps[*param.Loadpoint] = lp
}
lp[param.Key] = enc.Encode(param.Val)
}
}
// convert map to array
loadpoints := make([]map[string]any, len(lps))
for id, lp := range lps {
loadpoints[id] = lp
}
res["loadpoints"] = loadpoints
return res
}
// All provides a copy of the cached values
func (c *ParamCache) All() []Param {
c.mu.RLock()
defer c.mu.RUnlock()
return slices.Collect(maps.Values(c.val))
}
// Add entry to cache
func (c *ParamCache) Add(key string, param Param) {
c.mu.Lock()
defer c.mu.Unlock()
c.val[key] = param
}
// Get entry from cache
func (c *ParamCache) Get(key string) Param {
c.mu.RLock()
defer c.mu.RUnlock()
if val, ok := c.val[key]; ok {
return val
}
return Param{}
}

View file

@ -17,3 +17,7 @@ func TestParam(t *testing.T) {
p.Loadpoint = &lp
assert.Equal(t, "2.power", p.UniqueID())
}
func TestParamCache(t *testing.T) {
NewParamCache().Add("foo", Param{})
}