diff --git a/cmd/root.go b/cmd/root.go index 54fe77e39..c005de6e2 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -399,7 +399,7 @@ func runRoot(cmd *cobra.Command, args []string) { // setup messaging var pushChan chan messenger.Event if err == nil { - pushChan, err = configureMessengers(&conf.Messaging, &conf.MessagingEvents, site.Vehicles(), valueChan, cache) + pushChan, err = configureMessengers(&conf.Messaging, &conf.MessagingEvents, site.Vehicles()) err = wrapErrorWithClass(ClassMessenger, err) } diff --git a/cmd/setup.go b/cmd/setup.go index cfdfca3ab..d4831e240 100644 --- a/cmd/setup.go +++ b/cmd/setup.go @@ -973,7 +973,7 @@ func configureEEBus(conf *eebus.Config) error { return nil } -func configureMessengers(confMessaging *globalconfig.Messaging, confEvents *globalconfig.MessagingEvents, vehicles messenger.Vehicles, valueChan chan<- util.Param, cache *util.ParamCache) (chan messenger.Event, error) { +func configureMessengers(confMessaging *globalconfig.Messaging, confEvents *globalconfig.MessagingEvents, vehicles messenger.Vehicles) (chan messenger.Event, error) { // yaml config from file if len(confMessaging.Events) != 0 || len(confMessaging.Services) != 0 { yamlSource.messaging = globalconfig.YamlSourceFile @@ -1002,7 +1002,8 @@ func configureMessengers(confMessaging *globalconfig.Messaging, confEvents *glob } } - messageChan := make(chan messenger.Event, 1) + // events are queued by the value cache, keep enough slack to not stall it + messageChan := make(chan messenger.Event, 16) var eg errgroup.Group @@ -1043,7 +1044,7 @@ func configureMessengers(confMessaging *globalconfig.Messaging, confEvents *glob events = confMessaging.Events } - messageHub, err := messenger.NewHub(events, vehicles, cache) + messageHub, err := messenger.NewHub(events, vehicles) if err != nil { return messageChan, fmt.Errorf("failed configuring push services: %w", err) @@ -1055,7 +1056,7 @@ func configureMessengers(confMessaging *globalconfig.Messaging, confEvents *glob } } - go messageHub.Run(messageChan, valueChan) + go messageHub.Run(messageChan) return messageChan, nil } diff --git a/core/site.go b/core/site.go index c25e2db86..9f390bf59 100644 --- a/core/site.go +++ b/core/site.go @@ -1234,6 +1234,21 @@ func (site *Site) prepare() { vehicle.ClearPlanLocks = site.clearPlanLocks } +// pushEvent queues the event in the value stream. The cache attaches its state +// when the event reaches its position, so the message renders exactly the +// values published before the event was raised. +func (site *Site) pushEvent(ev messenger.Event) { + pushChan := site.pushChan + if pushChan == nil { + return + } + + site.valueChan <- util.Param{Val: util.Snapshot(func(state []util.Param) { + ev.State = state + pushChan <- ev + })} +} + // Prepare attaches communication channels to site and loadpoints func (site *Site) Prepare(valueChan chan<- util.Param, pushChan chan<- messenger.Event) { site.pushChan = pushChan @@ -1272,7 +1287,7 @@ func (site *Site) Prepare(valueChan chan<- util.Param, pushChan chan<- messenger site.valueChan <- param case ev := <-lpPushChan: ev.Loadpoint = &id - pushChan <- ev + site.pushEvent(ev) } } }(id) diff --git a/core/site_optimizer.go b/core/site_optimizer.go index f22b5bf2c..e3202e0e5 100644 --- a/core/site_optimizer.go +++ b/core/site_optimizer.go @@ -627,10 +627,8 @@ func (site *Site) optimizerUpdate(battery []types.Measurement) error { site.publishSuggestions() // notify on actionable suggestion changes (advisory only, see #31903) - if site.pushChan != nil { - for _, ev := range site.diffSuggestions(site.pendingSuggestions(details.BatteryDetails)) { - site.pushChan <- ev - } + for _, ev := range site.diffSuggestions(site.pendingSuggestions(details.BatteryDetails)) { + site.pushEvent(ev) } return nil diff --git a/docs/agents/core-domain.md b/docs/agents/core-domain.md index 42ba9ac3f..14585dafc 100644 --- a/docs/agents/core-domain.md +++ b/docs/agents/core-domain.md @@ -122,7 +122,7 @@ effectivePrice = gridPrice * (1 - greenShare) + feedInPrice * greenShare |---------|-------|--------|---------| | `valueChan` | Site | Unbounded (`chanx.NewUnboundedChan`) | State changes -> DB + UI (ordering) | | `lpUpdateChan` | Site | 1 | Early loadpoint update requests | -| `pushChan` | Loadpoint | Buffered | User notifications | +| `pushChan` | Loadpoint | 16 | User notifications, queued via `valueChan` so the message renders the state at event time | ## Tariff Integration diff --git a/messenger/hub.go b/messenger/hub.go index cd50bc9b7..22bf694e2 100644 --- a/messenger/hub.go +++ b/messenger/hub.go @@ -18,6 +18,7 @@ type Event struct { Loadpoint *int // optional loadpoint id Event string Attributes map[string]any // optional event-specific template attributes + State []util.Param // cache state at the time the event was raised } type Vehicles interface { @@ -29,12 +30,11 @@ type Vehicles interface { type Hub struct { definitions globalconfig.MessagingEvents sender []api.Messenger - cache *util.ParamCache vehicles Vehicles } // NewHub creates push hub with definitions and receiver -func NewHub(cc globalconfig.MessagingEvents, vv Vehicles, cache *util.ParamCache) (*Hub, error) { +func NewHub(cc globalconfig.MessagingEvents, vv Vehicles) (*Hub, error) { // keep only enabled events filtered := make(globalconfig.MessagingEvents, len(cc)) @@ -56,7 +56,6 @@ func NewHub(cc globalconfig.MessagingEvents, vv Vehicles, cache *util.ParamCache h := &Hub{ definitions: filtered, - cache: cache, vehicles: vv, } @@ -77,8 +76,8 @@ func (h *Hub) apply(ev Event, tmpl string) (string, error) { attr["loadpoint"] = *ev.Loadpoint + 1 } - // get all values from cache - for _, p := range h.cache.All() { + // get all values from the event's cache state + for _, p := range ev.State { if p.Loadpoint == nil || ev.Loadpoint == p.Loadpoint { val := p.Val @@ -114,7 +113,7 @@ func (h *Hub) apply(ev Event, tmpl string) (string, error) { } // Run is the Hub's main publishing loop -func (h *Hub) Run(events <-chan Event, valueChan chan<- util.Param) { +func (h *Hub) Run(events <-chan Event) { log := util.NewLogger("push") for ev := range events { @@ -127,11 +126,6 @@ func (h *Hub) Run(events <-chan Event, valueChan chan<- util.Param) { 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/param.go b/util/param.go index 252fdcc56..a03eb30df 100644 --- a/util/param.go +++ b/util/param.go @@ -31,15 +31,9 @@ type ParamCache 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) -} +// Snapshot requests a copy of the cache state at the parameter's position in +// the stream. It runs on the cache's goroutine and must not block for long. +type Snapshot func([]Param) // NewCache creates cache func NewParamCache() *ParamCache { @@ -51,8 +45,8 @@ func NewParamCache() *ParamCache { // Run adds input channel's values to cache func (c *ParamCache) Run(in <-chan Param) { for p := range in { - if flushC, ok := p.Val.(flush); ok { - close(flushC) + if snapshot, ok := p.Val.(Snapshot); ok { + snapshot(c.All()) continue } diff --git a/util/param_test.go b/util/param_test.go index 1ab9c66b2..15f45a2ec 100644 --- a/util/param_test.go +++ b/util/param_test.go @@ -21,3 +21,20 @@ func TestParam(t *testing.T) { func TestParamCache(t *testing.T) { NewParamCache().Add("foo", Param{}) } + +func TestParamCacheSnapshot(t *testing.T) { + in := make(chan Param) + go NewParamCache().Run(in) + + in <- Param{Key: "before", Val: 1} + + res := make(chan []Param, 1) + in <- Param{Val: Snapshot(func(state []Param) { res <- state })} + + // published after the snapshot request, must not be included + in <- Param{Key: "after", Val: 2} + + state := <-res + assert.Len(t, state, 1) + assert.Equal(t, "before", state[0].Key) +}