From 421a993a343b8a72baf773db870aff5f14cfdd10 Mon Sep 17 00:00:00 2001 From: andig Date: Sat, 24 Sep 2022 13:03:45 +0200 Subject: [PATCH] Ocpp: don't rely on charger sending boot notification (#4567) --- charger/ocpp.go | 32 +++++++++++++++--------- charger/ocpp/cp.go | 54 +++++++++++++++++++++++++---------------- charger/ocpp/cp_core.go | 10 -------- charger/ocpp/cs.go | 25 ++++++------------- cmd/ocpp/main.go | 53 ++++++++++++++++++++++++++++++++++------ 5 files changed, 106 insertions(+), 68 deletions(-) diff --git a/charger/ocpp.go b/charger/ocpp.go index 0e965b56e..24a853346 100644 --- a/charger/ocpp.go +++ b/charger/ocpp.go @@ -46,8 +46,10 @@ func NewOCPPFromConfig(other map[string]interface{}) (api.Charger, error) { MeterInterval time.Duration MeterValues string InitialReset core.ResetType + Timeout time.Duration }{ Connector: 1, + Timeout: time.Minute, } if err := util.DecodeOther(other, &cc); err != nil { @@ -63,7 +65,7 @@ func NewOCPPFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, fmt.Errorf("unknown configuration option detected for reset: %s", cc.InitialReset) } - c, err := NewOCPP(cc.StationId, cc.Connector, cc.IdTag, cc.MeterValues, cc.MeterInterval, cc.InitialReset) + c, err := NewOCPP(cc.StationId, cc.Connector, cc.IdTag, cc.MeterValues, cc.MeterInterval, cc.InitialReset, cc.Timeout) if err != nil { return c, err } @@ -89,28 +91,36 @@ func NewOCPPFromConfig(other map[string]interface{}) (api.Charger, error) { //go:generate go run ../cmd/tools/decorate.go -f decorateOCPP -b *OCPP -r api.Charger -t "api.Meter,CurrentPower,func() (float64, error)" -t "api.MeterEnergy,TotalEnergy,func() (float64, error)" -t "api.MeterCurrent,Currents,func() (float64, float64, float64, error)" // NewOCPP creates OCPP charger -func NewOCPP(id string, connector int, idtag string, meterValues string, meterInterval time.Duration, initialReset core.ResetType) (*OCPP, error) { - cp, err := ocpp.Instance().Register(id) - if err != nil { - return nil, err - } - +func NewOCPP(id string, connector int, idtag string, meterValues string, meterInterval time.Duration, initialReset core.ResetType, timeout time.Duration) (*OCPP, error) { unit := "ocpp" if id != "" { unit = id } + log := util.NewLogger(unit) + + cp := ocpp.NewChargePoint(log, id, timeout) + if err := ocpp.Instance().Register(id, cp); err != nil { + return nil, err + } c := &OCPP{ - log: util.NewLogger(unit), + log: log, cp: cp, connector: connector, idtag: idtag, } - if err := cp.Boot(); err != nil { - return nil, err + c.log.DEBUG.Printf("waiting for chargepoint: %v", timeout) + + select { + case <-time.After(timeout): + return nil, api.ErrTimeout + case <-cp.HasConnected(): } + // see who's there + ocpp.Instance().TriggerMessageRequest(cp.ID(), core.BootNotificationFeatureName) + var ( rc = make(chan error, 1) options []core.ConfigurationKey @@ -118,7 +128,7 @@ func NewOCPP(id string, connector int, idtag string, meterValues string, meterIn ) // configured id may be empty, use registered id below - err = ocpp.Instance().GetConfiguration(cp.ID(), func(resp *core.GetConfigurationConfirmation, err error) { + err := ocpp.Instance().GetConfiguration(cp.ID(), func(resp *core.GetConfigurationConfirmation, err error) { if err != nil { rc <- err return diff --git a/charger/ocpp/cp.go b/charger/ocpp/cp.go index 30503c975..6d38b9d15 100644 --- a/charger/ocpp/cp.go +++ b/charger/ocpp/cp.go @@ -13,8 +13,6 @@ import ( "github.com/lorenzodonini/ocpp-go/ocpp1.6/types" ) -const timeout = 2 * time.Minute - // Meter Profile Key const ( KeyMeterValuesSampledData = "MeterValuesSampledData" @@ -48,13 +46,15 @@ type smartChargingProfile struct { } type CP struct { - mu sync.Mutex - log *util.Logger - id string + mu sync.Mutex + log *util.Logger + once sync.Once - bootC, statusC chan struct{} - updated time.Time - status *core.StatusNotificationRequest + id string + + connectC, statusC chan struct{} + updated time.Time + status *core.StatusNotificationRequest timeout time.Duration meterUpdated time.Time @@ -67,6 +67,17 @@ type CP struct { txnId int } +func NewChargePoint(log *util.Logger, id string, timeout time.Duration) *CP { + return &CP{ + log: log, + id: id, + connectC: make(chan struct{}), + statusC: make(chan struct{}), + measurements: make(map[string]types.SampledValue), + timeout: timeout, + } +} + func (cp *CP) ID() string { cp.mu.Lock() defer cp.mu.Unlock() @@ -85,6 +96,19 @@ func (cp *CP) RegisterID(id string) { cp.id = id } +func (cp *CP) Connect() { + cp.mu.Lock() + defer cp.mu.Unlock() + + 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) @@ -217,18 +241,6 @@ func parseIntOption(key SmartchargingChargeProfileKey, options map[string]core.C return val, nil } -// Boot waits for the CP to register itself -func (cp *CP) Boot() error { - cp.log.DEBUG.Printf("waiting for chargepoint: %v", timeout) - - select { - case <-cp.bootC: - return nil - case <-time.After(timeout): - return api.ErrTimeout - } -} - // TransactionID returns the current transaction id func (cp *CP) TransactionID() int { cp.mu.Lock() @@ -242,7 +254,7 @@ func (cp *CP) Status() (api.ChargeStatus, error) { res := api.StatusNone - if time.Since(cp.updated) > timeout { + if time.Since(cp.updated) > cp.timeout { return res, api.ErrTimeout } diff --git a/charger/ocpp/cp_core.go b/charger/ocpp/cp_core.go index 56b6cdea4..925a0e618 100644 --- a/charger/ocpp/cp_core.go +++ b/charger/ocpp/cp_core.go @@ -29,16 +29,6 @@ func (cp *CP) Authorize(request *core.AuthorizeRequest) (*core.AuthorizeConfirma func (cp *CP) BootNotification(request *core.BootNotificationRequest) (*core.BootNotificationConfirmation, error) { cp.log.TRACE.Printf("%T: %+v", request, request) - if request != nil { - cp.log.DEBUG.Printf("chargepoint: %+v", request) - } - - // signal boot - select { - case cp.bootC <- struct{}{}: - default: - } - res := &core.BootNotificationConfirmation{ CurrentTime: types.NewDateTime(time.Now()), Interval: 60, // TODO diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go index c591e8fd3..b0a5ade7c 100644 --- a/charger/ocpp/cs.go +++ b/charger/ocpp/cs.go @@ -7,7 +7,6 @@ import ( "github.com/evcc-io/evcc/util" ocpp16 "github.com/lorenzodonini/ocpp-go/ocpp1.6" - "github.com/lorenzodonini/ocpp-go/ocpp1.6/types" ) type CS struct { @@ -17,30 +16,17 @@ type CS struct { cps map[string]*CP } -func (cs *CS) Register(id string) (*CP, error) { - unit := "ocpp" - if id != "" { - unit = id - } - - cp := &CP{ - id: id, - log: util.NewLogger(unit), - bootC: make(chan struct{}), - statusC: make(chan struct{}), - measurements: make(map[string]types.SampledValue), - } - +func (cs *CS) Register(id string, cp *CP) error { cs.mu.Lock() defer cs.mu.Unlock() if _, ok := cs.cps[id]; ok && id == "" { - return nil, errors.New("cannot have >1 chargepoint with empty station id") + return errors.New("cannot have >1 chargepoint with empty station id") } cs.cps[id] = cp - return cp, nil + return nil } // errorHandler logs error channel @@ -62,7 +48,7 @@ func (cs *CS) NewChargePoint(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 { if cp, ok := cs.cps[""]; ok { cs.log.INFO.Printf("chargepoint connected, registering: %s", chargePoint.ID()) @@ -71,12 +57,15 @@ func (cs *CS) NewChargePoint(chargePoint ocpp16.ChargePointConnection) { cs.cps[chargePoint.ID()] = cp delete(cs.cps, "") + cp.Connect() + return } cs.log.WARN.Printf("chargepoint connected, ignoring: %s", chargePoint.ID()) } else { cs.log.DEBUG.Printf("chargepoint connected: %s", chargePoint.ID()) + cp.Connect() } } diff --git a/cmd/ocpp/main.go b/cmd/ocpp/main.go index bc51dc65c..665fbd779 100644 --- a/cmd/ocpp/main.go +++ b/cmd/ocpp/main.go @@ -4,10 +4,17 @@ import ( "fmt" "log" + "github.com/gorilla/websocket" ocpp16 "github.com/lorenzodonini/ocpp-go/ocpp1.6" "github.com/lorenzodonini/ocpp-go/ocpp1.6/core" + "github.com/lorenzodonini/ocpp-go/ocpp1.6/firmware" + "github.com/lorenzodonini/ocpp-go/ocpp1.6/localauth" "github.com/lorenzodonini/ocpp-go/ocpp1.6/remotetrigger" + "github.com/lorenzodonini/ocpp-go/ocpp1.6/reservation" + "github.com/lorenzodonini/ocpp-go/ocpp1.6/smartcharging" "github.com/lorenzodonini/ocpp-go/ocpp1.6/types" + "github.com/lorenzodonini/ocpp-go/ocppj" + "github.com/lorenzodonini/ocpp-go/ws" ) const ( @@ -16,9 +23,35 @@ const ( ) func main() { - chargePoint := ocpp16.NewChargePoint(chargePointId, nil, nil) + // chargePoint := ocpp16.NewChargePoint(chargePointId, nil, nil) - // Set a handler for all callback functions + // create websocket client + client := ws.NewClient() + client.AddOption(func(dialer *websocket.Dialer) { + // Look for v1.6 subprotocol and add it, if not found + alreadyExists := false + for _, proto := range dialer.Subprotocols { + if proto == types.V16Subprotocol { + alreadyExists = true + break + } + } + if !alreadyExists { + dialer.Subprotocols = append(dialer.Subprotocols, types.V16Subprotocol) + } + }) + + // create chargepoint with connection tracking + endpoint := ocppj.NewClient(chargePointId, client, nil, nil, core.Profile, localauth.Profile, firmware.Profile, reservation.Profile, remotetrigger.Profile, smartcharging.Profile) + endpoint.SetOnReconnectedHandler(func() { + fmt.Println("reconnect") + }) + endpoint.SetOnDisconnectedHandler(func(err error) { + fmt.Println("disconnect") + }) + chargePoint := ocpp16.NewChargePoint(chargePointId, endpoint, client) + + // set a handler for all callback functions handler := &ChargePointHandler{triggerC: make(chan remotetrigger.MessageTrigger, 1)} chargePoint.SetCoreHandler(handler) chargePoint.SetRemoteTriggerHandler(handler) @@ -33,6 +66,13 @@ func main() { for msg := range handler.triggerC { fmt.Println("msg:", msg) switch msg { + case core.BootNotificationFeatureName: + if res, err := chargePoint.BootNotification("demo", "evcc"); err != nil { + log.Println("BootNotification:", err) + } else { + log.Println("BootNotification:", res) + } + case core.StatusNotificationFeatureName: if res, err := chargePoint.StatusNotification(1, core.NoError, core.ChargePointStatusAvailable); err != nil { log.Println("StatusNotification:", err) @@ -41,12 +81,14 @@ func main() { } case core.MeterValuesFeatureName: - if _, err := chargePoint.MeterValues(1, []types.MeterValue{ + if res, err := chargePoint.MeterValues(1, []types.MeterValue{ {SampledValue: []types.SampledValue{ {Measurand: types.MeasurandPowerActiveImport, Value: "1000"}, }}, }); err != nil { log.Println("MeterValues:", err) + } else { + log.Println("MeterValues:", res) } } } @@ -58,11 +100,6 @@ func main() { } 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 {} }