diff --git a/charger/eebus-ohpcf.go b/charger/eebus-ohpcf.go index 14cb898e3..7f8798453 100644 --- a/charger/eebus-ohpcf.go +++ b/charger/eebus-ohpcf.go @@ -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() } diff --git a/charger/eebus-ohpcf_lpc_test.go b/charger/eebus-ohpcf_lpc_test.go index 2fd3cc1af..119eced6f 100644 --- a/charger/eebus-ohpcf_lpc_test.go +++ b/charger/eebus-ohpcf_lpc_test.go @@ -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, diff --git a/hems/eebus/eebus.go b/hems/eebus/eebus.go index f8d6fea74..4aa43592f 100644 --- a/hems/eebus/eebus.go +++ b/hems/eebus/eebus.go @@ -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,14 +194,21 @@ 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) { - if err := c.run(); err != nil { - c.log.ERROR.Println(err) - } + for tick := time.Tick(c.interval); ; { + select { + case <-tick: + if err := c.run(); err != nil { + c.log.ERROR.Println(err) + } - if c.publishFunc != nil { - c.publishFunc() + if c.publishFunc != nil { + c.publishFunc() + } + + case <-c.ctx.Done(): + return } } } diff --git a/hems/eebus/eebus_test.go b/hems/eebus/eebus_test.go index 26d3fcc16..da1f44a93 100644 --- a/hems/eebus/eebus_test.go +++ b/hems/eebus/eebus_test.go @@ -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") + } +} diff --git a/meter/eebus.go b/meter/eebus.go index 9d4d42451..69262af24 100644 --- a/meter/eebus.go +++ b/meter/eebus.go @@ -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(), diff --git a/meter/eebus_events.go b/meter/eebus_events.go index bdba484bc..57edc2cd4 100644 --- a/meter/eebus_events.go +++ b/meter/eebus_events.go @@ -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()) }) } } diff --git a/meter/eebus_lpc_lpp_test.go b/meter/eebus_lpc_lpp_test.go index 72e5e2fa6..2087d86d2 100644 --- a/meter/eebus_lpc_lpp_test.go +++ b/meter/eebus_lpc_lpp_test.go @@ -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, diff --git a/server/eebus/helper.go b/server/eebus/helper.go index 0c09302f2..46316335e 100644 --- a/server/eebus/helper.go +++ b/server/eebus/helper.go @@ -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) } }