chore: flush cache without using timeouts (#6303)
This commit is contained in:
parent
2080452ae4
commit
2011b0b91d
4 changed files with 25 additions and 8 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
11
push/hub.go
11
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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue