From 6bbfeffb8f201d112f08ffd908b2caaca6cc5204 Mon Sep 17 00:00:00 2001 From: andig Date: Tue, 13 Sep 2022 18:01:41 +0200 Subject: [PATCH] Ocpp: fix chargepoint registration and startup (#4420) --- charger/ocpp.go | 15 +++++--- charger/ocpp/cp.go | 73 ++++++++++++++++++++++++++------------ charger/ocpp/cp_core.go | 12 ++++--- charger/ocpp/cs.go | 16 ++++----- charger/ocpp/cs_core.go | 58 +++++++++++++++--------------- cmd/ocpp/handler.go | 78 +++++++++++++++++++++++++++++++++++++++++ cmd/ocpp/main.go | 68 +++++++++++++++++++++++++++++++++++ 7 files changed, 250 insertions(+), 70 deletions(-) create mode 100644 cmd/ocpp/handler.go create mode 100644 cmd/ocpp/main.go diff --git a/charger/ocpp.go b/charger/ocpp.go index 9591f2f3f..19fe8b95f 100644 --- a/charger/ocpp.go +++ b/charger/ocpp.go @@ -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 diff --git a/charger/ocpp/cp.go b/charger/ocpp/cp.go index cfdbf1dd4..30503c975 100644 --- a/charger/ocpp/cp.go +++ b/charger/ocpp/cp.go @@ -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 { diff --git a/charger/ocpp/cp_core.go b/charger/ocpp/cp_core.go index 3e0ad74b2..56b6cdea4 100644 --- a/charger/ocpp/cp_core.go +++ b/charger/ocpp/cp_core.go @@ -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 { diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go index 22c25a618..c591e8fd3 100644 --- a/charger/ocpp/cs.go +++ b/charger/ocpp/cs.go @@ -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()) } diff --git a/charger/ocpp/cs_core.go b/charger/ocpp/cs_core.go index 7cfb0dd65..b579a633b 100644 --- a/charger/ocpp/cs_core.go +++ b/charger/ocpp/cs_core.go @@ -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 } diff --git a/cmd/ocpp/handler.go b/cmd/ocpp/handler.go new file mode 100644 index 000000000..13ee89396 --- /dev/null +++ b/cmd/ocpp/handler.go @@ -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 +} diff --git a/cmd/ocpp/main.go b/cmd/ocpp/main.go new file mode 100644 index 000000000..bc51dc65c --- /dev/null +++ b/cmd/ocpp/main.go @@ -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 {} +}