EEBus: abort run loop and limit assertion on context cancellation (#32494)
This commit is contained in:
parent
4a878838e7
commit
fed3fb731c
8 changed files with 54 additions and 12 deletions
|
|
@ -189,7 +189,7 @@ func (c *EEBusOHPCF) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spine
|
|||
c.egLpcEntity = entity
|
||||
|
||||
// [LPC-913]: state the limit to the newly available CS
|
||||
go eebus.AssertLimit(c.log, func() error { return c.Dim(c.lastDimmed()) })
|
||||
go eebus.AssertLimit(c.ctx, c.log, func() error { return c.Dim(c.lastDimmed()) })
|
||||
}
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ func newOHPCFEGCharger(t *testing.T) (*EEBusOHPCF, *egmocks.EgLPCInterface, spin
|
|||
entity := spinemocks.NewEntityRemoteInterface(t)
|
||||
|
||||
c := &EEBusOHPCF{
|
||||
ctx: t.Context(),
|
||||
log: util.NewLogger("eebus-ohpcf-test"),
|
||||
eg: &eebus.EnergyGuard{EgLPCInterface: lpc},
|
||||
egLpcEntity: entity,
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ func init() {
|
|||
|
||||
type EEBus struct {
|
||||
mux sync.RWMutex
|
||||
ctx context.Context // device lifetime, aborts Run
|
||||
log *util.Logger
|
||||
|
||||
*eebus.Connector
|
||||
|
|
@ -104,6 +105,7 @@ func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(b
|
|||
}
|
||||
|
||||
c := &EEBus{
|
||||
ctx: ctx,
|
||||
log: util.NewLogger("eebus"),
|
||||
site: site,
|
||||
passthrough: passthrough,
|
||||
|
|
@ -192,8 +194,11 @@ func (c *EEBus) Connect(connected bool) {
|
|||
}
|
||||
}
|
||||
|
||||
// Run applies limits until the device context is cancelled
|
||||
func (c *EEBus) Run() {
|
||||
for range time.Tick(c.interval) {
|
||||
for tick := time.Tick(c.interval); ; {
|
||||
select {
|
||||
case <-tick:
|
||||
if err := c.run(); err != nil {
|
||||
c.log.ERROR.Println(err)
|
||||
}
|
||||
|
|
@ -201,6 +206,10 @@ func (c *EEBus) Run() {
|
|||
if c.publishFunc != nil {
|
||||
c.publishFunc()
|
||||
}
|
||||
|
||||
case <-c.ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package eebus
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
|
|
@ -40,7 +41,9 @@ func newTestEEBus(t *testing.T) *EEBus {
|
|||
|
||||
failsafeProduction := testFailsafeProduction
|
||||
return &EEBus{
|
||||
ctx: t.Context(),
|
||||
log: util.NewLogger("test"),
|
||||
interval: time.Millisecond,
|
||||
site: &stubSite{},
|
||||
Connector: eebus.NewConnector(),
|
||||
heartbeat: util.NewValue[struct{}](time.Hour),
|
||||
|
|
@ -227,3 +230,27 @@ func TestEEBusEdgeTriggered(t *testing.T) {
|
|||
require.Equal(t, 1, calls, "passthrough must fire once on the edge, not every tick")
|
||||
assertConsumptionLimit(t, c, 3000)
|
||||
}
|
||||
|
||||
// TestRunAborts verifies the run loop terminates when the device context is
|
||||
// cancelled instead of ticking for the lifetime of the process.
|
||||
func TestRunAborts(t *testing.T) {
|
||||
c := newTestEEBus(t)
|
||||
c.interval = time.Hour // only the context can end the loop
|
||||
|
||||
ctx, cancel := context.WithCancel(t.Context())
|
||||
c.ctx = ctx
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
c.Run()
|
||||
close(done)
|
||||
}()
|
||||
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("Run did not return on cancelled context")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ import (
|
|||
// Uses MPC (Monitoring & Power Consumption) for all other cases (default)
|
||||
// Additionally supports LPC (Limitation of Power Consumption) and LPP (Limitation of Power Production)
|
||||
type EEBus struct {
|
||||
ctx context.Context // device lifetime, aborts limit retries
|
||||
log *util.Logger
|
||||
|
||||
connector *eebus.Connector
|
||||
|
|
@ -114,6 +115,7 @@ func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api.
|
|||
}
|
||||
|
||||
c := &EEBus{
|
||||
ctx: ctx,
|
||||
log: util.NewLogger("eebus-" + useCase),
|
||||
ma: ma,
|
||||
eg: inst.EnergyGuard(),
|
||||
|
|
|
|||
|
|
@ -81,7 +81,7 @@ func (c *EEBus) egLpcUseCaseSupportUpdate(entity spineapi.EntityRemoteInterface)
|
|||
c.egLpcEntity = entity
|
||||
|
||||
// [LPC-913]: state the limit to the newly available CS
|
||||
go eebus.AssertLimit(c.log, func() error { return c.Dim(c.lastDimmed()) })
|
||||
go eebus.AssertLimit(c.ctx, c.log, func() error { return c.Dim(c.lastDimmed()) })
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -98,6 +98,6 @@ func (c *EEBus) egLppUseCaseSupportUpdate(entity spineapi.EntityRemoteInterface)
|
|||
c.egLppEntity = entity
|
||||
|
||||
// [LPP-913]: state the limit to the newly available CS
|
||||
go eebus.AssertLimit(c.log, func() error { return c.SetCurtailPercent(c.lastCurtailPercent()) })
|
||||
go eebus.AssertLimit(c.ctx, c.log, func() error { return c.SetCurtailPercent(c.lastCurtailPercent()) })
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -30,6 +30,7 @@ func newEGMeter(t *testing.T) (*EEBus, *egmocks.EgLPCInterface, *egmocks.EgLPPIn
|
|||
entity := spinemocks.NewEntityRemoteInterface(t)
|
||||
|
||||
c := &EEBus{
|
||||
ctx: t.Context(),
|
||||
log: util.NewLogger("eebus-eg-test"),
|
||||
eg: &eebus.EnergyGuard{EgLPCInterface: lpc, EgLPPInterface: lpp},
|
||||
egLpcEntity: entity,
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package eebus
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
|
|
@ -52,10 +53,11 @@ const limitTimeout = 50 * time.Second
|
|||
|
||||
// AssertLimit states the current limit to a newly available Controllable System
|
||||
// ([LPC-913]/[LPP-913]). Blocks while retrying- the CS ignores writes that do not
|
||||
// follow a heartbeat and may reject them while still in state "init".
|
||||
func AssertLimit(log *util.Logger, write func() error) {
|
||||
// follow a heartbeat and may reject them while still in state "init". Retrying
|
||||
// stops when ctx is cancelled, i.e. when the device is gone.
|
||||
func AssertLimit(ctx context.Context, log *util.Logger, write func() error) {
|
||||
bo := backoff.NewExponentialBackOff(backoff.WithMaxElapsedTime(limitTimeout))
|
||||
if err := backoff.Retry(write, bo); err != nil {
|
||||
if err := backoff.Retry(write, backoff.WithContext(bo, ctx)); err != nil && ctx.Err() == nil {
|
||||
log.DEBUG.Printf("assert limit: %v", err)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue