From 4a878838e764d33102309873d990a438311a080e Mon Sep 17 00:00:00 2001 From: andig Date: Mon, 3 Aug 2026 18:13:37 +0200 Subject: [PATCH] EEBus: state the limit to a controllable system on connect (#32491) --- charger/eebus-ohpcf.go | 23 ++++++++++-- charger/eebus-ohpcf_lpc_test.go | 28 +++++++++++++++ meter/eebus.go | 54 ++++++++++++++++++++++------ meter/eebus_events.go | 6 ++++ meter/eebus_lpc_lpp_test.go | 62 ++++++++++++++++++++++++++++++++- server/eebus/helper.go | 15 ++++++++ 6 files changed, 175 insertions(+), 13 deletions(-) diff --git a/charger/eebus-ohpcf.go b/charger/eebus-ohpcf.go index ff145c03b..14cb898e3 100644 --- a/charger/eebus-ohpcf.go +++ b/charger/eebus-ohpcf.go @@ -43,6 +43,7 @@ type EEBusOHPCF struct { egLpcEntity spineapi.EntityRemoteInterface enabled bool reboosting bool + dimmed bool // last limit written, re-stated on reconnect connector *eebus.Connector } @@ -186,6 +187,9 @@ func (c *EEBusOHPCF) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spine // use most specific selector if c.egLpcEntity == nil || len(entity.Address().Entity) < len(c.egLpcEntity.Address().Entity) { 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()) }) } c.mu.Unlock() } @@ -212,6 +216,13 @@ func (c *EEBusOHPCF) lastEnabled() bool { return c.enabled } +func (c *EEBusOHPCF) lastDimmed() bool { + c.mu.RLock() + defer c.mu.RUnlock() + + return c.dimmed +} + // ohpcfStatus maps the compressor process state to a charge status: running is // consuming (C), any other connected state is standby (B). Disconnected (A) is handled in Status. func ohpcfStatus(state ucapi.CompressorPowerConsumptionStateType) api.ChargeStatus { @@ -405,9 +416,17 @@ func (c *EEBusOHPCF) Dim(dim bool) error { } // TODO: change api.Dimmer to make the limit configurable; use a fixed 0W safe limit for now - return eebus.Await(func(cb func(model.ResultDataType, model.MsgCounterType)) (*model.MsgCounterType, error) { + if err := eebus.Await(func(cb func(model.ResultDataType, model.MsgCounterType)) (*model.MsgCounterType, error) { return c.eg.EgLPCInterface.WriteConsumptionLimit(entity, ucapi.LoadLimit{Value: 0, IsActive: dim}, cb) - }) + }); err != nil { + return err + } + + c.mu.Lock() + c.dimmed = dim + c.mu.Unlock() + + return nil } // apply issues the command to align the optional consumption with the on/off diff --git a/charger/eebus-ohpcf_lpc_test.go b/charger/eebus-ohpcf_lpc_test.go index 2c526f01a..2fd3cc1af 100644 --- a/charger/eebus-ohpcf_lpc_test.go +++ b/charger/eebus-ohpcf_lpc_test.go @@ -5,9 +5,11 @@ package charger import ( "testing" + "time" eebusapi "github.com/enbility/eebus-go/api" ucapi "github.com/enbility/eebus-go/usecases/api" + eglpc "github.com/enbility/eebus-go/usecases/eg/lpc" egmocks "github.com/enbility/eebus-go/usecases/mocks" spineapi "github.com/enbility/spine-go/api" spinemocks "github.com/enbility/spine-go/mocks" @@ -61,6 +63,32 @@ func TestOHPCF_LPC_EGMessages_ConsumptionLimit(t *testing.T) { } } +// ATC_COM_PT_EGMessages_002 (LPC-913): after initial connection the EG states a +// limit to the CS - deactivated when nothing is being limited. +func TestOHPCF_LPC_InitialLimit(t *testing.T) { + c, lpc, entity := newOHPCFEGCharger(t) + c.egLpcEntity = nil + + written := make(chan ucapi.LoadLimit, 1) + lpc.EXPECT().IsScenarioAvailableAtEntity(entity, eebus.LPCLimit).Return(true) + lpc.EXPECT(). + WriteConsumptionLimit(entity, mock.Anything, mock.Anything). + Run(func(_ spineapi.EntityRemoteInterface, limit ucapi.LoadLimit, cb func(model.ResultDataType, model.MsgCounterType)) { + written <- limit + cb(model.ResultDataType{}, 0) + }). + Return(new(model.MsgCounterType), nil) + + c.UseCaseEvent(nil, entity, eglpc.UseCaseSupportUpdate) + + select { + case limit := <-written: + assert.Equal(t, ucapi.LoadLimit{Value: 0, IsActive: false}, limit) + case <-time.After(time.Second): + t.Fatal("no limit written after connect") + } +} + // A rejected write (NACK) must surface as an error, not silent success. func TestOHPCF_LPC_Dim_WriteRejected(t *testing.T) { c, lpc, entity := newOHPCFEGCharger(t) diff --git a/meter/eebus.go b/meter/eebus.go index 4da224985..9d4d42451 100644 --- a/meter/eebus.go +++ b/meter/eebus.go @@ -35,6 +35,9 @@ type EEBus struct { maEntity spineapi.EntityRemoteInterface egLpcEntity spineapi.EntityRemoteInterface egLppEntity spineapi.EntityRemoteInterface + + dimmed bool // last limits written, re-stated on reconnect// last limits written, re-stated on reconnect + curtailPercent int } // maScenarios holds the spec scenario numbers for the active monitoring use case. @@ -111,12 +114,13 @@ func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage) (api. } c := &EEBus{ - log: util.NewLogger("eebus-" + useCase), - ma: ma, - eg: inst.EnergyGuard(), - mm: mm, - scenarios: scenarios, - connector: eebus.NewConnector(), + log: util.NewLogger("eebus-" + useCase), + ma: ma, + eg: inst.EnergyGuard(), + mm: mm, + scenarios: scenarios, + connector: eebus.NewConnector(), + curtailPercent: 100, } if err := inst.RegisterDevice(ski, ip, c); err != nil { @@ -166,6 +170,20 @@ func eebusReadValue[T any](uc eebusapi.UseCaseBaseInterface, entity spineapi.Ent return res, nil } +func (c *EEBus) lastDimmed() bool { + c.mu.Lock() + defer c.mu.Unlock() + + return c.dimmed +} + +func (c *EEBus) lastCurtailPercent() int { + c.mu.Lock() + defer c.mu.Unlock() + + return c.curtailPercent +} + func (c *EEBus) readValue(scenario uint, update func(entity spineapi.EntityRemoteInterface) (float64, error)) (float64, error) { c.mu.Lock() defer c.mu.Unlock() @@ -258,9 +276,17 @@ func (c *EEBus) Dim(dim bool) error { return api.ErrNotAvailable } - return eebus.Await(func(cb func(model.ResultDataType, model.MsgCounterType)) (*model.MsgCounterType, error) { + if err := eebus.Await(func(cb func(model.ResultDataType, model.MsgCounterType)) (*model.MsgCounterType, error) { return c.eg.EgLPCInterface.WriteConsumptionLimit(entity, ucapi.LoadLimit{Value: value, IsActive: dim}, cb) - }) + }); err != nil { + return err + } + + c.mu.Lock() + c.dimmed = dim + c.mu.Unlock() + + return nil } var _ api.Curtailer = (*EEBus)(nil) @@ -311,7 +337,15 @@ func (c *EEBus) SetCurtailPercent(percent int) error { } } - return eebus.Await(func(cb func(model.ResultDataType, model.MsgCounterType)) (*model.MsgCounterType, error) { + if err := eebus.Await(func(cb func(model.ResultDataType, model.MsgCounterType)) (*model.MsgCounterType, error) { return c.eg.EgLPPInterface.WriteProductionLimit(entity, ucapi.LoadLimit{Value: value, IsActive: curtail}, cb) - }) + }); err != nil { + return err + } + + c.mu.Lock() + c.curtailPercent = percent + c.mu.Unlock() + + return nil } diff --git a/meter/eebus_events.go b/meter/eebus_events.go index f67db0819..bdba484bc 100644 --- a/meter/eebus_events.go +++ b/meter/eebus_events.go @@ -79,6 +79,9 @@ func (c *EEBus) egLpcUseCaseSupportUpdate(entity spineapi.EntityRemoteInterface) // prefer the shallowest (device-level) entity if c.egLpcEntity == nil || len(entity.Address().Entity) < len(c.egLpcEntity.Address().Entity) { 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()) }) } } @@ -93,5 +96,8 @@ func (c *EEBus) egLppUseCaseSupportUpdate(entity spineapi.EntityRemoteInterface) // prefer the shallowest (device-level) entity if c.egLppEntity == nil || len(entity.Address().Entity) < len(c.egLppEntity.Address().Entity) { 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()) }) } } diff --git a/meter/eebus_lpc_lpp_test.go b/meter/eebus_lpc_lpp_test.go index 9a05dd875..72e5e2fa6 100644 --- a/meter/eebus_lpc_lpp_test.go +++ b/meter/eebus_lpc_lpp_test.go @@ -5,8 +5,11 @@ package meter import ( "testing" + "time" ucapi "github.com/enbility/eebus-go/usecases/api" + eglpc "github.com/enbility/eebus-go/usecases/eg/lpc" + eglpp "github.com/enbility/eebus-go/usecases/eg/lpp" egmocks "github.com/enbility/eebus-go/usecases/mocks" spineapi "github.com/enbility/spine-go/api" spinemocks "github.com/enbility/spine-go/mocks" @@ -42,6 +45,27 @@ func ackWrite(_ spineapi.EntityRemoteInterface, _ ucapi.LoadLimit, cb func(model cb(model.ResultDataType{}, 0) } +// captureWrite acks a mocked write and reports the limit it carried. +func captureWrite(written chan<- ucapi.LoadLimit) func(spineapi.EntityRemoteInterface, ucapi.LoadLimit, func(model.ResultDataType, model.MsgCounterType)) { + return func(entity spineapi.EntityRemoteInterface, limit ucapi.LoadLimit, cb func(model.ResultDataType, model.MsgCounterType)) { + written <- limit + ackWrite(entity, limit, cb) + } +} + +// awaitLimit waits for a limit written by the background AssertLimit retry. +func awaitLimit(t *testing.T, written <-chan ucapi.LoadLimit) ucapi.LoadLimit { + t.Helper() + + select { + case limit := <-written: + return limit + case <-time.After(time.Second): + t.Fatal("no limit written after connect") + return ucapi.LoadLimit{} + } +} + // --- LPC: Dim/Dimmed (Active Power Consumption Limit) ------------------------- // ATC_COM_PT_EGMessages_001/003 (LPC-TS-001/001-2): the EG sends an activated, @@ -246,6 +270,43 @@ func TestLPP_CurtailedPercent_NoNominal(t *testing.T) { assert.ErrorIs(t, err, api.ErrNotAvailable) } +// ATC_COM_PT_EGMessages_002 (LPC-913/LPP-913): after initial connection the EG +// states a limit to the CS - deactivated when nothing is being limited. +func TestLPC_LPP_InitialLimit(t *testing.T) { + t.Run("consumption", func(t *testing.T) { + c, lpc, _, entity := newEGMeter(t) + c.egLpcEntity = nil + + written := make(chan ucapi.LoadLimit, 1) + lpc.EXPECT().IsScenarioAvailableAtEntity(entity, eebus.LPCLimit).Return(true) + lpc.EXPECT(). + WriteConsumptionLimit(entity, mock.Anything, mock.Anything). + Run(captureWrite(written)). + Return(new(model.MsgCounterType), nil) + + c.UseCaseEvent(nil, entity, eglpc.UseCaseSupportUpdate) + + assert.Equal(t, ucapi.LoadLimit{Value: 0, IsActive: false}, awaitLimit(t, written)) + }) + + t.Run("production", func(t *testing.T) { + c, _, lpp, entity := newEGMeter(t) + c.egLppEntity = nil + c.curtailPercent = 100 + + written := make(chan ucapi.LoadLimit, 1) + lpp.EXPECT().IsScenarioAvailableAtEntity(entity, eebus.LPPLimit).Return(true) + lpp.EXPECT(). + WriteProductionLimit(entity, mock.Anything, mock.Anything). + Run(captureWrite(written)). + Return(new(model.MsgCounterType), nil) + + c.UseCaseEvent(nil, entity, eglpp.UseCaseSupportUpdate) + + assert.Equal(t, ucapi.LoadLimit{Value: 0, IsActive: false}, awaitLimit(t, written)) + }) +} + // TestLPC_LPP_NonCoverage records the Controllable-System and connection/heartbeat // abstract test cases that belong to eebus-go and the evcc HEMS/charger, not the meter. func TestLPC_LPP_NonCoverage(t *testing.T) { @@ -254,7 +315,6 @@ func TestLPC_LPP_NonCoverage(t *testing.T) { "ATC_COM_PT_CSFS_001", // failsafe values → hems/eebus + eebus-go "ATC_COM_PT_EGHeartbeat_001", // heartbeat cadence → eebus-go "ATC_COM_PT_EGConnection_001", // connection setup → eebus-go - "ATC_COM_PT_EGMessages_002", // resend-after-reboot/NACK → eebus-go } { t.Run(atc, func(t *testing.T) { t.Skip("not applicable: covered by eebus-go or the evcc HEMS/charger, not the grid meter") diff --git a/server/eebus/helper.go b/server/eebus/helper.go index b84f71d06..0c09302f2 100644 --- a/server/eebus/helper.go +++ b/server/eebus/helper.go @@ -6,9 +6,11 @@ import ( "log" "time" + "github.com/cenkalti/backoff/v4" eebusapi "github.com/enbility/eebus-go/api" "github.com/enbility/spine-go/model" "github.com/evcc-io/evcc/api" + "github.com/evcc-io/evcc/util" ) func WrapError(err error) error { @@ -45,6 +47,19 @@ func Await(write func(func(model.ResultDataType, model.MsgCounterType)) (*model. } } +// limitTimeout keeps the retries inside the 60s the spec grants the Energy Guard. +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) { + bo := backoff.NewExponentialBackOff(backoff.WithMaxElapsedTime(limitTimeout)) + if err := backoff.Retry(write, bo); err != nil { + log.DEBUG.Printf("assert limit: %v", err) + } +} + func LogEntities(log *log.Logger, actor string, uc eebusapi.UseCaseInterface) { ss := uc.RemoteEntitiesScenarios() if len(ss) > 0 {