diff --git a/hems/eebus/eebus.go b/hems/eebus/eebus.go index fdcd550dd..11965fdf9 100644 --- a/hems/eebus/eebus.go +++ b/hems/eebus/eebus.go @@ -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 +} diff --git a/hems/eebus/events.go b/hems/eebus/events.go index 346cecc0b..3dbd6bfda 100644 --- a/hems/eebus/events.go +++ b/hems/eebus/events.go @@ -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(){} diff --git a/hems/eebus/types.go b/hems/eebus/types.go index f0a6623e2..d70bd6b4a 100644 --- a/hems/eebus/types.go +++ b/hems/eebus/types.go @@ -3,7 +3,6 @@ package eebus type status int const ( - StatusUnlimited status = iota - StatusLimited + StatusNormal status = iota StatusFailsafe ) diff --git a/meter/eebus.go b/meter/eebus.go index 2455429bd..3e965b0e0 100644 --- a/meter/eebus.go +++ b/meter/eebus.go @@ -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: diff --git a/server/eebus/eebus.go b/server/eebus/eebus.go index 55317673c..c6b4066da 100644 --- a/server/eebus/eebus.go +++ b/server/eebus/eebus.go @@ -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") } diff --git a/server/eebus/test/controlbox.go b/server/eebus/test/controlbox.go new file mode 100644 index 000000000..08df10fc4 --- /dev/null +++ b/server/eebus/test/controlbox.go @@ -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 +} diff --git a/server/eebus/test/cs_test.go b/server/eebus/test/cs_test.go new file mode 100644 index 000000000..df31cde9d --- /dev/null +++ b/server/eebus/test/cs_test.go @@ -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") +} diff --git a/server/eebus/types.go b/server/eebus/types.go index 1da233e5e..43aef380d 100644 --- a/server/eebus/types.go +++ b/server/eebus/types.go @@ -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