Ocpp: fix various connection issues (#6918)

This commit is contained in:
andig 2023-03-18 13:47:26 +01:00 • committed by GitHub
parent 41d1c29147
commit 52d63a2eba
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 291 additions and 82 deletions

View file

@ -18,8 +18,6 @@ import (
"github.com/samber/lo"
)
const statusTimeout = 30 * time.Second
// OCPP charger implementation
type OCPP struct {
log *util.Logger
@ -51,9 +49,6 @@ func NewOCPPFromConfig(other map[string]interface{}) (api.Charger, error) {
Timeout time.Duration
BootNotification *bool
GetConfiguration *bool
Meter interface{} // TODO deprecated
Quirks interface{} // TODO deprecated
InitialReset interface{} // TODO deprecated
}{
Connector: 1,
IdTag: defaultIdTag,
@ -65,15 +60,6 @@ func NewOCPPFromConfig(other map[string]interface{}) (api.Charger, error) {
return nil, err
}
// switch cc.InitialReset {
// case
// "",
// core.ResetTypeSoft,
// core.ResetTypeHard:
// default:
// return nil, fmt.Errorf("unknown configuration option detected for reset: %s", cc.InitialReset)
// }
boot := cc.BootNotification != nil && *cc.BootNotification
noConfig := cc.GetConfiguration != nil && !*cc.GetConfiguration
@ -184,7 +170,6 @@ func NewOCPP(id string, connector int, idtag string,
for _, opt := range resp.ConfigurationKey {
if opt.Value == nil {
c.log.ERROR.Printf("%s (%s): %s", opt.Key, rw[opt.Readonly], "nil")
continue
}
@ -259,22 +244,9 @@ func NewOCPP(id string, connector int, idtag string,
}
}
// TODO deprecate
// if initialReset != "" {
// t := core.ResetTypeSoft
// if initialReset == core.ResetTypeHard {
// t = core.ResetTypeHard
// }
// ocpp.Instance().TriggerResetRequest(cp.ID(), t)
// }
// request initial status
_ = cp.Initialized(statusTimeout)
// TODO: check for running transaction
return c, nil
return c, cp.Initialized()
}
// hasMeasurement checks if meterValuesSample contains given measurement
@ -317,7 +289,8 @@ func (c *OCPP) Status() (api.ChargeStatus, error) {
// Enabled implements the api.Charger interface
func (c *OCPP) Enabled() (bool, error) {
return c.cp.TransactionID() > 0, nil
txn, err := c.cp.TransactionID()
return txn > 0, err
}
// Enable implements the api.Charger interface
@ -337,13 +310,19 @@ func (c *OCPP) Enable(enable bool) error {
request.ChargingProfile = getTxChargingProfile(c.current, c.phases)
})
} else {
var txn int
txn, err = c.cp.TransactionID()
if err != nil {
return err
}
err = ocpp.Instance().RemoteStopTransaction(c.cp.ID(), func(resp *core.RemoteStopTransactionConfirmation, err error) {
if err == nil && resp != nil && resp.Status != types.RemoteStartStopStatusAccepted {
err = errors.New(string(resp.Status))
}
rc <- err
}, c.cp.TransactionID())
}, txn)
}
return c.wait(err, rc)

View file

@ -7,6 +7,7 @@ import (
"sync"
"time"
"github.com/benbjohnson/clock"
"github.com/evcc-io/evcc/api"
"github.com/evcc-io/evcc/util"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/core"
@ -36,17 +37,18 @@ const (
// Since ocpp-go interfaces at charge point level, we need to manage multiple connector separately
type CP struct {
mu sync.Mutex
log *util.Logger
once sync.Once
mu sync.Mutex
once sync.Once
clock clock.Clock // mockable time
log *util.Logger
id string
connector int
connectC, statusC chan struct{}
connected bool
status *core.StatusNotificationRequest
updated time.Time
meterUpdated time.Time
timeout time.Duration
@ -58,6 +60,7 @@ type CP struct {
func NewChargePoint(log *util.Logger, id string, connector int, timeout time.Duration) *CP {
return &CP{
clock: clock.New(),
log: log,
id: id,
connector: connector,
@ -68,6 +71,10 @@ func NewChargePoint(log *util.Logger, id string, connector int, timeout time.Dur
}
}
func (cp *CP) TestClock(clock clock.Clock) {
cp.clock = clock
}
func (cp *CP) ID() string {
cp.mu.Lock()
defer cp.mu.Unlock()
@ -86,24 +93,26 @@ func (cp *CP) RegisterID(id string) {
cp.id = id
}
func (cp *CP) Connect() {
func (cp *CP) connect(connect bool) {
cp.mu.Lock()
defer cp.mu.Unlock()
cp.once.Do(func() {
close(cp.connectC)
})
cp.connected = connect
if connect {
cp.once.Do(func() {
close(cp.connectC)
})
}
}
func (cp *CP) HasConnected() <-chan struct{} {
return cp.connectC
}
func (cp *CP) Initialized(timeout time.Duration) bool {
cp.log.DEBUG.Printf("waiting for chargepoint status: %v", timeout)
func (cp *CP) Initialized() error {
// trigger status
time.AfterFunc(5*time.Second, func() {
time.AfterFunc(cp.timeout/2, func() {
select {
case <-cp.statusC:
return
@ -115,32 +124,22 @@ func (cp *CP) Initialized(timeout time.Duration) bool {
// wait for status
select {
case <-cp.statusC:
cp.update()
return true
case <-time.After(timeout):
return false
}
}
// WatchDog triggers meter values messages if older than timeout.
// Must be wrapped in a goroutine.
func (cp *CP) WatchDog(timeout time.Duration) {
for ; true; <-time.NewTicker(timeout).C {
cp.mu.Lock()
update := cp.txnId != 0 && time.Since(cp.meterUpdated) > timeout
cp.mu.Unlock()
if update {
Instance().TriggerMessageRequest(cp.ID(), core.MeterValuesFeatureName)
}
return nil
case <-time.After(cp.timeout):
return api.ErrTimeout
}
}
// TransactionID returns the current transaction id
func (cp *CP) TransactionID() int {
func (cp *CP) TransactionID() (int, error) {
cp.mu.Lock()
defer cp.mu.Unlock()
return cp.txnId
if !cp.connected {
return 0, api.ErrTimeout
}
return cp.txnId, nil
}
func (cp *CP) Status() (api.ChargeStatus, error) {
@ -149,7 +148,7 @@ func (cp *CP) Status() (api.ChargeStatus, error) {
res := api.StatusNone
if time.Since(cp.updated) > cp.timeout {
if !cp.connected {
return res, api.ErrTimeout
}
@ -179,13 +178,31 @@ func (cp *CP) Status() (api.ChargeStatus, error) {
return res, nil
}
// WatchDog triggers meter values messages if older than timeout.
// Must be wrapped in a goroutine.
func (cp *CP) WatchDog(timeout time.Duration) {
for ; true; <-time.NewTicker(timeout).C {
cp.mu.Lock()
update := cp.txnId != 0 && cp.clock.Since(cp.meterUpdated) > timeout
cp.mu.Unlock()
if update {
Instance().TriggerMessageRequest(cp.ID(), core.MeterValuesFeatureName)
}
}
}
var _ api.Meter = (*CP)(nil)
func (cp *CP) CurrentPower() (float64, error) {
cp.mu.Lock()
defer cp.mu.Unlock()
if cp.txnId != 0 && cp.timeout > 0 && time.Since(cp.meterUpdated) > cp.timeout {
if !cp.connected {
return 0, api.ErrTimeout
}
if cp.txnId != 0 && cp.timeout > 0 && cp.clock.Since(cp.meterUpdated) > cp.timeout {
return 0, api.ErrNotAvailable
}
@ -203,7 +220,11 @@ func (cp *CP) TotalEnergy() (float64, error) {
cp.mu.Lock()
defer cp.mu.Unlock()
if cp.txnId != 0 && cp.timeout > 0 && time.Since(cp.meterUpdated) > cp.timeout {
if !cp.connected {
return 0, api.ErrTimeout
}
if cp.txnId != 0 && cp.timeout > 0 && cp.clock.Since(cp.meterUpdated) > cp.timeout {
return 0, api.ErrNotAvailable
}
@ -236,7 +257,11 @@ func (cp *CP) Currents() (float64, float64, float64, error) {
cp.mu.Lock()
defer cp.mu.Unlock()
if cp.txnId != 0 && cp.timeout > 0 && time.Since(cp.meterUpdated) > cp.timeout {
if !cp.connected {
return 0, 0, 0, api.ErrTimeout
}
if cp.txnId != 0 && cp.timeout > 0 && cp.clock.Since(cp.meterUpdated) > cp.timeout {
return 0, 0, 0, api.ErrNotAvailable
}

View file

@ -76,14 +76,7 @@ func (cp *CP) DataTransfer(request *core.DataTransferRequest) (*core.DataTransfe
return res, nil
}
func (cp *CP) update() {
cp.mu.Lock()
cp.updated = time.Now()
cp.mu.Unlock()
}
func (cp *CP) Heartbeat(request *core.HeartbeatRequest) (*core.HeartbeatConfirmation, error) {
cp.update()
res := &core.HeartbeatConfirmation{
CurrentTime: types.NewDateTime(time.Now()),
}

View file

@ -16,6 +16,8 @@ type CS struct {
cps map[string]*CP
}
// Register registers a chargepoint with the central system.
// The chargepoint identified by id may already be connected in which case initial connection is triggered.
func (cs *CS) Register(id string, cp *CP) error {
cs.mu.Lock()
defer cs.mu.Unlock()
@ -24,6 +26,11 @@ func (cs *CS) Register(id string, cp *CP) error {
return errors.New("cannot have >1 chargepoint with empty station id")
}
// trigger unknown chargepoint connected
if unknown, ok := cs.cps[id]; ok && unknown == nil {
cp.connect(true)
}
cs.cps[id] = cp
return nil
@ -49,6 +56,7 @@ func (cs *CS) NewChargePoint(chargePoint ocpp16.ChargePointConnection) {
defer cs.mu.Unlock()
if cp, err := cs.chargepointByID(chargePoint.ID()); err != nil {
// check for anonymous chargepoint
if cp, ok := cs.cps[""]; ok {
cs.log.INFO.Printf("chargepoint connected, registering: %s", chargePoint.ID())
@ -57,15 +65,23 @@ func (cs *CS) NewChargePoint(chargePoint ocpp16.ChargePointConnection) {
cs.cps[chargePoint.ID()] = cp
delete(cs.cps, "")
cp.Connect()
cp.connect(true)
return
}
cs.log.WARN.Printf("chargepoint connected, ignoring: %s", chargePoint.ID())
cs.log.WARN.Printf("chargepoint connected, unknown: %s", chargePoint.ID())
// register unknown chargepoint
// when chargepoint setup is complete, it will eventually be associated with the connected id
cs.cps[chargePoint.ID()] = nil
} else {
cs.log.DEBUG.Printf("chargepoint connected: %s", chargePoint.ID())
cp.Connect()
// trigger initial connection if chargepoint is already setup
if cp != nil {
cp.connect(true)
}
}
}
@ -73,10 +89,17 @@ func (cs *CS) ChargePointDisconnected(chargePoint ocpp16.ChargePointConnection)
cs.mu.Lock()
defer cs.mu.Unlock()
if _, err := cs.chargepointByID(chargePoint.ID()); err != nil {
if cp, err := cs.chargepointByID(chargePoint.ID()); err != nil {
cs.log.ERROR.Printf("chargepoint disconnected: %v", err)
} else {
cs.log.DEBUG.Printf("chargepoint disconnected: %s", chargePoint.ID())
if cp == nil {
// remove unknown chargepoint
delete(cs.cps, chargePoint.ID())
} else {
cp.connect(false)
}
}
}

View file

@ -1,6 +1,7 @@
package ocpp
import (
"sync"
"time"
"github.com/evcc-io/evcc/util"
@ -8,10 +9,13 @@ import (
"github.com/lorenzodonini/ocpp-go/ocppj"
)
var instance *CS
var (
once sync.Once
instance *CS
)
func Instance() *CS {
if instance == nil {
once.Do(func() {
cs := ocpp16.NewCentralSystem(nil, nil)
instance = &CS{
@ -27,11 +31,11 @@ func Instance() *CS {
cs.SetChargePointDisconnectedHandler(instance.ChargePointDisconnected)
cs.SetFirmwareManagementHandler(instance)
go Instance().errorHandler(cs.Errors())
go instance.errorHandler(cs.Errors())
go cs.Start(8887, "/{ws}")
time.Sleep(time.Second)
}
})
return instance
}

106
charger/ocpp_test.go Normal file
View file

@ -0,0 +1,106 @@
package charger
import (
"testing"
"time"
"github.com/benbjohnson/clock"
"github.com/evcc-io/evcc/charger/ocpp"
ocpp16 "github.com/lorenzodonini/ocpp-go/ocpp1.6"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/core"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/remotetrigger"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/types"
"github.com/stretchr/testify/suite"
)
const (
ocppTestUrl = "ws://localhost:8887"
ocppTestConnectTimeout = 10 * time.Second
ocppTestTimeout = 3 * time.Second
ocppTestConnector = 1
)
func TestOcpp(t *testing.T) {
suite.Run(t, new(ocppTestSuite))
}
type ocppTestSuite struct {
suite.Suite
cp ocpp16.ChargePoint
}
func (suite *ocppTestSuite) SetupSuite() {
// setup cs
suite.NotNil(ocpp.Instance())
// setup cp
cp := ocpp16.NewChargePoint("test", nil, nil)
// set a handler for all callback functions
triggerC := make(chan remotetrigger.MessageTrigger, 1)
handler := &ChargePointHandler{triggerC: triggerC}
cp.SetCoreHandler(handler)
cp.SetRemoteTriggerHandler(handler)
go func() {
for msg := range triggerC {
suite.handleTrigger(msg)
}
}()
suite.cp = cp
}
func (suite *ocppTestSuite) handleTrigger(msg remotetrigger.MessageTrigger) {
switch msg {
case core.BootNotificationFeatureName:
if res, err := suite.cp.BootNotification("demo", "evcc"); err != nil {
suite.T().Log("BootNotification:", err)
} else {
suite.T().Log("BootNotification:", res)
}
case core.StatusNotificationFeatureName:
if res, err := suite.cp.StatusNotification(ocppTestConnector, core.NoError, core.ChargePointStatusAvailable); err != nil {
suite.T().Log("StatusNotification:", err)
} else {
suite.T().Log("StatusNotification:", res)
}
case core.MeterValuesFeatureName:
if res, err := suite.cp.MeterValues(1, []types.MeterValue{
{SampledValue: []types.SampledValue{
{Measurand: types.MeasurandPowerActiveImport, Value: "1000"},
}},
}); err != nil {
suite.T().Log("MeterValues:", err)
} else {
suite.T().Log("MeterValues:", res)
}
default:
suite.T().Log(msg)
}
}
func (suite *ocppTestSuite) TestConnect() {
// start cp client
suite.NoError(suite.cp.Start(ocppTestUrl))
suite.True(suite.cp.IsConnected())
// start cp server
c, err := NewOCPP("test", ocppTestConnector, "", "", 0, false, false, ocppTestConnectTimeout, ocppTestTimeout)
suite.NoError(err)
if err != nil {
return
}
clock := clock.NewMock()
c.cp.TestClock(clock)
clock.Add(ocppTestTimeout)
_, err = c.Status()
suite.NoError(err)
}

View file

@ -0,0 +1,79 @@
package charger
import (
"fmt"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/core"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/remotetrigger"
"github.com/lorenzodonini/ocpp-go/ocpp1.6/types"
)
type ChargePointHandler struct {
triggerC chan remotetrigger.MessageTrigger
}
func (handler *ChargePointHandler) OnChangeAvailability(request *core.ChangeAvailabilityRequest) (confirmation *core.ChangeAvailabilityConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewChangeAvailabilityConfirmation(core.AvailabilityStatusAccepted), nil
}
func (handler *ChargePointHandler) OnChangeConfiguration(request *core.ChangeConfigurationRequest) (confirmation *core.ChangeConfigurationConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewChangeConfigurationConfirmation(core.ConfigurationStatusAccepted), nil
}
func (handler *ChargePointHandler) OnClearCache(request *core.ClearCacheRequest) (confirmation *core.ClearCacheConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewClearCacheConfirmation(core.ClearCacheStatusAccepted), nil
}
func (handler *ChargePointHandler) OnDataTransfer(request *core.DataTransferRequest) (confirmation *core.DataTransferConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewDataTransferConfirmation(core.DataTransferStatusAccepted), nil
}
func (handler *ChargePointHandler) OnGetConfiguration(request *core.GetConfigurationRequest) (confirmation *core.GetConfigurationConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
one := "1"
return core.NewGetConfigurationConfirmation([]core.ConfigurationKey{
{Key: "AuthorizationKey"},
{Key: "NumberOfConnectors", Value: &one},
{Key: "ChargeProfileMaxStackLevel", Value: &one},
{Key: "ChargingScheduleMaxPeriods", Value: &one},
{Key: "MaxChargingProfilesInstalled", Value: &one},
{Key: "ChargingScheduleAllowedChargingRateUnit", Value: &one},
}), nil
}
func (handler *ChargePointHandler) OnRemoteStartTransaction(request *core.RemoteStartTransactionRequest) (confirmation *core.RemoteStartTransactionConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewRemoteStartTransactionConfirmation(types.RemoteStartStopStatusAccepted), nil
}
func (handler *ChargePointHandler) OnRemoteStopTransaction(request *core.RemoteStopTransactionRequest) (confirmation *core.RemoteStopTransactionConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewRemoteStopTransactionConfirmation(types.RemoteStartStopStatusAccepted), nil
}
func (handler *ChargePointHandler) OnReset(request *core.ResetRequest) (confirmation *core.ResetConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewResetConfirmation(core.ResetStatusAccepted), nil
}
func (handler *ChargePointHandler) OnUnlockConnector(request *core.UnlockConnectorRequest) (confirmation *core.UnlockConnectorConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
return core.NewUnlockConnectorConfirmation(core.UnlockStatusUnlocked), nil
}
func (handler *ChargePointHandler) OnTriggerMessage(request *remotetrigger.TriggerMessageRequest) (confirmation *remotetrigger.TriggerMessageConfirmation, err error) {
fmt.Printf("%T %+v\n", request, request)
if c := handler.triggerC; request != nil && c != nil {
select {
case c <- request.RequestedMessage:
default:
}
}
return remotetrigger.NewTriggerMessageConfirmation(remotetrigger.TriggerMessageStatusAccepted), nil
}