EEBus: await control-write results across all use cases (#31350)

This commit is contained in:
andig 2026-06-30 14:14:32 +02:00 • committed by GitHub
parent 2ac2fae0a3
commit 258dcb41b5
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
5 changed files with 68 additions and 96 deletions

View file

@ -345,7 +345,9 @@ func (c *EEBus) writeCurrentLimitData(evEntity spineapi.EntityRemoteInterface, c
}
// always set overload protection limits (obligation)
if _, err := c.cem.OpEV.WriteLoadControlLimits(evEntity, limits, c.callbackResult("opEV limits")); err != nil {
if err := eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.cem.OpEV.WriteLoadControlLimits(evEntity, limits, cb)
}); err != nil {
return err
}
@ -395,27 +397,13 @@ func (c *EEBus) writeOscevLimits(evEntity spineapi.EntityRemoteInterface, curren
limits = append(limits, limit)
}
if _, err := c.cem.OscEV.WriteLoadControlLimits(evEntity, limits, c.callbackResult("oscEV limits")); err != nil {
if err := eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.cem.OscEV.WriteLoadControlLimits(evEntity, limits, cb)
}); err != nil {
c.log.DEBUG.Println("failed to write OSCEV limits:", err)
}
}
// callbackResult logs a rejected eebus write; a successful result is ignored.
func (c *EEBus) callbackResult(msg string) func(model.ResultDataType) {
return func(result model.ResultDataType) {
if result.ErrorNumber == nil || *result.ErrorNumber == 0 {
return
}
if result.Description != nil {
c.log.ERROR.Printf("%s: write rejected: %d (%s)", msg, *result.ErrorNumber, *result.Description)
return
}
c.log.ERROR.Printf("%s: write rejected: %d", msg, *result.ErrorNumber)
}
}
// MaxCurrent implements the api.Charger interface
func (c *EEBus) MaxCurrent(current int64) error {
return c.MaxCurrentMillis(float64(current))

View file

@ -3,7 +3,6 @@ package charger
import (
"context"
"errors"
"fmt"
"sync"
"time"
@ -326,13 +325,13 @@ func ohpcfControlAction(state ucapi.CompressorPowerConsumptionStateType, enable
// it aborts the process.
func (c *EEBusOHPCF) stop(entity spineapi.EntityRemoteInterface) error {
if pausable, err := c.cem.OHPCF.ConsumptionIsPausable(entity); err == nil && pausable {
return c.await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.cem.OHPCF.PausePowerConsumptionProcess(entity, cb)
})
}
if stoppable, err := c.cem.OHPCF.ConsumptionIsStoppable(entity); err == nil && stoppable {
return c.await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.cem.OHPCF.AbortPowerConsumptionProcess(entity, cb)
})
}
@ -340,34 +339,6 @@ func (c *EEBusOHPCF) stop(entity spineapi.EntityRemoteInterface) error {
return api.ErrNotAvailable
}
// ohpcfWriteTimeout bounds how long a control write waits for its result.
const ohpcfWriteTimeout = 10 * time.Second
// await runs a control write and waits for the heat pump's result, returning an
// error if the write is rejected or no result arrives within the timeout.
func (c *EEBusOHPCF) await(write func(func(model.ResultDataType)) (*model.MsgCounterType, error)) error {
res := make(chan model.ResultDataType, 1)
if _, err := write(func(r model.ResultDataType) { res <- r }); err != nil {
return err
}
select {
case r := <-res:
if r.ErrorNumber != nil && *r.ErrorNumber != 0 {
err := fmt.Errorf("write rejected: %d", *r.ErrorNumber)
if r.Description != nil {
err = fmt.Errorf("%w (%s)", err, *r.Description)
}
c.log.ERROR.Println(err)
return err
}
return nil
case <-time.After(ohpcfWriteTimeout):
return errors.New("write result timeout")
}
}
// MaxCurrent implements the api.Charger interface. OHPCF is on/off and cannot
// be modulated, so the offered current is ignored.
func (c *EEBusOHPCF) MaxCurrent(int64) error {
@ -413,7 +384,7 @@ 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 c.await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.eg.EgLPCInterface.WriteConsumptionLimit(entity, ucapi.LoadLimit{Value: 0, IsActive: dim}, cb)
})
}
@ -434,12 +405,12 @@ func (c *EEBusOHPCF) apply() error {
switch ohpcfControlAction(state, c.lastEnabled()) {
case ohpcfSchedule:
return c.await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
// 0 = start immediately (relative schedule, see SchedulePowerConsumptionProcess)
return c.cem.OHPCF.SchedulePowerConsumptionProcess(entity, 0, cb)
})
case ohpcfResume:
return c.await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.cem.OHPCF.ResumePowerConsumptionProcess(entity, cb)
})
case ohpcfStop:

View file

@ -8,6 +8,7 @@ import (
ucapi "github.com/enbility/eebus-go/usecases/api"
evcemuc "github.com/enbility/eebus-go/usecases/cem/evcem"
"github.com/enbility/eebus-go/usecases/mocks"
spineapi "github.com/enbility/spine-go/api"
spinemocks "github.com/enbility/spine-go/mocks"
"github.com/enbility/spine-go/model"
"github.com/evcc-io/evcc/api"
@ -131,6 +132,12 @@ func opevLimits3p(min, max, def float64) ([]float64, []float64, []float64, error
return []float64{min, min, min}, []float64{max, max, max}, []float64{def, def, def}, nil
}
// ackWrite makes a mocked WriteLoadControlLimits invoke its result callback with a
// success result, as the real eebus-go does, so eebus.Await completes.
func ackWrite(_ spineapi.EntityRemoteInterface, _ []ucapi.LoadLimitsPhase, resultCB func(model.ResultDataType)) {
resultCB(model.ResultDataType{})
}
func TestWriteCurrentLimitData_OpevOnly(t *testing.T) {
eebus, opev, oscev, evEntity := newTestEEBus(t)
_ = eebus
@ -140,7 +147,7 @@ func TestWriteCurrentLimitData_OpevOnly(t *testing.T) {
opev.EXPECT().CurrentLimits(evEntity).Return(opevLimits3p(6, 16, 0))
opev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
return len(limits) == 3 && limits[0].IsActive && limits[0].Value == 10
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
oscev.EXPECT().IsScenarioAvailableAtEntity(evEntity, uint(1)).Return(false)
@ -158,7 +165,7 @@ func TestWriteCurrentLimitData_OpevAndOscev(t *testing.T) {
opev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
// OPEV: active at 10A (below max of 16)
return len(limits) == 3 && limits[0].IsActive && limits[0].Value == 10
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
oscev.EXPECT().IsScenarioAvailableAtEntity(evEntity, uint(1)).Return(true)
oscev.EXPECT().LoadControlLimits(evEntity).Return([]ucapi.LoadLimitsPhase{}, nil)
@ -166,7 +173,7 @@ func TestWriteCurrentLimitData_OpevAndOscev(t *testing.T) {
oscev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
// OSCEV: active at 10A (>= min of 2, recommendation to charge)
return len(limits) == 3 && limits[0].IsActive && limits[0].Value == 10
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
err := eebus.writeCurrentLimitData(evEntity, 10)
require.NoError(t, err)
@ -182,7 +189,7 @@ func TestWriteCurrentLimitData_AtMax(t *testing.T) {
opev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
// OPEV: inactive at max (no restriction needed)
return len(limits) == 3 && !limits[0].IsActive && limits[0].Value == 16
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
oscev.EXPECT().IsScenarioAvailableAtEntity(evEntity, uint(1)).Return(true)
oscev.EXPECT().LoadControlLimits(evEntity).Return([]ucapi.LoadLimitsPhase{}, nil)
@ -190,7 +197,7 @@ func TestWriteCurrentLimitData_AtMax(t *testing.T) {
oscev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
// OSCEV: active at 16A (>= min, recommend charging)
return len(limits) == 3 && limits[0].IsActive && limits[0].Value == 16
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
err := eebus.writeCurrentLimitData(evEntity, 16)
require.NoError(t, err)
@ -206,7 +213,7 @@ func TestWriteCurrentLimitData_Disable(t *testing.T) {
opev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
// OPEV: active at 0A (hard stop)
return len(limits) == 3 && limits[0].IsActive && limits[0].Value == 0
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
oscev.EXPECT().IsScenarioAvailableAtEntity(evEntity, uint(1)).Return(true)
oscev.EXPECT().LoadControlLimits(evEntity).Return([]ucapi.LoadLimitsPhase{}, nil)
@ -214,7 +221,7 @@ func TestWriteCurrentLimitData_Disable(t *testing.T) {
oscev.EXPECT().WriteLoadControlLimits(evEntity, mock.MatchedBy(func(limits []ucapi.LoadLimitsPhase) bool {
// OSCEV: inactive at 0A (no recommendation, < min)
return len(limits) == 3 && !limits[0].IsActive && limits[0].Value == 0
}), mock.Anything).Return(nil, nil)
}), mock.Anything).Run(ackWrite).Return(nil, nil)
err := eebus.writeCurrentLimitData(evEntity, 0)
require.NoError(t, err)
@ -227,7 +234,7 @@ func TestWriteCurrentLimitData_OscevNoLimitData(t *testing.T) {
// OSCEV scenario available but no limit data (e.g. PCMP wallbox)
opev.EXPECT().IsScenarioAvailableAtEntity(evEntity, uint(1)).Return(true)
opev.EXPECT().CurrentLimits(evEntity).Return(opevLimits3p(6, 16, 0))
opev.EXPECT().WriteLoadControlLimits(evEntity, mock.Anything, mock.Anything).Return(nil, nil)
opev.EXPECT().WriteLoadControlLimits(evEntity, mock.Anything, mock.Anything).Run(ackWrite).Return(nil, nil)
oscev.EXPECT().IsScenarioAvailableAtEntity(evEntity, uint(1)).Return(true)
oscev.EXPECT().LoadControlLimits(evEntity).Return(nil, errors.New("data not available"))

View file

@ -4,7 +4,6 @@ import (
"context"
"errors"
"fmt"
"strings"
"sync"
"time"
@ -250,18 +249,16 @@ func (c *EEBus) Dim(dim bool) error {
}
c.mu.Lock()
defer c.mu.Unlock()
entity := c.egLpcEntity
c.mu.Unlock()
if c.egLpcEntity == nil || !c.eg.EgLPCInterface.IsScenarioAvailableAtEntity(c.egLpcEntity, eebus.LPCLimit) {
if entity == nil || !c.eg.EgLPCInterface.IsScenarioAvailableAtEntity(entity, eebus.LPCLimit) {
return api.ErrNotAvailable
}
_, err := c.eg.EgLPCInterface.WriteConsumptionLimit(c.egLpcEntity, ucapi.LoadLimit{
Value: value,
IsActive: dim,
}, c.callbackResult("consumption limit"))
return err
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.eg.EgLPCInterface.WriteConsumptionLimit(entity, ucapi.LoadLimit{Value: value, IsActive: dim}, cb)
})
}
var _ api.Curtailer = (*EEBus)(nil)
@ -294,35 +291,14 @@ func (c *EEBus) Curtail(curtail bool) error {
}
c.mu.Lock()
defer c.mu.Unlock()
entity := c.egLppEntity
c.mu.Unlock()
if c.egLppEntity == nil || !c.eg.EgLPPInterface.IsScenarioAvailableAtEntity(c.egLppEntity, eebus.LPPLimit) {
if entity == nil || !c.eg.EgLPPInterface.IsScenarioAvailableAtEntity(entity, eebus.LPPLimit) {
return api.ErrNotAvailable
}
_, err := c.eg.EgLPPInterface.WriteProductionLimit(c.egLppEntity, ucapi.LoadLimit{
Value: value,
IsActive: curtail,
}, c.callbackResult("production limit"))
return err
}
func (c *EEBus) callbackResult(typ string) func(result model.ResultDataType) {
return func(result model.ResultDataType) {
sb := new(strings.Builder)
if result.ErrorNumber != nil {
fmt.Fprint(sb, *result.ErrorNumber)
}
if result.Description != nil {
if sb.Len() > 0 {
fmt.Print(sb, ":")
}
fmt.Print(sb, *result.Description)
}
if sb.Len() > 0 {
c.log.ERROR.Printf("%s: %s", typ, sb.String())
}
}
return eebus.Await(func(cb func(model.ResultDataType)) (*model.MsgCounterType, error) {
return c.eg.EgLPPInterface.WriteProductionLimit(entity, ucapi.LoadLimit{Value: value, IsActive: curtail}, cb)
})
}

View file

@ -2,9 +2,12 @@ package eebus
import (
"errors"
"fmt"
"log"
"time"
eebusapi "github.com/enbility/eebus-go/api"
"github.com/enbility/spine-go/model"
"github.com/evcc-io/evcc/api"
)
@ -15,6 +18,33 @@ func WrapError(err error) error {
return err
}
// WriteTimeout bounds how long an awaited eebus write waits for its result.
const WriteTimeout = 10 * time.Second
// Await runs a control write and waits for the remote device's result, returning
// an error if the write is rejected or no result arrives within WriteTimeout.
func Await(write func(func(model.ResultDataType)) (*model.MsgCounterType, error)) error {
res := make(chan model.ResultDataType, 1)
if _, err := write(func(r model.ResultDataType) { res <- r }); err != nil {
return err
}
select {
case r := <-res:
if r.ErrorNumber != nil && *r.ErrorNumber != 0 {
err := fmt.Errorf("write rejected: %d", *r.ErrorNumber)
if r.Description != nil {
err = fmt.Errorf("%w (%s)", err, *r.Description)
}
return err
}
return nil
case <-time.After(WriteTimeout):
return errors.New("write result timeout")
}
}
func LogEntities(log *log.Logger, actor string, uc eebusapi.UseCaseInterface) {
ss := uc.RemoteEntitiesScenarios()
if len(ss) > 0 {