Ocpp: serialise setup (#16262)
This commit is contained in:
parent
92e4f46e11
commit
df353c6244
3 changed files with 60 additions and 32 deletions
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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()),
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue