Ocpp: fix chargepoint registration and startup (#4420)

This commit is contained in:
andig 2022-09-13 18:01:41 +02:00 • committed by GitHub
parent 2492e4b0ae
commit 6bbfeffb8f
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 250 additions and 70 deletions

View file

@ -18,6 +18,8 @@ import (
"github.com/samber/lo"
)
const statusTimeout = 30 * time.Second
// OCPP charger implementation
type OCPP struct {
log *util.Logger
@ -94,7 +96,7 @@ func NewOCPP(id string, connector int, idtag string, meterValues string, meterIn
return nil, err
}
unit := "ocpp-cp"
unit := "ocpp"
if id != "" {
unit = id
}
@ -107,7 +109,6 @@ func NewOCPP(id string, connector int, idtag string, meterValues string, meterIn
idtag: idtag,
}
c.log.DEBUG.Println("waiting for chargepoint to register")
if err := cp.Boot(); err != nil {
return nil, err
}
@ -118,7 +119,8 @@ func NewOCPP(id string, connector int, idtag string, meterValues string, meterIn
meterSampleInterval time.Duration
)
err = ocpp.Instance().GetConfiguration(id, func(resp *core.GetConfigurationConfirmation, err error) {
// configured id may be empty, use registered id below
err = ocpp.Instance().GetConfiguration(cp.ID(), func(resp *core.GetConfigurationConfirmation, err error) {
if err != nil {
rc <- err
return
@ -187,7 +189,7 @@ func NewOCPP(id string, connector int, idtag string, meterValues string, meterIn
// get initial meter values and configure sample rate
if c.hasMeasurement("Power.Active.Import") || c.hasMeasurement("Energy.Active.Import.Register") {
ocpp.Instance().TriggerMeterValuesRequest(cp)
ocpp.Instance().TriggerMessageRequest(cp.ID(), core.MeterValuesFeatureName)
if meterSampleInterval > meterInterval && meterInterval > 0 {
if err := c.configure(ocpp.KeyMeterValueSampleInterval, strconv.Itoa(int(meterInterval.Seconds()))); err != nil {
@ -206,9 +208,12 @@ func NewOCPP(id string, connector int, idtag string, meterValues string, meterIn
t = core.ResetTypeHard
}
ocpp.Instance().TriggerResetRequest(cp, t)
ocpp.Instance().TriggerResetRequest(cp.ID(), t)
}
// request initial status
_ = cp.Initialized(statusTimeout)
// TODO: check for running transaction
return c, nil

View file

@ -52,10 +52,9 @@ type CP struct {
log *util.Logger
id string
updated time.Time
initialized *sync.Cond
boot *core.BootNotificationRequest
status *core.StatusNotificationRequest
bootC, statusC chan struct{}
updated time.Time
status *core.StatusNotificationRequest
timeout time.Duration
meterUpdated time.Time
@ -68,6 +67,47 @@ type CP struct {
txnId int
}
func (cp *CP) ID() string {
cp.mu.Lock()
defer cp.mu.Unlock()
return cp.id
}
func (cp *CP) RegisterID(id string) {
cp.mu.Lock()
defer cp.mu.Unlock()
if cp.id != "" {
panic("ocpp: cannot re-register id")
}
cp.id = id
}
func (cp *CP) Initialized(timeout time.Duration) bool {
cp.log.DEBUG.Printf("waiting for chargepoint status: %v", timeout)
// trigger status
time.AfterFunc(5*time.Second, func() {
select {
case <-cp.statusC:
return
default:
Instance().TriggerMessageRequest(cp.ID(), core.StatusNotificationFeatureName)
}
})
// wait for status
select {
case <-cp.statusC:
cp.update()
return true
case <-time.After(timeout):
return false
}
}
func (cp *CP) WatchDog(timeout time.Duration) {
cp.timeout = timeout
@ -78,7 +118,7 @@ func (cp *CP) WatchDog(timeout time.Duration) {
cp.mu.Unlock()
if update {
Instance().TriggerMeterValuesRequest(cp)
Instance().TriggerMessageRequest(cp.ID(), core.MeterValuesFeatureName)
}
}
}()
@ -148,8 +188,8 @@ func detectSmartChargingCapabilities(options map[string]core.ConfigurationKey) (
{ // optional
var supported bool
opt, found := options[string(KeyConnectorSwitch3to1PhaseSupported)]
if found {
if opt, ok := options[string(KeyConnectorSwitch3to1PhaseSupported)]; ok {
var err error
supported, err = strconv.ParseBool(*opt.Value)
if err != nil {
@ -179,21 +219,10 @@ func parseIntOption(key SmartchargingChargeProfileKey, options map[string]core.C
// Boot waits for the CP to register itself
func (cp *CP) Boot() error {
bootC := make(chan struct{})
go func() {
cp.mu.Lock()
defer cp.mu.Unlock()
for cp.boot == nil || cp.status == nil {
cp.initialized.Wait()
}
close(bootC)
}()
cp.log.DEBUG.Printf("waiting for chargepoint: %v", timeout)
select {
case <-bootC:
cp.update()
case <-cp.bootC:
return nil
case <-time.After(timeout):
return api.ErrTimeout
@ -213,14 +242,12 @@ func (cp *CP) Status() (api.ChargeStatus, error) {
res := api.StatusNone
cp.log.TRACE.Printf("current transaction ID (last update): %d (%s)", cp.txnId, cp.updated.Format(time.RFC3339))
if time.Since(cp.updated) > timeout {
return res, api.ErrTimeout
}
if cp.status.ErrorCode != core.NoError {
cp.log.DEBUG.Printf("chargepoint error: %s: %s", cp.status.ErrorCode, cp.status.Info)
return res, fmt.Errorf("%s: %s", cp.status.ErrorCode, cp.status.Info)
}
switch cp.status.Status {

View file

@ -30,11 +30,13 @@ func (cp *CP) BootNotification(request *core.BootNotificationRequest) (*core.Boo
cp.log.TRACE.Printf("%T: %+v", request, request)
if request != nil {
cp.mu.Lock()
defer cp.mu.Unlock()
cp.log.DEBUG.Printf("chargepoint: %+v", request)
}
cp.boot = request
cp.initialized.Broadcast()
// signal boot
select {
case cp.bootC <- struct{}{}:
default:
}
res := &core.BootNotificationConfirmation{
@ -71,7 +73,7 @@ func (cp *CP) StatusNotification(request *core.StatusNotificationRequest) (*core
if cp.status == nil {
cp.status = request
cp.initialized.Broadcast()
close(cp.statusC) // signal initial status received
} else if request.Timestamp == nil || cp.timestampValid(request.Timestamp.Time) {
cp.status = request
} else {

View file

@ -18,7 +18,7 @@ type CS struct {
}
func (cs *CS) Register(id string) (*CP, error) {
unit := "ocpp-cp"
unit := "ocpp"
if id != "" {
unit = id
}
@ -26,11 +26,11 @@ func (cs *CS) Register(id string) (*CP, error) {
cp := &CP{
id: id,
log: util.NewLogger(unit),
bootC: make(chan struct{}),
statusC: make(chan struct{}),
measurements: make(map[string]types.SampledValue),
}
cp.initialized = sync.NewCond(&cp.mu)
cs.mu.Lock()
defer cs.mu.Unlock()
@ -63,18 +63,18 @@ func (cs *CS) NewChargePoint(chargePoint ocpp16.ChargePointConnection) {
defer cs.mu.Unlock()
if _, err := cs.chargepointByID(chargePoint.ID()); err != nil {
if auto, ok := cs.cps[""]; ok {
cs.log.INFO.Printf("unknown chargepoint connected, registering: %s", chargePoint.ID())
if cp, ok := cs.cps[""]; ok {
cs.log.INFO.Printf("chargepoint connected, registering: %s", chargePoint.ID())
// update id
auto.id = chargePoint.ID()
cs.cps[chargePoint.ID()] = auto
cp.RegisterID(chargePoint.ID())
cs.cps[chargePoint.ID()] = cp
delete(cs.cps, "")
return
}
cs.log.WARN.Printf("unknown chargepoint connected, ignored: %s", chargePoint.ID())
cs.log.WARN.Printf("chargepoint connected, ignoring: %s", chargePoint.ID())
} else {
cs.log.DEBUG.Printf("chargepoint connected: %s", chargePoint.ID())
}

View file

@ -8,8 +8,8 @@ import (
// cs actions
func (cs *CS) TriggerResetRequest(cp *CP, resetType core.ResetType) {
if err := cs.Reset(cp.id, func(request *core.ResetConfirmation, err error) {
func (cs *CS) TriggerResetRequest(id string, resetType core.ResetType) {
if err := cs.Reset(id, func(request *core.ResetConfirmation, err error) {
log := cs.log.TRACE
if err == nil && request != nil && request.Status != core.ResetStatusAccepted {
log = cs.log.ERROR
@ -20,14 +20,14 @@ func (cs *CS) TriggerResetRequest(cp *CP, resetType core.ResetType) {
status = request.Status
}
log.Printf("TriggerReset for %s: %+v", cp.id, status)
log.Printf("TriggerReset for %s: %+v", id, status)
}, resetType); err != nil {
cs.log.ERROR.Printf("send TriggerReset for %s failed: %v", cp.id, err)
cs.log.ERROR.Printf("send TriggerReset for %s failed: %v", id, err)
}
}
func (cs *CS) TriggerMeterValuesRequest(cp *CP) {
if err := cs.TriggerMessage(cp.id, func(request *remotetrigger.TriggerMessageConfirmation, err error) {
func (cs *CS) TriggerMessageRequest(id string, requestedMessage remotetrigger.MessageTrigger) {
if err := cs.TriggerMessage(id, func(request *remotetrigger.TriggerMessageConfirmation, err error) {
log := cs.log.TRACE
if err == nil && request != nil && request.Status != remotetrigger.TriggerMessageStatusAccepted {
log = cs.log.ERROR
@ -38,16 +38,16 @@ func (cs *CS) TriggerMeterValuesRequest(cp *CP) {
status = request.Status
}
log.Printf("TriggerMessage %s for %s: %+v", core.MeterValuesFeatureName, cp.id, status)
}, core.MeterValuesFeatureName); err != nil {
cs.log.ERROR.Printf("send TriggerMessage for %s failed: %v", cp.id, err)
log.Printf("TriggerMessage %s for %s: %+v", requestedMessage, id, status)
}, requestedMessage); err != nil {
cs.log.ERROR.Printf("send TriggerMessage %s for %s failed: %v", requestedMessage, id, err)
}
}
// cp actions
func (cs *CS) OnAuthorize(chargePointId string, request *core.AuthorizeRequest) (*core.AuthorizeConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnAuthorize(id string, request *core.AuthorizeRequest) (*core.AuthorizeConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -55,8 +55,8 @@ func (cs *CS) OnAuthorize(chargePointId string, request *core.AuthorizeRequest)
return cp.Authorize(request)
}
func (cs *CS) OnBootNotification(chargePointId string, request *core.BootNotificationRequest) (*core.BootNotificationConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnBootNotification(id string, request *core.BootNotificationRequest) (*core.BootNotificationConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -64,8 +64,8 @@ func (cs *CS) OnBootNotification(chargePointId string, request *core.BootNotific
return cp.BootNotification(request)
}
func (cs *CS) OnDataTransfer(chargePointId string, request *core.DataTransferRequest) (*core.DataTransferConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnDataTransfer(id string, request *core.DataTransferRequest) (*core.DataTransferConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -73,8 +73,8 @@ func (cs *CS) OnDataTransfer(chargePointId string, request *core.DataTransferReq
return cp.DataTransfer(request)
}
func (cs *CS) OnHeartbeat(chargePointId string, request *core.HeartbeatRequest) (*core.HeartbeatConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnHeartbeat(id string, request *core.HeartbeatRequest) (*core.HeartbeatConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -82,8 +82,8 @@ func (cs *CS) OnHeartbeat(chargePointId string, request *core.HeartbeatRequest)
return cp.Heartbeat(request)
}
func (cs *CS) OnMeterValues(chargePointId string, request *core.MeterValuesRequest) (*core.MeterValuesConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnMeterValues(id string, request *core.MeterValuesRequest) (*core.MeterValuesConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -91,8 +91,8 @@ func (cs *CS) OnMeterValues(chargePointId string, request *core.MeterValuesReque
return cp.MeterValues(request)
}
func (cs *CS) OnStatusNotification(chargePointId string, request *core.StatusNotificationRequest) (*core.StatusNotificationConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnStatusNotification(id string, request *core.StatusNotificationRequest) (*core.StatusNotificationConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -100,8 +100,8 @@ func (cs *CS) OnStatusNotification(chargePointId string, request *core.StatusNot
return cp.StatusNotification(request)
}
func (cs *CS) OnStartTransaction(chargePointId string, request *core.StartTransactionRequest) (*core.StartTransactionConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnStartTransaction(id string, request *core.StartTransactionRequest) (*core.StartTransactionConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -109,8 +109,8 @@ func (cs *CS) OnStartTransaction(chargePointId string, request *core.StartTransa
return cp.StartTransaction(request)
}
func (cs *CS) OnStopTransaction(chargePointId string, request *core.StopTransactionRequest) (*core.StopTransactionConfirmation, error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnStopTransaction(id string, request *core.StopTransactionRequest) (*core.StopTransactionConfirmation, error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -118,8 +118,8 @@ func (cs *CS) OnStopTransaction(chargePointId string, request *core.StopTransact
return cp.StopTransaction(request)
}
func (cs *CS) OnDiagnosticsStatusNotification(chargePointId string, request *firmware.DiagnosticsStatusNotificationRequest) (confirmation *firmware.DiagnosticsStatusNotificationConfirmation, err error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnDiagnosticsStatusNotification(id string, request *firmware.DiagnosticsStatusNotificationRequest) (confirmation *firmware.DiagnosticsStatusNotificationConfirmation, err error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}
@ -127,8 +127,8 @@ func (cs *CS) OnDiagnosticsStatusNotification(chargePointId string, request *fir
return cp.DiagnosticStatusNotification(request)
}
func (cs *CS) OnFirmwareStatusNotification(chargePointId string, request *firmware.FirmwareStatusNotificationRequest) (confirmation *firmware.FirmwareStatusNotificationConfirmation, err error) {
cp, err := cs.chargepointByID(chargePointId)
func (cs *CS) OnFirmwareStatusNotification(id string, request *firmware.FirmwareStatusNotificationRequest) (confirmation *firmware.FirmwareStatusNotificationConfirmation, err error) {
cp, err := cs.chargepointByID(id)
if err != nil {
return nil, err
}

78
cmd/ocpp/handler.go Normal file
View file

@ -0,0 +1,78 @@
package main
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: "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
}

68
cmd/ocpp/main.go Normal file
View file

@ -0,0 +1,68 @@
package main
import (
"fmt"
"log"
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"
)
const (
chargePointId = "cp0001"
url = "ws://localhost:8887"
)
func main() {
chargePoint := ocpp16.NewChargePoint(chargePointId, nil, nil)
// Set a handler for all callback functions
handler := &ChargePointHandler{triggerC: make(chan remotetrigger.MessageTrigger, 1)}
chargePoint.SetCoreHandler(handler)
chargePoint.SetRemoteTriggerHandler(handler)
go func() {
for err := range chargePoint.Errors() {
fmt.Println(err)
}
}()
go func() {
for msg := range handler.triggerC {
fmt.Println("msg:", msg)
switch msg {
case core.StatusNotificationFeatureName:
if res, err := chargePoint.StatusNotification(1, core.NoError, core.ChargePointStatusAvailable); err != nil {
log.Println("StatusNotification:", err)
} else {
log.Println("StatusNotification:", res)
}
case core.MeterValuesFeatureName:
if _, err := chargePoint.MeterValues(1, []types.MeterValue{
{SampledValue: []types.SampledValue{
{Measurand: types.MeasurandPowerActiveImport, Value: "1000"},
}},
}); err != nil {
log.Println("MeterValues:", err)
}
}
}
}()
// Connects to central system
if err := chargePoint.Start(url); err != nil {
log.Fatal(err)
}
log.Printf("connected to central system at %v", url)
if res, err := chargePoint.BootNotification("model1", "vendor1"); err != nil {
log.Fatal("BootNotification", err)
} else {
log.Printf("status: %v, interval: %v", res.Status, res.Interval)
}
select {}
}