Fix rendering notifications from the state at event time (#32401)
This commit is contained in:
parent
ad255ff54c
commit
a4da257bb5
8 changed files with 52 additions and 33 deletions
|
|
@ -399,7 +399,7 @@ func runRoot(cmd *cobra.Command, args []string) {
|
||||||
// setup messaging
|
// setup messaging
|
||||||
var pushChan chan messenger.Event
|
var pushChan chan messenger.Event
|
||||||
if err == nil {
|
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)
|
err = wrapErrorWithClass(ClassMessenger, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -973,7 +973,7 @@ func configureEEBus(conf *eebus.Config) error {
|
||||||
return nil
|
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
|
// yaml config from file
|
||||||
if len(confMessaging.Events) != 0 || len(confMessaging.Services) != 0 {
|
if len(confMessaging.Events) != 0 || len(confMessaging.Services) != 0 {
|
||||||
yamlSource.messaging = globalconfig.YamlSourceFile
|
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
|
var eg errgroup.Group
|
||||||
|
|
||||||
|
|
@ -1043,7 +1044,7 @@ func configureMessengers(confMessaging *globalconfig.Messaging, confEvents *glob
|
||||||
events = confMessaging.Events
|
events = confMessaging.Events
|
||||||
}
|
}
|
||||||
|
|
||||||
messageHub, err := messenger.NewHub(events, vehicles, cache)
|
messageHub, err := messenger.NewHub(events, vehicles)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return messageChan, fmt.Errorf("failed configuring push services: %w", err)
|
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
|
return messageChan, nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
17
core/site.go
17
core/site.go
|
|
@ -1234,6 +1234,21 @@ func (site *Site) prepare() {
|
||||||
vehicle.ClearPlanLocks = site.clearPlanLocks
|
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
|
// Prepare attaches communication channels to site and loadpoints
|
||||||
func (site *Site) Prepare(valueChan chan<- util.Param, pushChan chan<- messenger.Event) {
|
func (site *Site) Prepare(valueChan chan<- util.Param, pushChan chan<- messenger.Event) {
|
||||||
site.pushChan = pushChan
|
site.pushChan = pushChan
|
||||||
|
|
@ -1272,7 +1287,7 @@ func (site *Site) Prepare(valueChan chan<- util.Param, pushChan chan<- messenger
|
||||||
site.valueChan <- param
|
site.valueChan <- param
|
||||||
case ev := <-lpPushChan:
|
case ev := <-lpPushChan:
|
||||||
ev.Loadpoint = &id
|
ev.Loadpoint = &id
|
||||||
pushChan <- ev
|
site.pushEvent(ev)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}(id)
|
}(id)
|
||||||
|
|
|
||||||
|
|
@ -627,10 +627,8 @@ func (site *Site) optimizerUpdate(battery []types.Measurement) error {
|
||||||
site.publishSuggestions()
|
site.publishSuggestions()
|
||||||
|
|
||||||
// notify on actionable suggestion changes (advisory only, see #31903)
|
// notify on actionable suggestion changes (advisory only, see #31903)
|
||||||
if site.pushChan != nil {
|
for _, ev := range site.diffSuggestions(site.pendingSuggestions(details.BatteryDetails)) {
|
||||||
for _, ev := range site.diffSuggestions(site.pendingSuggestions(details.BatteryDetails)) {
|
site.pushEvent(ev)
|
||||||
site.pushChan <- ev
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
|
|
|
||||||
|
|
@ -122,7 +122,7 @@ effectivePrice = gridPrice * (1 - greenShare) + feedInPrice * greenShare
|
||||||
|---------|-------|--------|---------|
|
|---------|-------|--------|---------|
|
||||||
| `valueChan` | Site | Unbounded (`chanx.NewUnboundedChan`) | State changes -> DB + UI (ordering) |
|
| `valueChan` | Site | Unbounded (`chanx.NewUnboundedChan`) | State changes -> DB + UI (ordering) |
|
||||||
| `lpUpdateChan` | Site | 1 | Early loadpoint update requests |
|
| `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
|
## Tariff Integration
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ type Event struct {
|
||||||
Loadpoint *int // optional loadpoint id
|
Loadpoint *int // optional loadpoint id
|
||||||
Event string
|
Event string
|
||||||
Attributes map[string]any // optional event-specific template attributes
|
Attributes map[string]any // optional event-specific template attributes
|
||||||
|
State []util.Param // cache state at the time the event was raised
|
||||||
}
|
}
|
||||||
|
|
||||||
type Vehicles interface {
|
type Vehicles interface {
|
||||||
|
|
@ -29,12 +30,11 @@ type Vehicles interface {
|
||||||
type Hub struct {
|
type Hub struct {
|
||||||
definitions globalconfig.MessagingEvents
|
definitions globalconfig.MessagingEvents
|
||||||
sender []api.Messenger
|
sender []api.Messenger
|
||||||
cache *util.ParamCache
|
|
||||||
vehicles Vehicles
|
vehicles Vehicles
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewHub creates push hub with definitions and receiver
|
// 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
|
// keep only enabled events
|
||||||
filtered := make(globalconfig.MessagingEvents, len(cc))
|
filtered := make(globalconfig.MessagingEvents, len(cc))
|
||||||
|
|
||||||
|
|
@ -56,7 +56,6 @@ func NewHub(cc globalconfig.MessagingEvents, vv Vehicles, cache *util.ParamCache
|
||||||
|
|
||||||
h := &Hub{
|
h := &Hub{
|
||||||
definitions: filtered,
|
definitions: filtered,
|
||||||
cache: cache,
|
|
||||||
vehicles: vv,
|
vehicles: vv,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -77,8 +76,8 @@ func (h *Hub) apply(ev Event, tmpl string) (string, error) {
|
||||||
attr["loadpoint"] = *ev.Loadpoint + 1
|
attr["loadpoint"] = *ev.Loadpoint + 1
|
||||||
}
|
}
|
||||||
|
|
||||||
// get all values from cache
|
// get all values from the event's cache state
|
||||||
for _, p := range h.cache.All() {
|
for _, p := range ev.State {
|
||||||
if p.Loadpoint == nil || ev.Loadpoint == p.Loadpoint {
|
if p.Loadpoint == nil || ev.Loadpoint == p.Loadpoint {
|
||||||
val := p.Val
|
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
|
// 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")
|
log := util.NewLogger("push")
|
||||||
|
|
||||||
for ev := range events {
|
for ev := range events {
|
||||||
|
|
@ -127,11 +126,6 @@ func (h *Hub) Run(events <-chan Event, valueChan chan<- util.Param) {
|
||||||
continue
|
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)
|
title, err := h.apply(ev, definition.Title)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.ERROR.Printf("invalid title template for %s: %v", ev.Event, err)
|
log.ERROR.Printf("invalid title template for %s: %v", ev.Event, err)
|
||||||
|
|
|
||||||
|
|
@ -31,15 +31,9 @@ type ParamCache struct {
|
||||||
val map[string]Param
|
val map[string]Param
|
||||||
}
|
}
|
||||||
|
|
||||||
// flush is the value type used as parameter for flushing the cache.
|
// Snapshot requests a copy of the cache state at the parameter's position in
|
||||||
// Flushing is implemented by closing the channel. At this time, it is guaranteed
|
// the stream. It runs on the cache's goroutine and must not block for long.
|
||||||
// that the cache has catched up processing all pending messages.
|
type Snapshot func([]Param)
|
||||||
type flush chan struct{}
|
|
||||||
|
|
||||||
// Flusher returns a new flush channel
|
|
||||||
func Flusher() flush {
|
|
||||||
return make(flush)
|
|
||||||
}
|
|
||||||
|
|
||||||
// NewCache creates cache
|
// NewCache creates cache
|
||||||
func NewParamCache() *ParamCache {
|
func NewParamCache() *ParamCache {
|
||||||
|
|
@ -51,8 +45,8 @@ func NewParamCache() *ParamCache {
|
||||||
// Run adds input channel's values to cache
|
// Run adds input channel's values to cache
|
||||||
func (c *ParamCache) Run(in <-chan Param) {
|
func (c *ParamCache) Run(in <-chan Param) {
|
||||||
for p := range in {
|
for p := range in {
|
||||||
if flushC, ok := p.Val.(flush); ok {
|
if snapshot, ok := p.Val.(Snapshot); ok {
|
||||||
close(flushC)
|
snapshot(c.All())
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -21,3 +21,20 @@ func TestParam(t *testing.T) {
|
||||||
func TestParamCache(t *testing.T) {
|
func TestParamCache(t *testing.T) {
|
||||||
NewParamCache().Add("foo", Param{})
|
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)
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue