EEBus HEMS: add controllable system limitation of power production (experimental) (#26226)

This commit is contained in:
andig 2025-12-29 19:07:24 +01:00 • committed by GitHub
parent 155029b258
commit 5f47a55237
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 517 additions and 182 deletions

View file

@ -24,29 +24,37 @@ type EEBus struct {
*eebus.Connector
cs *eebus.ControllableSystem
ma *eebus.MonitoringAppliance
eg *eebus.EnergyGuard
root api.Circuit
passthrough func(bool) error
smartgridID uint
status status
statusUpdated time.Time
consumptionLimit *ucapi.LoadLimit // LPC-041
failsafeLimit float64
failsafeDuration time.Duration
smartgridConsumptionId uint
consumptionLimit ucapi.LoadLimit // LPC-041
consumptionLimitActivated time.Time
failsafeConsumptionLimit float64
smartgridProductionId uint
productionLimit ucapi.LoadLimit
productionLimitActivated time.Time
failsafeProductionLimit float64
heartbeat *util.Value[struct{}]
interval time.Duration
}
type Limits struct {
ContractualConsumptionNominalMax float64
ConsumptionLimit float64
FailsafeConsumptionActivePowerLimit float64
FailsafeDurationMinimum time.Duration
ProductionNominalMax float64
FailsafeProductionActivePowerLimit float64
FailsafeDurationMinimum time.Duration
}
// NewFromConfig creates an EEBus HEMS from generic config
@ -59,9 +67,12 @@ func NewFromConfig(ctx context.Context, other map[string]any, site site.API) (*E
}{
Limits: Limits{
ContractualConsumptionNominalMax: 24800,
ConsumptionLimit: 0,
FailsafeConsumptionActivePowerLimit: 4200,
FailsafeDurationMinimum: 2 * time.Hour,
ProductionNominalMax: 0,
FailsafeProductionActivePowerLimit: 0,
FailsafeDurationMinimum: 2 * time.Hour,
},
Interval: 10 * time.Second,
}
@ -82,21 +93,21 @@ func NewFromConfig(ctx context.Context, other map[string]any, site site.API) (*E
}
// register LPC circuit if not already registered
lpc, err := shared.GetOrCreateCircuit("lpc", "eebus")
gridcontrol, err := shared.GetOrCreateCircuit("gridcontrol", "eebus")
if err != nil {
return nil, err
}
// wrap old root with new pc parent
if err := root.Wrap(lpc); err != nil {
// wrap old root with new grid control parent
if err := root.Wrap(gridcontrol); err != nil {
return nil, err
}
site.SetCircuit(lpc)
site.SetCircuit(gridcontrol)
return NewEEBus(ctx, cc.Ski, cc.Limits, passthroughS, lpc, cc.Interval)
return NewEEBus(ctx, cc.Ski, cc.Limits, passthroughS, gridcontrol, cc.Interval)
}
// NewEEBus creates EEBus charger
// NewEEBus creates EEBus HEMS
func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(bool) error, root api.Circuit, interval time.Duration) (*EEBus, error) {
if eebus.Instance == nil {
return nil, errors.New("eebus not configured")
@ -107,19 +118,13 @@ func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(b
root: root,
passthrough: passthrough,
cs: eebus.Instance.ControllableSystem(),
ma: eebus.Instance.MonitoringAppliance(),
eg: eebus.Instance.EnergyGuard(),
Connector: eebus.NewConnector(),
heartbeat: util.NewValue[struct{}](2 * time.Minute), // LPC-031
interval: interval,
consumptionLimit: &ucapi.LoadLimit{
Value: limits.ConsumptionLimit,
IsChangeable: true,
},
failsafeLimit: limits.FailsafeConsumptionActivePowerLimit,
failsafeDuration: limits.FailsafeDurationMinimum,
failsafeDuration: limits.FailsafeDurationMinimum,
failsafeConsumptionLimit: limits.FailsafeConsumptionActivePowerLimit,
failsafeProductionLimit: limits.FailsafeProductionActivePowerLimit,
}
// simulate a received heartbeat
@ -137,37 +142,38 @@ func NewEEBus(ctx context.Context, ski string, limits Limits, passthrough func(b
// controllable system
for _, s := range c.cs.CsLPCInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Println("CS LPC RemoteEntitiesScenarios:", s.Scenarios)
c.log.DEBUG.Printf("ski %s CS LPC scenarios: %v", s.Entity.Device().Ski(), s.Scenarios)
}
for _, s := range c.cs.CsLPPInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Println("CS LPP RemoteEntitiesScenarios:", s.Scenarios)
}
// monitoring appliance
for _, s := range c.ma.MaMPCInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Println("MA MPC RemoteEntitiesScenarios:", s.Scenarios)
}
for _, s := range c.ma.MaMGCPInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Println("MA MGCP RemoteEntitiesScenarios:", s.Scenarios)
}
// energy guard
for _, s := range c.eg.EgLPCInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Println("EG LPC RemoteEntitiesScenarios:", s.Scenarios)
c.log.DEBUG.Printf("ski %s CS LPP scenarios: %v", s.Entity.Device().Ski(), s.Scenarios)
}
// set initial values
if err := c.cs.CsLPCInterface.SetConsumptionNominalMax(limits.ContractualConsumptionNominalMax); err != nil {
c.log.ERROR.Println("CS LPC SetConsumptionNominalMax:", err)
}
if err := c.cs.CsLPCInterface.SetConsumptionLimit(*c.consumptionLimit); err != nil {
c.log.ERROR.Println("CS LPC SetConsumptionLimit:", err)
if c.failsafeConsumptionLimit > 0 {
if err := c.cs.CsLPCInterface.SetFailsafeConsumptionActivePowerLimit(c.failsafeConsumptionLimit, true); err != nil {
c.log.ERROR.Println("CS LPC SetFailsafeConsumptionActivePowerLimit:", err)
}
}
if err := c.cs.CsLPCInterface.SetFailsafeConsumptionActivePowerLimit(c.failsafeLimit, true); err != nil {
c.log.ERROR.Println("CS LPC SetFailsafeConsumptionActivePowerLimit:", err)
if err := c.cs.CsLPPInterface.SetProductionNominalMax(limits.ProductionNominalMax); err != nil {
c.log.ERROR.Println("CS LPP SetProductionNominalMax:", err)
}
if err := c.cs.CsLPCInterface.SetFailsafeDurationMinimum(c.failsafeDuration, true); err != nil {
c.log.ERROR.Println("CS LPC SetFailsafeDurationMinimum:", err)
if c.failsafeProductionLimit > 0 {
if err := c.cs.CsLPPInterface.SetFailsafeProductionActivePowerLimit(c.failsafeProductionLimit, true); err != nil {
c.log.ERROR.Println("CS LPP SetFailsafeProductionActivePowerLimit:", err)
}
}
if c.failsafeDuration > 0 {
if err := c.cs.CsLPCInterface.SetFailsafeDurationMinimum(c.failsafeDuration, true); err != nil {
c.log.ERROR.Println("CS LPC SetFailsafeDurationMinimum:", err)
}
if err := c.cs.CsLPPInterface.SetFailsafeDurationMinimum(c.failsafeDuration, true); err != nil {
c.log.ERROR.Println("CS LPP SetFailsafeDurationMinimum:", err)
}
}
return c, nil
@ -185,7 +191,6 @@ func (c *EEBus) Run() {
}
}
// TODO check state machine against spec
func (c *EEBus) run() error {
c.mux.Lock()
defer c.mux.Unlock()
@ -197,100 +202,124 @@ func (c *EEBus) run() error {
if heartbeatErr != nil && c.status != StatusFailsafe {
// LPC-914/2
c.log.WARN.Println("missing heartbeat- entering failsafe mode")
c.setStatusAndLimit(StatusFailsafe, c.failsafeLimit)
c.setStatusAndLimit(StatusFailsafe, c.failsafeConsumptionLimit, c.failsafeProductionLimit)
return nil
}
// TODO
// status init
// status Unlimited/controlled
// status Unlimited/autonomous
switch c.status {
case StatusUnlimited:
// LPC-914/1
if c.consumptionLimit != nil && c.consumptionLimit.IsActive {
c.log.WARN.Println("active consumption limit")
c.setStatusAndLimit(StatusLimited, c.consumptionLimit.Value)
}
case StatusLimited:
// limit updated?
if !c.consumptionLimit.IsActive {
c.log.WARN.Println("inactive consumption limit")
c.setStatusAndLimit(StatusUnlimited, 0)
break
}
c.setLimit(c.consumptionLimit.Value)
// LPC-914/1
if d := c.consumptionLimit.Duration; d > 0 && time.Since(c.statusUpdated) > d {
c.consumptionLimit = nil
c.log.DEBUG.Println("limit duration exceeded- return to normal")
c.setStatusAndLimit(StatusUnlimited, 0)
}
case StatusFailsafe:
if c.status == StatusFailsafe {
// LPC-914/2
if d := c.failsafeDuration; heartbeatErr == nil || time.Since(c.statusUpdated) > d {
c.log.DEBUG.Println("heartbeat returned and failsafe duration exceeded- return to normal")
c.setStatusAndLimit(StatusUnlimited, 0)
if heartbeatErr != nil || time.Since(c.statusUpdated) <= c.failsafeDuration {
return nil
}
c.log.DEBUG.Println("heartbeat returned or failsafe duration exceeded- leaving failsafe mode")
c.setStatusAndLimit(StatusNormal, 0, 0)
}
// LPC-914/1
if c.consumptionLimitActivated.IsZero() {
if c.consumptionLimit.IsActive {
c.log.WARN.Println("activating consumption limit")
c.setConsumptionLimit(c.consumptionLimit.Value)
}
} else {
if time.Since(c.consumptionLimitActivated) > c.consumptionLimit.Duration {
c.log.DEBUG.Println("consumption limit duration exceeded")
c.setConsumptionLimit(0)
}
}
// LPP
if c.productionLimitActivated.IsZero() {
if c.productionLimit.IsActive {
c.log.WARN.Println("activating production limit")
c.setProductionLimit(c.productionLimit.Value)
}
} else {
if time.Since(c.productionLimitActivated) > c.productionLimit.Duration {
c.log.DEBUG.Println("production limit duration exceeded")
c.setProductionLimit(0)
}
}
return nil
}
func (c *EEBus) setStatusAndLimit(status status, limit float64) {
func (c *EEBus) setStatusAndLimit(status status, consumption, production float64) {
c.status = status
c.statusUpdated = time.Now()
c.setLimit(limit)
if err := c.updateSession(limit); err != nil {
c.log.ERROR.Printf("smartgrid session: %v", err)
}
c.setConsumptionLimit(consumption)
c.setProductionLimit(production)
}
// TODO keep in sync across HEMS implementations
func (c *EEBus) updateSession(limit float64) error {
// start session
if limit > 0 && c.smartgridID == 0 {
var power *float64
if p := c.root.GetChargePower(); p > 0 {
power = lo.ToPtr(p)
}
func (c *EEBus) setConsumptionLimit(limit float64) {
active := limit > 0
sid, err := smartgrid.StartManage(smartgrid.Dim, power, limit)
if err != nil {
return err
}
c.smartgridID = sid
if active {
c.consumptionLimitActivated = time.Now()
} else {
c.consumptionLimitActivated = time.Time{}
}
// stop session
if limit == 0 && c.smartgridID != 0 {
if err := smartgrid.StopManage(c.smartgridID); err != nil {
return err
}
c.smartgridID = 0
}
return nil
}
func (c *EEBus) setLimit(limit float64) {
c.root.Dim(limit > 0)
c.root.Dim(active)
c.root.SetMaxPower(limit)
if err := c.updateSession(&c.smartgridConsumptionId, smartgrid.Dim, limit); err != nil {
c.log.ERROR.Printf("smartgrid dim session: %v", err)
}
if c.passthrough != nil {
if err := c.passthrough(limit > 0); err != nil {
c.log.ERROR.Printf("passthrough failed: %v", err)
}
}
}
func (c *EEBus) setProductionLimit(limit float64) {
active := limit > 0
if active {
c.productionLimitActivated = time.Now()
} else {
c.productionLimitActivated = time.Time{}
}
c.root.Curtail(active)
// TODO make ProductionNominalMax configurable (Site kWp)
// c.root.SetMaxProduction(limit)
if err := c.updateSession(&c.smartgridProductionId, smartgrid.Curtail, limit); err != nil {
c.log.ERROR.Printf("smartgrid curtail session: %v", err)
}
}
// TODO keep in sync across HEMS implementations
func (c *EEBus) updateSession(id *uint, typ smartgrid.Type, limit float64) error {
// start session
if limit > 0 && *id == 0 {
var power *float64
if p := c.root.GetChargePower(); p > 0 {
power = lo.ToPtr(p)
}
sid, err := smartgrid.StartManage(typ, power, limit)
if err != nil {
return err
}
*id = sid
}
// stop session
if limit == 0 && *id != 0 {
if err := smartgrid.StopManage(*id); err != nil {
return err
}
*id = 0
}
return nil
}

View file

@ -1,8 +1,11 @@
package eebus
import (
"time"
eebusapi "github.com/enbility/eebus-go/api"
"github.com/enbility/eebus-go/usecases/cs/lpc"
"github.com/enbility/eebus-go/usecases/cs/lpp"
spineapi "github.com/enbility/spine-go/api"
"github.com/evcc-io/evcc/server/eebus"
)
@ -18,7 +21,7 @@ func (c *EEBus) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spineapi.E
//
// Use Case LPC, Scenario 1
case lpc.DataUpdateLimit:
c.dataUpdateLimit()
c.updateConsumptionLimit()
// An incoming load control obligation limit needs to be approved or denied
//
@ -27,7 +30,7 @@ func (c *EEBus) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spineapi.E
//
// Use Case LPC, Scenario 1
case lpc.WriteApprovalRequired:
c.writeApprovalRequired()
c.consumptionWriteApprovalRequired()
// Failsafe limit for the consumed active (real) power of the
// Controllable System data update received
@ -36,7 +39,7 @@ func (c *EEBus) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spineapi.E
//
// Use Case LPC, Scenario 2
case lpc.DataUpdateFailsafeConsumptionActivePowerLimit:
c.dataUpdateFailsafeConsumptionActivePowerLimit()
c.updateFailsafeConsumptionActivePowerLimit()
// Minimum time the Controllable System remains in "failsafe state" unless conditions
// specified in this Use Case permit leaving the "failsafe state" data update received
@ -45,60 +48,60 @@ func (c *EEBus) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spineapi.E
//
// Use Case LPC, Scenario 2
case lpc.DataUpdateFailsafeDurationMinimum:
c.dataUpdateFailsafeDurationMinimum()
c.updateFailsafeConsumptionDurationMinimum()
// Indicates a notify heartbeat event the application should care of.
// E.g. going into or out of the Failsafe state
//
// Use Case LPC, Scenario 3
case lpc.DataUpdateHeartbeat:
c.dataUpdateHeartbeat()
c.updateHeartbeat()
// // Load control obligation limit data update received
// //
// // Use `ProductionLimit` to get the current data
// //
// // Use Case LPC, Scenario 1
// case lpp.DataUpdateLimit:
// c.dataUpdateLimit()
// Load control obligation limit data update received
//
// Use `ProductionLimit` to get the current data
//
// Use Case LPC, Scenario 1
case lpp.DataUpdateLimit:
c.updateProductionLimit()
// // An incoming load control obligation limit needs to be approved or denied
// //
// // Use `PendingProductionLimits` to get the currently pending write approval requests
// // and invoke `ApproveOrDenyProductionLimit` for each
// //
// // Use Case LPC, Scenario 1
// case lpp.WriteApprovalRequired:
// c.writeApprovalRequired()
// An incoming load control obligation limit needs to be approved or denied
//
// Use `PendingProductionLimits` to get the currently pending write approval requests
// and invoke `ApproveOrDenyProductionLimit` for each
//
// Use Case LPP, Scenario 1
case lpp.WriteApprovalRequired:
c.productionWriteApprovalRequired()
// // Failsafe limit for the produced active (real) power of the
// // Controllable System data update received
// //
// // Use `FailsafeProductionActivePowerLimit` to get the current data
// //
// // Use Case LPC, Scenario 2
// case lpp.DataUpdateFailsafeProductionActivePowerLimit:
// c.dataUpdateFailsafeProductionActivePowerLimit()
// Failsafe limit for the produced active (real) power of the
// Controllable System data update received
//
// Use `FailsafeProductionActivePowerLimit` to get the current data
//
// Use Case LPP, Scenario 2
case lpp.DataUpdateFailsafeProductionActivePowerLimit:
c.updateFailsafeProductionActivePowerLimit()
// // Minimum time the Controllable System remains in "failsafe state" unless conditions
// // specified in this Use Case permit leaving the "failsafe state" data update received
// //
// // Use `FailsafeDurationMinimum` to get the current data
// //
// // Use Case LPC, Scenario 2
// case lpp.DataUpdateFailsafeDurationMinimum:
// c.dataUpdateFailsafeDurationMinimum()
// Minimum time the Controllable System remains in "failsafe state" unless conditions
// specified in this Use Case permit leaving the "failsafe state" data update received
//
// Use `FailsafeDurationMinimum` to get the current data
//
// Use Case LPP, Scenario 2
case lpp.DataUpdateFailsafeDurationMinimum:
c.updateFailsafeProductionDurationMinimum()
// // Indicates a notify heartbeat event the application should care of.
// // E.g. going into or out of the Failsafe state
// //
// // Use Case LPP, Scenario 3
// case lpp.DataUpdateHeartbeat:
// c.dataUpdateHeartbeat()
// Indicates a notify heartbeat event the application should care of.
// E.g. going into or out of the Failsafe state
//
// Use Case LPP, Scenario 3
case lpp.DataUpdateHeartbeat:
c.updateHeartbeat()
}
}
func (c *EEBus) dataUpdateLimit() {
func (c *EEBus) updateConsumptionLimit() {
limit, err := c.cs.CsLPCInterface.ConsumptionLimit()
if err != nil {
c.log.ERROR.Println("CS LPC ConsumptionLimit:", err)
@ -108,10 +111,25 @@ func (c *EEBus) dataUpdateLimit() {
c.mux.Lock()
defer c.mux.Unlock()
c.consumptionLimit = &limit
c.consumptionLimit = limit
c.statusUpdated = time.Now()
}
func (c *EEBus) writeApprovalRequired() {
func (c *EEBus) updateProductionLimit() {
limit, err := c.cs.CsLPPInterface.ProductionLimit()
if err != nil {
c.log.ERROR.Println("CS LPP ProductionLimit:", err)
return
}
c.mux.Lock()
defer c.mux.Unlock()
c.productionLimit = limit
c.statusUpdated = time.Now()
}
func (c *EEBus) consumptionWriteApprovalRequired() {
for msg, limit := range c.cs.CsLPCInterface.PendingConsumptionLimits() {
c.log.DEBUG.Println("CS LPC PendingConsumptionLimit:", msg, limit)
if limit.Value < 0 {
@ -122,12 +140,27 @@ func (c *EEBus) writeApprovalRequired() {
c.cs.CsLPCInterface.ApproveOrDenyConsumptionLimit(msg, true, "")
c.mux.Lock()
c.consumptionLimit = &limit
c.consumptionLimit = limit
c.mux.Unlock()
}
}
func (c *EEBus) dataUpdateFailsafeConsumptionActivePowerLimit() {
func (c *EEBus) productionWriteApprovalRequired() {
for msg, limit := range c.cs.CsLPPInterface.PendingProductionLimits() {
c.log.DEBUG.Println("CS LPP PendingProductionLimit:", msg, limit)
if limit.Value > 0 {
c.cs.CsLPPInterface.ApproveOrDenyProductionLimit(msg, false, "positive limit")
continue
}
c.cs.CsLPPInterface.ApproveOrDenyProductionLimit(msg, true, "")
c.mux.Lock()
c.productionLimit = limit
c.mux.Unlock()
}
}
func (c *EEBus) updateFailsafeConsumptionActivePowerLimit() {
limit, _, err := c.cs.CsLPCInterface.FailsafeConsumptionActivePowerLimit()
if err != nil {
c.log.ERROR.Println("CS LPC FailsafeConsumptionActivePowerLimit:", err)
@ -137,10 +170,23 @@ func (c *EEBus) dataUpdateFailsafeConsumptionActivePowerLimit() {
c.mux.Lock()
defer c.mux.Unlock()
c.failsafeLimit = limit
c.failsafeConsumptionLimit = limit
}
func (c *EEBus) dataUpdateFailsafeDurationMinimum() {
func (c *EEBus) updateFailsafeProductionActivePowerLimit() {
limit, _, err := c.cs.CsLPPInterface.FailsafeProductionActivePowerLimit()
if err != nil {
c.log.ERROR.Println("CS LPP FailsafeProductionActivePowerLimit:", err)
return
}
c.mux.Lock()
defer c.mux.Unlock()
c.failsafeProductionLimit = limit
}
func (c *EEBus) updateFailsafeConsumptionDurationMinimum() {
duration, _, err := c.cs.CsLPCInterface.FailsafeDurationMinimum()
if err != nil {
c.log.ERROR.Println("CS LPC FailsafeDurationMinimum:", err)
@ -153,15 +199,22 @@ func (c *EEBus) dataUpdateFailsafeDurationMinimum() {
c.failsafeDuration = duration
}
func (c *EEBus) dataUpdateHeartbeat() {
func (c *EEBus) updateFailsafeProductionDurationMinimum() {
duration, _, err := c.cs.CsLPPInterface.FailsafeDurationMinimum()
if err != nil {
c.log.ERROR.Println("CS LPP FailsafeDurationMinimum:", err)
return
}
c.mux.Lock()
defer c.mux.Unlock()
c.failsafeDuration = duration
}
func (c *EEBus) updateHeartbeat() {
c.mux.Lock()
defer c.mux.Unlock()
c.heartbeat.Set(struct{}{})
}
// func (c *EEBus)dataUpdateLimit(){}
// func (c *EEBus)writeApprovalRequired(){}
// func (c *EEBus)dataUpdateFailsafeProductionActivePowerLimit(){}
// func (c *EEBus)dataUpdateFailsafeDurationMinimum(){}
// func (c *EEBus)dataUpdateHeartbeat(){}

View file

@ -3,7 +3,6 @@ package eebus
type status int
const (
StatusUnlimited status = iota
StatusLimited
StatusNormal status = iota
StatusFailsafe
)

View file

@ -113,6 +113,19 @@ func NewEEBus(ctx context.Context, ski, ip string, usage *templates.Usage, timeo
return nil, err
}
// monitoring appliance
for _, s := range c.ma.MaMPCInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Printf("ski %s MA MPC scenarios: %v", s.Entity.Device().Ski(), s.Scenarios)
}
for _, s := range c.ma.MaMGCPInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Printf("ski %s MA MGCP scenarios: %v", s.Entity.Device().Ski(), s.Scenarios)
}
// energy guard
for _, s := range c.eg.EgLPCInterface.RemoteEntitiesScenarios() {
c.log.DEBUG.Printf("ski %s EG LPC scenarios: %v", s.Entity.Device().Ski(), s.Scenarios)
}
return c, nil
}
@ -120,8 +133,6 @@ var _ eebus.Device = (*EEBus)(nil)
// UseCaseEvent implements the eebus.Device interface
func (c *EEBus) UseCaseEvent(_ spineapi.DeviceRemoteInterface, entity spineapi.EntityRemoteInterface, event eebusapi.EventType) {
c.log.TRACE.Printf("recv: %s", event)
switch event {
// Monitoring Appliance
case mpc.DataUpdatePower, mgcp.DataUpdatePower:

View file

@ -77,7 +77,7 @@ type EEBus struct {
mux sync.Mutex
log *util.Logger
ski string
Ski string
clients map[string][]Device
}
@ -143,7 +143,7 @@ func NewServer(other Config) (*EEBus, error) {
c := &EEBus{
log: log,
ski: ski,
Ski: ski,
clients: make(map[string][]Device),
}
@ -201,7 +201,7 @@ func (c *EEBus) RegisterDevice(ski, ip string, device Device) error {
ski = shiputil.NormalizeSKI(ski)
c.log.TRACE.Printf("registering ski: %s", ski)
if ski == c.ski {
if ski == c.Ski {
return errors.New("device ski can not be identical to host ski")
}

View file

@ -0,0 +1,168 @@
package eebus
import (
"context"
"fmt"
"sync"
"time"
"github.com/enbility/eebus-go/api"
"github.com/enbility/eebus-go/service"
ucapi "github.com/enbility/eebus-go/usecases/api"
"github.com/enbility/eebus-go/usecases/eg/lpc"
"github.com/enbility/eebus-go/usecases/eg/lpp"
shipapi "github.com/enbility/ship-go/api"
"github.com/enbility/ship-go/cert"
spineapi "github.com/enbility/spine-go/api"
"github.com/enbility/spine-go/model"
server "github.com/evcc-io/evcc/server/eebus"
)
type controlbox struct {
mu sync.Mutex
ski string
myService *service.Service
uclpc ucapi.EgLPCInterface
uclpp ucapi.EgLPPInterface
remoteEntities map[api.EventType][]spineapi.EntityRemoteInterface
remoteEventC chan<- api.EventType
isConnected bool
}
func createControlbox(ctx context.Context, remoteSki string, port int) (*controlbox, error) {
certificate, err := cert.CreateCertificate("Demo", "Demo", "DE", "Demo-Unit-01")
if err != nil {
return nil, err
}
ski, err := server.SkiFromCert(certificate)
if err != nil {
return nil, err
}
h := &controlbox{
ski: ski,
}
configuration, err := api.NewConfiguration(
"Demo", "Demo", "ControlBox", "123456789",
// []shipapi.DeviceCategoryType{shipapi.DeviceCategoryTypeGridConnectionHub},
model.DeviceTypeTypeElectricitySupplySystem,
[]model.EntityTypeType{model.EntityTypeTypeGridGuard},
port, certificate, time.Second*60)
if err != nil {
return nil, err
}
configuration.SetAlternateIdentifier("Demo-ControlBox-123456789")
h.myService = service.NewService(configuration, h)
// h.myService.SetLogging(h)
if err = h.myService.Setup(); err != nil {
return nil, err
}
localEntity := h.myService.LocalDevice().EntityForType(model.EntityTypeTypeGridGuard)
h.uclpc = lpc.NewLPC(localEntity, h.OnLPCEvent)
h.myService.AddUseCase(h.uclpc)
h.uclpp = lpp.NewLPP(localEntity, h.OnLPPEvent)
h.myService.AddUseCase(h.uclpp)
h.myService.RegisterRemoteSKI(remoteSki)
h.myService.Start()
go func() {
<-ctx.Done()
h.myService.Shutdown()
}()
return h, nil
}
func (h *controlbox) remoteEntity(event api.EventType) []spineapi.EntityRemoteInterface {
h.mu.Lock()
defer h.mu.Unlock()
return h.remoteEntities[event]
}
func (h *controlbox) registerRemoteEntity(entity spineapi.EntityRemoteInterface, event api.EventType) {
h.mu.Lock()
defer h.mu.Unlock()
defer func() {
if h.remoteEventC != nil {
h.remoteEventC <- event
}
}()
if h.remoteEntities == nil {
h.remoteEntities = make(map[api.EventType][]spineapi.EntityRemoteInterface)
}
h.remoteEntities[event] = append(h.remoteEntities[event], entity)
}
// LPC
func (h *controlbox) OnLPCEvent(ski string, device spineapi.DeviceRemoteInterface, entity spineapi.EntityRemoteInterface, event api.EventType) {
if !h.isConnected {
return
}
switch event {
case lpc.UseCaseSupportUpdate:
h.registerRemoteEntity(entity, event)
// case lpc.DataUpdateLimit:
// if currentLimit, err := h.uclpc.ConsumptionLimit(entity); err == nil {
// fmt.Println("New consumption limit received", currentLimit.Value, "W")
// }
default:
fmt.Println("lpc:", event)
}
}
// LPP
func (h *controlbox) OnLPPEvent(ski string, device spineapi.DeviceRemoteInterface, entity spineapi.EntityRemoteInterface, event api.EventType) {
if !h.isConnected {
return
}
switch event {
case lpp.UseCaseSupportUpdate:
h.registerRemoteEntity(entity, event)
// case lpp.DataUpdateLimit:
// if currentLimit, err := h.uclpp.ConsumptionLimit(entity); err == nil {
// fmt.Println("New consumption limit received", currentLimit.Value, "W")
// }
default:
fmt.Println("lpp:", event)
}
}
// EEBUSServiceHandler
func (h *controlbox) RemoteSKIConnected(service api.ServiceInterface, ski string) {
h.isConnected = true
}
func (h *controlbox) RemoteSKIDisconnected(service api.ServiceInterface, ski string) {
h.isConnected = false
}
func (h *controlbox) VisibleRemoteServicesUpdated(service api.ServiceInterface, entries []shipapi.RemoteService) {
}
func (h *controlbox) ServiceShipIDUpdate(ski string, shipdID string) {
}
func (h *controlbox) ServicePairingDetailUpdate(ski string, detail *shipapi.ConnectionStateDetail) {
}
func (h *controlbox) AllowWaitingForTrust(ski string) bool {
return true
}

View file

@ -0,0 +1,73 @@
package eebus
import (
"testing"
"time"
"github.com/enbility/eebus-go/api"
ucapi "github.com/enbility/eebus-go/usecases/api"
"github.com/enbility/eebus-go/usecases/eg/lpc"
"github.com/enbility/ship-go/cert"
"github.com/enbility/spine-go/model"
"github.com/evcc-io/evcc/core/circuit"
"github.com/evcc-io/evcc/hems/eebus"
hems "github.com/evcc-io/evcc/hems/eebus"
server "github.com/evcc-io/evcc/server/eebus"
"github.com/evcc-io/evcc/util"
"github.com/stretchr/testify/require"
)
const remotePort = 9001
func TestEEBus(t *testing.T) {
t.Skip()
util.LogLevel("error", map[string]string{"eebus": "trace"})
certificate, err := cert.CreateCertificate("Demo", "Demo", "DE", "Demo-Server-01")
require.NoError(t, err, "certificate")
public, private, err := server.GetX509KeyPair(certificate)
require.NoError(t, err, "decode certificate")
srv, err := server.NewServer(server.Config{
Certificate: server.Certificate{
Public: public,
Private: private,
},
})
require.NoError(t, err, "server")
require.NotEmpty(t, srv.Ski, "server ski")
server.Instance = srv
go srv.Run()
box, err := createControlbox(t.Context(), server.Instance.Ski, remotePort)
require.NoError(t, err, "controlbox")
eventC := make(chan api.EventType, 1)
box.remoteEventC = eventC
gridcontrol, err := circuit.New(util.NewLogger("gridcontrol"), "gridcontrol", 0, 0, nil, time.Minute)
require.NoError(t, err)
hems, err := hems.NewEEBus(t.Context(), box.ski, eebus.Limits{}, nil, gridcontrol, time.Second)
require.NoError(t, err, "hems")
go hems.Run()
<-eventC
t.Log(box.remoteEntities)
srvEntity := box.remoteEntity(lpc.UseCaseSupportUpdate)[0]
_, err = box.uclpc.WriteConsumptionLimit(srvEntity, ucapi.LoadLimit{
IsActive: true,
Value: 1,
}, func(result model.ResultDataType) {
t.Logf("lpc result: %v", result)
})
// TODO no error
require.NoError(t, err, "consumption limit")
}

View file

@ -10,13 +10,15 @@ const (
// used as common name in cert generation
var DeviceCode = util.Getenv("EEBUS_DEVICE_CODE", "EVCC_HEMS_01")
type Certificate struct {
Public, Private string
}
type Config struct {
URI string
ShipID string
Interfaces []string
Certificate struct {
Public, Private string
}
Certificate Certificate
}
// Configured returns true if the EEbus server is configured