From df353c6244ff5fd881e5cca629f27eac4f011d90 Mon Sep 17 00:00:00 2001 From: andig Date: Sun, 22 Sep 2024 18:26:53 +0200 Subject: [PATCH] Ocpp: serialise setup (#16262) --- charger/ocpp.go | 42 ++++++++++++++-------------------- charger/ocpp/cs.go | 49 ++++++++++++++++++++++++++++++++++------ charger/ocpp/instance.go | 1 + 3 files changed, 60 insertions(+), 32 deletions(-) diff --git a/charger/ocpp.go b/charger/ocpp.go index e2a333973..84b42c6ad 100644 --- a/charger/ocpp.go +++ b/charger/ocpp.go @@ -120,34 +120,26 @@ func NewOCPP(id string, connector int, idTag string, stackLevelZero, remoteStart bool, connectTimeout time.Duration, ) (*OCPP, error) { - unit := "ocpp" - if id != "" { - unit = id - } - unit = fmt.Sprintf("%s-%d", unit, connector) + log := util.NewLogger(fmt.Sprintf("%s-%d", lo.CoalesceOrEmpty(id, "ocpp"), connector)) - log := util.NewLogger(unit) + cp, err := ocpp.Instance().RegisterChargepoint(id, + func() *ocpp.CP { + return ocpp.NewChargePoint(log, id) + }, + func(cp *ocpp.CP) error { + log.DEBUG.Printf("waiting for chargepoint: %v", connectTimeout) - cp, err := ocpp.Instance().ChargepointByID(id) + select { + case <-time.After(connectTimeout): + return api.ErrTimeout + case <-cp.HasConnected(): + } + + return cp.Setup(meterValues, meterInterval) + }, + ) if err != nil { - cp = ocpp.NewChargePoint(log, id) - - // should not error - if err := ocpp.Instance().Register(id, cp); err != nil { - return nil, err - } - - log.DEBUG.Printf("waiting for chargepoint: %v", connectTimeout) - - select { - case <-time.After(connectTimeout): - return nil, api.ErrTimeout - case <-cp.HasConnected(): - } - - if err := cp.Setup(meterValues, meterInterval); err != nil { - return nil, err - } + return nil, err } if cp.NumberOfConnectors > 0 && connector > cp.NumberOfConnectors { diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go index 72be4ac00..565c63f30 100644 --- a/charger/ocpp/cs.go +++ b/charger/ocpp/cs.go @@ -14,26 +14,35 @@ type CS struct { log *util.Logger ocpp16.CentralSystem cps map[string]*CP + init map[string]*sync.Mutex txnId int } // Register registers a charge point with the central system. // The charge point identified by id may already be connected in which case initial connection is triggered. -func (cs *CS) Register(id string, cp *CP) error { +func (cs *CS) register(id string, new *CP) error { cs.mu.Lock() defer cs.mu.Unlock() - if _, ok := cs.cps[id]; ok && id == "" { + cp, ok := cs.cps[id] + + // case 1: charge point neither registered nor physically connected + if !ok { + cs.cps[id] = new + return nil + } + + // case 2: duplicate registration of id empty + if id == "" { return errors.New("cannot have >1 charge point with empty station id") } - // trigger unknown charge point connected - if unknown, ok := cs.cps[id]; ok && unknown == nil { - cp.connect(true) + // case 3: charge point not registered but physically already connected + if cp == nil { + cs.cps[id] = new + new.connect(true) } - cs.cps[id] = cp - return nil } @@ -58,6 +67,32 @@ func (cs *CS) ChargepointByID(id string) (*CP, error) { return cp, nil } +func (cs *CS) RegisterChargepoint(id string, newfun func() *CP, init func(*CP) error) (*CP, error) { + cs.mu.Lock() + cpmu, ok := cs.init[id] + if !ok { + cpmu = new(sync.Mutex) + cs.init[id] = cpmu + } + cs.mu.Unlock() + + // serialise on chargepoint id + cpmu.Lock() + defer cpmu.Unlock() + + cp, err := cs.ChargepointByID(id) + if err != nil { + cp = newfun() + } + + // should not error + if err := cs.register(id, cp); err != nil { + return nil, err + } + + return cp, init(cp) +} + // NewChargePoint implements ocpp16.ChargePointConnectionHandler func (cs *CS) NewChargePoint(chargePoint ocpp16.ChargePointConnection) { cs.mu.Lock() diff --git a/charger/ocpp/instance.go b/charger/ocpp/instance.go index ec6315f13..64086802a 100644 --- a/charger/ocpp/instance.go +++ b/charger/ocpp/instance.go @@ -41,6 +41,7 @@ func Instance() *CS { instance = &CS{ log: log, cps: make(map[string]*CP), + init: make(map[string]*sync.Mutex), CentralSystem: cs, txnId: int(time.Now().UTC().Unix()), }