diff --git a/charger/ocpp.go b/charger/ocpp.go index b38605ba1..f2b514d66 100644 --- a/charger/ocpp.go +++ b/charger/ocpp.go @@ -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) diff --git a/charger/ocpp/cp.go b/charger/ocpp/cp.go index ab732cc71..4cb2b014e 100644 --- a/charger/ocpp/cp.go +++ b/charger/ocpp/cp.go @@ -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 } diff --git a/charger/ocpp/cp_core.go b/charger/ocpp/cp_core.go index 4cace1a94..ccbe6a825 100644 --- a/charger/ocpp/cp_core.go +++ b/charger/ocpp/cp_core.go @@ -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()), } diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go index d5674995a..a3a188570 100644 --- a/charger/ocpp/cs.go +++ b/charger/ocpp/cs.go @@ -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) + } } } diff --git a/charger/ocpp/instance.go b/charger/ocpp/instance.go index a0ae7d513..ea15bb125 100644 --- a/charger/ocpp/instance.go +++ b/charger/ocpp/instance.go @@ -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 } diff --git a/charger/ocpp_test.go b/charger/ocpp_test.go new file mode 100644 index 000000000..3f98d750e --- /dev/null +++ b/charger/ocpp_test.go @@ -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) +} diff --git a/charger/ocpp_test_handler.go b/charger/ocpp_test_handler.go new file mode 100644 index 000000000..889b9dcc0 --- /dev/null +++ b/charger/ocpp_test_handler.go @@ -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 +}