From 2011b0b91de4a7ab302f3ec6f4c00f0694576feb Mon Sep 17 00:00:00 2001 From: andig Date: Sat, 18 Feb 2023 14:01:40 +0100 Subject: [PATCH] chore: flush cache without using timeouts (#6303) --- cmd/root.go | 2 +- cmd/setup.go | 4 ++-- push/hub.go | 11 ++++++----- util/cache.go | 16 ++++++++++++++++ 4 files changed, 25 insertions(+), 8 deletions(-) diff --git a/cmd/root.go b/cmd/root.go index 8a89f010f..c4ed0023c 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -205,7 +205,7 @@ func runRoot(cmd *cobra.Command, args []string) { // setup messaging var pushChan chan push.Event if err == nil { - pushChan, err = configureMessengers(conf.Messaging, cache) + pushChan, err = configureMessengers(conf.Messaging, valueChan, cache) } // run shutdown functions on stop diff --git a/cmd/setup.go b/cmd/setup.go index 828542c13..3764e0913 100644 --- a/cmd/setup.go +++ b/cmd/setup.go @@ -203,7 +203,7 @@ func configureEEBus(conf map[string]interface{}) error { } // setup messaging -func configureMessengers(conf messagingConfig, cache *util.Cache) (chan push.Event, error) { +func configureMessengers(conf messagingConfig, valueChan chan util.Param, cache *util.Cache) (chan push.Event, error) { messageChan := make(chan push.Event, 1) messageHub, err := push.NewHub(conf.Events, cache) @@ -219,7 +219,7 @@ func configureMessengers(conf messagingConfig, cache *util.Cache) (chan push.Eve messageHub.Add(impl) } - go messageHub.Run(messageChan) + go messageHub.Run(messageChan, valueChan) return messageChan, nil } diff --git a/push/hub.go b/push/hub.go index 2bff7dae7..0663dbc8d 100644 --- a/push/hub.go +++ b/push/hub.go @@ -3,7 +3,6 @@ package push import ( "strings" "text/template" - "time" "github.com/Masterminds/sprig/v3" "github.com/evcc-io/evcc/util" @@ -70,9 +69,6 @@ func (h *Hub) Add(sender Messenger) { func (h *Hub) apply(ev Event, tmpl *template.Template) (string, error) { attr := make(map[string]interface{}) - // let cache catch up, refs reverted https://github.com/evcc-io/evcc/pull/445 - time.Sleep(100 * time.Millisecond) - // loadpoint id if ev.Loadpoint != nil { attr["loadpoint"] = *ev.Loadpoint + 1 @@ -95,7 +91,7 @@ func (h *Hub) apply(ev Event, tmpl *template.Template) (string, error) { } // Run is the Hub's main publishing loop -func (h *Hub) Run(events <-chan Event) { +func (h *Hub) Run(events <-chan Event, valueChan chan util.Param) { log := util.NewLogger("push") for ev := range events { @@ -108,6 +104,11 @@ func (h *Hub) Run(events <-chan Event) { continue } + // let cache catch up, refs https://github.com/evcc-io/evcc/pull/445 + flushC := util.Flusher() + valueChan <- util.Param{Val: flushC} + <-flushC + title, err := h.apply(ev, definition.Title) if err != nil { log.ERROR.Printf("invalid title template for %s: %v", ev.Event, err) diff --git a/util/cache.go b/util/cache.go index 6435bd02e..2818d4c34 100644 --- a/util/cache.go +++ b/util/cache.go @@ -11,6 +11,16 @@ type Cache struct { 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{ @@ -23,10 +33,16 @@ 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) }