From cff845195f8dcbe0f2ac2a235b7eef61feee7299 Mon Sep 17 00:00:00 2001 From: andig Date: Fri, 6 Oct 2023 14:00:11 +0200 Subject: [PATCH] Ocpp: support multiple connectors (#10187) --- charger/ocpp.go | 71 +++++---- charger/ocpp/connector.go | 260 ++++++++++++++++++++++++++++++ charger/ocpp/connector_core.go | 124 +++++++++++++++ charger/ocpp/cp.go | 280 ++++++--------------------------- charger/ocpp/cp_core.go | 138 +++++----------- charger/ocpp/cs.go | 5 +- charger/ocpp/cs_core.go | 26 ++- charger/ocpp_test.go | 84 +++++----- 8 files changed, 576 insertions(+), 412 deletions(-) create mode 100644 charger/ocpp/connector.go create mode 100644 charger/ocpp/connector_core.go diff --git a/charger/ocpp.go b/charger/ocpp.go index 27601b6a9..74f34cd26 100644 --- a/charger/ocpp.go +++ b/charger/ocpp.go @@ -22,8 +22,7 @@ import ( // OCPP charger implementation type OCPP struct { log *util.Logger - cp *ocpp.CP - connector int + conn *ocpp.Connector idtag string enabled bool phases int @@ -113,19 +112,30 @@ func NewOCPP(id string, connector int, idtag string, if id != "" { unit = id } + unit = fmt.Sprintf("%s-%d", unit, connector) + log := util.NewLogger(unit) - cp := ocpp.NewChargePoint(log, id, connector, timeout) - if err := ocpp.Instance().Register(id, cp); err != nil { + cp, err := ocpp.Instance().ChargepointByID(id) + if err != nil { + cp = ocpp.NewChargePoint(log, id) + + // should not error + if err := ocpp.Instance().Register(id, cp); err != nil { + return nil, err + } + } + + conn, err := ocpp.NewConnector(log, connector, cp, timeout) + if err != nil { return nil, err } c := &OCPP{ - log: log, - cp: cp, - connector: connector, - idtag: idtag, - timeout: timeout, + log: log, + conn: conn, + idtag: idtag, + timeout: timeout, } c.log.DEBUG.Printf("waiting for chargepoint: %v", connectTimeout) @@ -138,7 +148,7 @@ func NewOCPP(id string, connector int, idtag string, // see who's there if boot { - ocpp.Instance().TriggerMessageRequest(cp.ID(), core.BootNotificationFeatureName) + conn.TriggerMessageRequest(core.BootNotificationFeatureName) } var ( @@ -188,8 +198,8 @@ func NewOCPP(id string, connector int, idtag string, switch opt.Key { case ocpp.KeyNumberOfConnectors: var val int - if val, err = strconv.Atoi(*opt.Value); err == nil && c.connector > val { - err = fmt.Errorf("connector %d exceeds max available connectors: %d", c.connector, val) + if val, err = strconv.Atoi(*opt.Value); err == nil && connector > val { + err = fmt.Errorf("connector %d exceeds max available connectors: %d", connector, val) } case ocpp.KeyMeterValuesSampledData: @@ -244,7 +254,7 @@ func NewOCPP(id string, connector int, idtag string, // get initial meter values and configure sample rate if c.hasMeasurement(types.MeasurandPowerActiveImport) || c.hasMeasurement(types.MeasurandEnergyActiveImportRegister) { - ocpp.Instance().TriggerMeterValuesRequest(cp.ID(), cp.Connector()) + conn.TriggerMessageRequest(core.MeterValuesFeatureName) if meterInterval > 0 && meterInterval != meterSampleInterval { if err := c.configure(ocpp.KeyMeterValueSampleInterval, strconv.Itoa(int(meterInterval.Seconds()))); err != nil { @@ -255,13 +265,13 @@ func NewOCPP(id string, connector int, idtag string, // HACK: setup watchdog for meter values if not happy with config if meterInterval > 0 { c.log.DEBUG.Println("enabling meter watchdog") - go cp.WatchDog(meterInterval) + go conn.WatchDog(meterInterval) } } // TODO: check for running transaction - return c, cp.Initialized() + return c, conn.Initialized() } // hasMeasurement checks if meterValuesSample contains given measurement @@ -273,7 +283,7 @@ func (c *OCPP) hasMeasurement(val types.Measurand) bool { func (c *OCPP) configure(key, val string) error { rc := make(chan error, 1) - err := ocpp.Instance().ChangeConfiguration(c.cp.ID(), func(resp *core.ChangeConfigurationConfirmation, err error) { + err := ocpp.Instance().ChangeConfiguration(c.conn.ChargePoint().ID(), func(resp *core.ChangeConfigurationConfirmation, err error) { if err == nil && resp != nil && resp.Status != core.ConfigurationStatusAccepted { rc <- fmt.Errorf("ChangeConfiguration failed: %s", resp.Status) } @@ -299,7 +309,7 @@ func (c *OCPP) wait(err error, rc chan error) error { // Status implements the api.Charger interface func (c *OCPP) Status() (api.ChargeStatus, error) { - return c.cp.Status() + return c.conn.Status() } // Enabled implements the api.Charger interface @@ -310,7 +320,7 @@ func (c *OCPP) Enabled() (bool, error) { // Enable implements the api.Charger interface func (c *OCPP) Enable(enable bool) (err error) { rc := make(chan error, 1) - txn, err := c.cp.TransactionID() + txn, err := c.conn.TransactionID() defer func() { if err == nil { @@ -323,14 +333,15 @@ func (c *OCPP) Enable(enable bool) (err error) { return errors.New("cannot enable: transaction already running") } - err = ocpp.Instance().RemoteStartTransaction(c.cp.ID(), func(resp *core.RemoteStartTransactionConfirmation, err error) { + err = ocpp.Instance().RemoteStartTransaction(c.conn.ChargePoint().ID(), func(resp *core.RemoteStartTransactionConfirmation, err error) { if err == nil && resp != nil && resp.Status != types.RemoteStartStopStatusAccepted { err = errors.New(string(resp.Status)) } rc <- err }, c.idtag, func(request *core.RemoteStartTransactionRequest) { - request.ConnectorId = &c.connector + connector := c.conn.ID() + request.ConnectorId = &connector request.ChargingProfile = c.getTxChargingProfile(c.current, 0) }) } else { @@ -348,7 +359,7 @@ func (c *OCPP) Enable(enable bool) (err error) { return nil } - err = ocpp.Instance().RemoteStopTransaction(c.cp.ID(), func(resp *core.RemoteStopTransactionConfirmation, err error) { + err = ocpp.Instance().RemoteStopTransaction(c.conn.ChargePoint().ID(), func(resp *core.RemoteStopTransactionConfirmation, err error) { if err == nil && resp != nil && resp.Status != types.RemoteStartStopStatusAccepted { err = errors.New(string(resp.Status)) } @@ -360,15 +371,17 @@ func (c *OCPP) Enable(enable bool) (err error) { return c.wait(err, rc) } -func (c *OCPP) setChargingProfile(connectorId int, profile *types.ChargingProfile) error { +func (c *OCPP) setChargingProfile(profile *types.ChargingProfile) error { + connector := c.conn.ID() + rc := make(chan error, 1) - err := ocpp.Instance().SetChargingProfile(c.cp.ID(), func(resp *smartcharging.SetChargingProfileConfirmation, err error) { + err := ocpp.Instance().SetChargingProfile(c.conn.ChargePoint().ID(), func(resp *smartcharging.SetChargingProfileConfirmation, err error) { if err == nil && resp != nil && resp.Status != smartcharging.ChargingProfileStatusAccepted { err = errors.New(string(resp.Status)) } rc <- err - }, connectorId, profile) + }, connector, profile) return c.wait(err, rc) } @@ -380,14 +393,14 @@ func (c *OCPP) updatePeriod(current float64) error { return err } - txn, err := c.cp.TransactionID() + txn, err := c.conn.TransactionID() if err != nil { return err } current = math.Trunc(10*current) / 10 - err = c.setChargingProfile(c.connector, c.getTxChargingProfile(current, txn)) + err = c.setChargingProfile(c.getTxChargingProfile(current, txn)) if err != nil { err = fmt.Errorf("set charging profile: %w", err) } @@ -445,17 +458,17 @@ func (c *OCPP) MaxCurrentMillis(current float64) error { // CurrentPower implements the api.Meter interface func (c *OCPP) currentPower() (float64, error) { - return c.cp.CurrentPower() + return c.conn.CurrentPower() } // TotalEnergy implements the api.MeterTotal interface func (c *OCPP) totalEnergy() (float64, error) { - return c.cp.TotalEnergy() + return c.conn.TotalEnergy() } // Currents implements the api.PhaseCurrents interface func (c *OCPP) currents() (float64, float64, float64, error) { - return c.cp.Currents() + return c.conn.Currents() } // Phases1p3p implements the api.PhaseSwitcher interface diff --git a/charger/ocpp/connector.go b/charger/ocpp/connector.go new file mode 100644 index 000000000..fb5d45532 --- /dev/null +++ b/charger/ocpp/connector.go @@ -0,0 +1,260 @@ +package ocpp + +import ( + "fmt" + "strconv" + "strings" + "sync" + "time" + + "github.com/benbjohnson/clock" + "github.com/evcc-io/evcc/api" + "github.com/evcc-io/evcc/util" + "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 Connector struct { + log *util.Logger + mu sync.Mutex + clock clock.Clock // mockable time + cp *CP + id int + + status *core.StatusNotificationRequest + statusC chan struct{} + + meterUpdated time.Time + measurements map[string]types.SampledValue + timeout time.Duration + + txnCount int // change initial value to the last known global transaction. Needs persistence + txnId int +} + +func NewConnector(log *util.Logger, id int, cp *CP, timeout time.Duration) (*Connector, error) { + conn := &Connector{ + log: log, + cp: cp, + id: id, + clock: clock.New(), + statusC: make(chan struct{}), + measurements: make(map[string]types.SampledValue), + timeout: timeout, + } + + err := cp.registerConnector(id, conn) + + return conn, err +} + +func (conn *Connector) TestClock(clock clock.Clock) { + conn.clock = clock +} + +func (conn *Connector) ChargePoint() *CP { + return conn.cp +} + +func (conn *Connector) ID() int { + return conn.id +} + +func (conn *Connector) TriggerMessageRequest(feature remotetrigger.MessageTrigger, f ...func(request *remotetrigger.TriggerMessageRequest)) { + Instance().TriggerMessageRequest(conn.cp.ID(), feature, func(request *remotetrigger.TriggerMessageRequest) { + request.ConnectorId = &conn.id + for _, f := range f { + f(request) + } + }) +} + +// WatchDog triggers meter values messages if older than timeout. +// Must be wrapped in a goroutine. +func (conn *Connector) WatchDog(timeout time.Duration) { + for ; true; <-time.Tick(timeout) { + conn.mu.Lock() + update := conn.txnId != 0 && conn.clock.Since(conn.meterUpdated) > timeout + conn.mu.Unlock() + + if update { + conn.TriggerMessageRequest(core.MeterValuesFeatureName) + } + } +} + +// Initialized waits for initial charge point status notification +func (conn *Connector) Initialized() error { + trigger := time.After(conn.timeout / 2) + timeout := time.After(conn.timeout) + for { + select { + case <-conn.statusC: + return nil + + case <-trigger: + conn.TriggerMessageRequest(core.StatusNotificationFeatureName) + + case <-timeout: + return api.ErrTimeout + } + } +} + +// TransactionID returns the current transaction id +func (conn *Connector) TransactionID() (int, error) { + if !conn.cp.Connected() { + return 0, api.ErrTimeout + } + + conn.mu.Lock() + defer conn.mu.Unlock() + + return conn.txnId, nil +} + +func (conn *Connector) Status() (api.ChargeStatus, error) { + conn.mu.Lock() + defer conn.mu.Unlock() + + res := api.StatusNone + + if !conn.cp.Connected() { + return res, api.ErrTimeout + } + + if conn.status.ErrorCode != core.NoError { + return res, fmt.Errorf("%s: %s", conn.status.ErrorCode, conn.status.Info) + } + + switch conn.status.Status { + case core.ChargePointStatusAvailable, // "Available" + core.ChargePointStatusUnavailable: // "Unavailable" + res = api.StatusA + case + core.ChargePointStatusPreparing, // "Preparing" + core.ChargePointStatusSuspendedEVSE, // "SuspendedEVSE" + core.ChargePointStatusSuspendedEV, // "SuspendedEV" + core.ChargePointStatusFinishing: // "Finishing" + res = api.StatusB + case core.ChargePointStatusCharging: // "Charging" + res = api.StatusC + case core.ChargePointStatusReserved, // "Reserved" + core.ChargePointStatusFaulted: // "Faulted" + return api.StatusF, fmt.Errorf("chargepoint status: %s", conn.status.ErrorCode) + default: + return api.StatusNone, fmt.Errorf("invalid chargepoint status: %s", conn.status.Status) + } + + return res, nil +} + +// isMeterTimeout checks if meter values are outdated. +// Must only be called while holding lock. +func (conn *Connector) isMeterTimeout() bool { + return conn.timeout > 0 && conn.clock.Since(conn.meterUpdated) > conn.timeout +} + +var _ api.Meter = (*Connector)(nil) + +func (conn *Connector) CurrentPower() (float64, error) { + if !conn.cp.Connected() { + return 0, api.ErrTimeout + } + + conn.mu.Lock() + defer conn.mu.Unlock() + + // zero value on timeout when not charging + if conn.isMeterTimeout() { + if conn.txnId != 0 { + return 0, api.ErrTimeout + } + + return 0, nil + } + + if m, ok := conn.measurements[string(types.MeasurandPowerActiveImport)]; ok { + f, err := strconv.ParseFloat(m.Value, 64) + return scale(f, m.Unit), err + } + + return 0, api.ErrNotAvailable +} + +var _ api.MeterEnergy = (*Connector)(nil) + +func (conn *Connector) TotalEnergy() (float64, error) { + if !conn.cp.Connected() { + return 0, api.ErrTimeout + } + + conn.mu.Lock() + defer conn.mu.Unlock() + + // fallthrough for last value on timeout when not charging + if conn.txnId != 0 && conn.isMeterTimeout() { + return 0, api.ErrTimeout + } + + if m, ok := conn.measurements[string(types.MeasurandEnergyActiveImportRegister)]; ok { + f, err := strconv.ParseFloat(m.Value, 64) + return scale(f, m.Unit) / 1e3, err + } + + return 0, api.ErrNotAvailable +} + +func scale(f float64, scale types.UnitOfMeasure) float64 { + switch { + case strings.HasPrefix(string(scale), "k"): + return f * 1e3 + case strings.HasPrefix(string(scale), "m"): + return f / 1e3 + default: + return f + } +} + +func getKeyCurrentPhase(phase int) string { + return string(types.MeasurandCurrentImport) + "@L" + strconv.Itoa(phase) +} + +var _ api.PhaseCurrents = (*Connector)(nil) + +func (conn *Connector) Currents() (float64, float64, float64, error) { + if !conn.cp.Connected() { + return 0, 0, 0, api.ErrTimeout + } + + conn.mu.Lock() + defer conn.mu.Unlock() + + // zero value on timeout when not charging + if conn.isMeterTimeout() { + if conn.txnId != 0 { + return 0, 0, 0, api.ErrTimeout + } + + return 0, 0, 0, nil + } + + currents := make([]float64, 0, 3) + + for phase := 1; phase <= 3; phase++ { + m, ok := conn.measurements[getKeyCurrentPhase(phase)] + if !ok { + return 0, 0, 0, api.ErrNotAvailable + } + + f, err := strconv.ParseFloat(m.Value, 64) + if err != nil { + return 0, 0, 0, fmt.Errorf("invalid current for phase %d: %w", phase, err) + } + + currents = append(currents, scale(f, m.Unit)) + } + + return currents[0], currents[1], currents[2], nil +} diff --git a/charger/ocpp/connector_core.go b/charger/ocpp/connector_core.go new file mode 100644 index 000000000..70796fd27 --- /dev/null +++ b/charger/ocpp/connector_core.go @@ -0,0 +1,124 @@ +package ocpp + +import ( + "time" + + "github.com/lorenzodonini/ocpp-go/ocpp1.6/core" + "github.com/lorenzodonini/ocpp-go/ocpp1.6/types" +) + +// timestampValid returns false if status timestamps are outdated +func (conn *Connector) timestampValid(t time.Time) bool { + // reject if expired + if conn.clock.Since(t) > messageExpiry { + return false + } + + // assume having a timestamp is better than not + if conn.status.Timestamp == nil { + return true + } + + // reject older values than we already have + return t.After(conn.status.Timestamp.Time) +} + +func (conn *Connector) StatusNotification(request *core.StatusNotificationRequest) (*core.StatusNotificationConfirmation, error) { + conn.mu.Lock() + defer conn.mu.Unlock() + + if conn.status == nil { + conn.status = request + close(conn.statusC) // signal initial status received + } else if request.Timestamp == nil || conn.timestampValid(request.Timestamp.Time) { + conn.status = request + } else { + conn.log.TRACE.Printf("ignoring status: %s < %s", request.Timestamp.Time, conn.status.Timestamp) + } + + return new(core.StatusNotificationConfirmation), nil +} + +func (conn *Connector) MeterValues(request *core.MeterValuesRequest) (*core.MeterValuesConfirmation, error) { + conn.mu.Lock() + defer conn.mu.Unlock() + + if request.TransactionId != nil && conn.txnId == 0 { + conn.log.DEBUG.Printf("hijacking transaction: %d", *request.TransactionId) + conn.txnId = *request.TransactionId + } + + for _, meterValue := range request.MeterValue { + // ignore old meter value requests + if meterValue.Timestamp.Time.After(conn.meterUpdated) { + for _, sample := range meterValue.SampledValue { + conn.measurements[getSampleKey(sample)] = sample + conn.meterUpdated = conn.clock.Now() + } + } + } + + return new(core.MeterValuesConfirmation), nil +} + +func getSampleKey(s types.SampledValue) string { + if s.Phase != "" { + return string(s.Measurand) + "@" + string(s.Phase) + } + + return string(s.Measurand) +} + +func (conn *Connector) StartTransaction(request *core.StartTransactionRequest) (*core.StartTransactionConfirmation, error) { + conn.mu.Lock() + defer conn.mu.Unlock() + + // expired request + if request.Timestamp != nil && conn.clock.Since(request.Timestamp.Time) < transactionExpiry { + res := &core.StartTransactionConfirmation{ + IdTagInfo: &types.IdTagInfo{ + Status: types.AuthorizationStatusExpired, + }, + } + + return res, nil + } + + conn.txnCount++ + conn.txnId = conn.txnCount + + res := &core.StartTransactionConfirmation{ + IdTagInfo: &types.IdTagInfo{ + Status: types.AuthorizationStatusAccepted, + }, + TransactionId: conn.txnId, + } + + return res, nil +} + +func (conn *Connector) StopTransaction(request *core.StopTransactionRequest) (*core.StopTransactionConfirmation, error) { + conn.mu.Lock() + defer conn.mu.Unlock() + + // expired request + if request.Timestamp != nil && conn.clock.Since(request.Timestamp.Time) > transactionExpiry { + res := &core.StopTransactionConfirmation{ + IdTagInfo: &types.IdTagInfo{ + Status: types.AuthorizationStatusExpired, // accept + }, + } + + return res, nil + } + + conn.txnId = 0 + + res := &core.StopTransactionConfirmation{ + IdTagInfo: &types.IdTagInfo{ + Status: types.AuthorizationStatusAccepted, // accept + }, + } + + return res, nil +} diff --git a/charger/ocpp/cp.go b/charger/ocpp/cp.go index 14b6b891c..b51ea983a 100644 --- a/charger/ocpp/cp.go +++ b/charger/ocpp/cp.go @@ -2,59 +2,67 @@ package ocpp import ( "fmt" - "strconv" - "strings" "sync" - "time" - "github.com/benbjohnson/clock" - "github.com/evcc-io/evcc/api" "github.com/evcc-io/evcc/util" - "github.com/lorenzodonini/ocpp-go/ocpp1.6/core" - "github.com/lorenzodonini/ocpp-go/ocpp1.6/remotetrigger" - "github.com/lorenzodonini/ocpp-go/ocpp1.6/types" ) // TODO support multiple connectors // Since ocpp-go interfaces at charge point level, we need to manage multiple connector separately type CP struct { - mu sync.Mutex - once sync.Once - clock clock.Clock // mockable time - log *util.Logger + mu sync.Mutex + once sync.Once + log *util.Logger - id string - connector int + id string - connectC, statusC chan struct{} - connected bool - status *core.StatusNotificationRequest + connected bool + connectC chan struct{} - meterUpdated time.Time - timeout time.Duration - - measurements map[string]types.SampledValue - - txnCount int // change initial value to the last known global transaction. Needs persistence - txnId int + connectors map[int]*Connector } -func NewChargePoint(log *util.Logger, id string, connector int, timeout time.Duration) *CP { +func NewChargePoint(log *util.Logger, id string) *CP { return &CP{ - clock: clock.New(), - log: log, - id: id, - connector: connector, - connectC: make(chan struct{}), - statusC: make(chan struct{}), - measurements: make(map[string]types.SampledValue), - timeout: timeout, + log: log, + id: id, + + connectC: make(chan struct{}), + connectors: make(map[int]*Connector), } } -func (cp *CP) TestClock(clock clock.Clock) { - cp.clock = clock +func (cp *CP) registerConnector(id int, conn *Connector) error { + cp.mu.Lock() + defer cp.mu.Unlock() + + if _, ok := cp.connectors[id]; ok { + return fmt.Errorf("connector already registered: %d", id) + } + + cp.connectors[id] = conn + return nil +} + +func (cp *CP) connectorByID(id int) *Connector { + cp.mu.Lock() + defer cp.mu.Unlock() + + return cp.connectors[id] +} + +func (cp *CP) connectorByTransactionID(id int) *Connector { + cp.mu.Lock() + defer cp.mu.Unlock() + + for _, conn := range cp.connectors { + if txn, err := conn.TransactionID(); err == nil && txn == id { + return conn + } + } + + return nil } func (cp *CP) ID() string { @@ -75,10 +83,6 @@ func (cp *CP) RegisterID(id string) { cp.id = id } -func (cp *CP) Connector() int { - return cp.connector -} - func (cp *CP) connect(connect bool) { cp.mu.Lock() defer cp.mu.Unlock() @@ -92,197 +96,13 @@ func (cp *CP) connect(connect bool) { } } +func (cp *CP) Connected() bool { + cp.mu.Lock() + defer cp.mu.Unlock() + + return cp.connected +} + func (cp *CP) HasConnected() <-chan struct{} { return cp.connectC } - -func (cp *CP) Initialized() error { - // trigger status - time.AfterFunc(cp.timeout/2, func() { - select { - case <-cp.statusC: - return - default: - Instance().TriggerMessageRequest(cp.ID(), core.StatusNotificationFeatureName, func(request *remotetrigger.TriggerMessageRequest) { - request.ConnectorId = &cp.connector - }) - } - }) - - // wait for status - select { - case <-cp.statusC: - return nil - case <-time.After(cp.timeout): - return api.ErrTimeout - } -} - -// TransactionID returns the current transaction id -func (cp *CP) TransactionID() (int, error) { - cp.mu.Lock() - defer cp.mu.Unlock() - - if !cp.connected { - return 0, api.ErrTimeout - } - - return cp.txnId, nil -} - -func (cp *CP) Status() (api.ChargeStatus, error) { - cp.mu.Lock() - defer cp.mu.Unlock() - - res := api.StatusNone - - if !cp.connected { - return res, api.ErrTimeout - } - - if cp.status.ErrorCode != core.NoError { - return res, fmt.Errorf("%s: %s", cp.status.ErrorCode, cp.status.Info) - } - - switch cp.status.Status { - case core.ChargePointStatusAvailable, // "Available" - core.ChargePointStatusUnavailable: // "Unavailable" - res = api.StatusA - case - core.ChargePointStatusPreparing, // "Preparing" - core.ChargePointStatusSuspendedEVSE, // "SuspendedEVSE" - core.ChargePointStatusSuspendedEV, // "SuspendedEV" - core.ChargePointStatusFinishing: // "Finishing" - res = api.StatusB - case core.ChargePointStatusCharging: // "Charging" - res = api.StatusC - case core.ChargePointStatusReserved, // "Reserved" - core.ChargePointStatusFaulted: // "Faulted" - return api.StatusF, fmt.Errorf("chargepoint status: %s", cp.status.ErrorCode) - default: - return api.StatusNone, fmt.Errorf("invalid chargepoint status: %s", cp.status.Status) - } - - return res, nil -} - -// WatchDog triggers meter values messages if older than timeout. -// Must be wrapped in a goroutine. -func (cp *CP) WatchDog(timeout time.Duration) { - for ; true; <-time.Tick(timeout) { - cp.mu.Lock() - update := cp.txnId != 0 && cp.clock.Since(cp.meterUpdated) > timeout - cp.mu.Unlock() - - if update { - Instance().TriggerMeterValuesRequest(cp.ID(), cp.Connector()) - } - } -} - -func (cp *CP) isTimeout() bool { - return cp.timeout > 0 && cp.clock.Since(cp.meterUpdated) > cp.timeout -} - -var _ api.Meter = (*CP)(nil) - -func (cp *CP) CurrentPower() (float64, error) { - cp.mu.Lock() - defer cp.mu.Unlock() - - if !cp.connected { - return 0, api.ErrTimeout - } - - // zero value on timeout when not charging - if cp.isTimeout() { - if cp.txnId != 0 { - return 0, api.ErrTimeout - } - - return 0, nil - } - - if m, ok := cp.measurements[string(types.MeasurandPowerActiveImport)]; ok { - f, err := strconv.ParseFloat(m.Value, 64) - return scale(f, m.Unit), err - } - - return 0, api.ErrNotAvailable -} - -var _ api.MeterEnergy = (*CP)(nil) - -func (cp *CP) TotalEnergy() (float64, error) { - cp.mu.Lock() - defer cp.mu.Unlock() - - if !cp.connected { - return 0, api.ErrTimeout - } - - // fallthrough for last value on timeout when not charging - if cp.txnId != 0 && cp.isTimeout() { - return 0, api.ErrTimeout - } - - if m, ok := cp.measurements[string(types.MeasurandEnergyActiveImportRegister)]; ok { - f, err := strconv.ParseFloat(m.Value, 64) - return scale(f, m.Unit) / 1e3, err - } - - return 0, api.ErrNotAvailable -} - -func scale(f float64, scale types.UnitOfMeasure) float64 { - switch { - case strings.HasPrefix(string(scale), "k"): - return f * 1e3 - case strings.HasPrefix(string(scale), "m"): - return f / 1e3 - default: - return f - } -} - -func getKeyCurrentPhase(phase int) string { - return string(types.MeasurandCurrentImport) + "@L" + strconv.Itoa(phase) -} - -var _ api.PhaseCurrents = (*CP)(nil) - -func (cp *CP) Currents() (float64, float64, float64, error) { - cp.mu.Lock() - defer cp.mu.Unlock() - - if !cp.connected { - return 0, 0, 0, api.ErrTimeout - } - - // zero value on timeout when not charging - if cp.isTimeout() { - if cp.txnId != 0 { - return 0, 0, 0, api.ErrTimeout - } - - return 0, 0, 0, nil - } - - currents := make([]float64, 0, 3) - - for phase := 1; phase <= 3; phase++ { - m, ok := cp.measurements[getKeyCurrentPhase(phase)] - if !ok { - return 0, 0, 0, api.ErrNotAvailable - } - - f, err := strconv.ParseFloat(m.Value, 64) - if err != nil { - return 0, 0, 0, fmt.Errorf("invalid current for phase %d: %w", phase, err) - } - - currents = append(currents, scale(f, m.Unit)) - } - - return currents[0], currents[1], currents[2], nil -} diff --git a/charger/ocpp/cp_core.go b/charger/ocpp/cp_core.go index 0f8bb99f1..4d405709d 100644 --- a/charger/ocpp/cp_core.go +++ b/charger/ocpp/cp_core.go @@ -1,6 +1,7 @@ package ocpp import ( + "errors" "time" "github.com/lorenzodonini/ocpp-go/ocpp1.6/core" @@ -13,6 +14,12 @@ const ( transactionExpiry = time.Hour ) +var ( + ErrInvalidRequest = errors.New("invalid request") + ErrInvalidConnector = errors.New("invalid connector") + ErrInvalidTransaction = errors.New("invalid transaction") +) + func (cp *CP) Authorize(request *core.AuthorizeRequest) (*core.AuthorizeConfirmation, error) { // TODO check if this authorizes foreign RFID tags res := &core.AuthorizeConfirmation{ @@ -26,7 +33,7 @@ func (cp *CP) Authorize(request *core.AuthorizeRequest) (*core.AuthorizeConfirma func (cp *CP) BootNotification(request *core.BootNotificationRequest) (*core.BootNotificationConfirmation, error) { res := &core.BootNotificationConfirmation{ - CurrentTime: types.NewDateTime(cp.clock.Now()), + CurrentTime: types.NewDateTime(time.Now()), Interval: 60, // TODO Status: core.RegistrationStatusAccepted, } @@ -34,38 +41,25 @@ func (cp *CP) BootNotification(request *core.BootNotificationRequest) (*core.Boo return res, nil } -// timestampValid returns false if status timestamps are outdated -func (cp *CP) timestampValid(t time.Time) bool { - // reject if expired - if time.Since(t) > messageExpiry { - return false - } +func (cp *CP) DiagnosticStatusNotification(request *firmware.DiagnosticsStatusNotificationRequest) (*firmware.DiagnosticsStatusNotificationConfirmation, error) { + return new(firmware.DiagnosticsStatusNotificationConfirmation), nil +} - // assume having a timestamp is better than not - if cp.status.Timestamp == nil { - return true - } - - // reject older values than we already have - return !t.Before(cp.status.Timestamp.Time) +func (cp *CP) FirmwareStatusNotification(request *firmware.FirmwareStatusNotificationRequest) (*firmware.FirmwareStatusNotificationConfirmation, error) { + return new(firmware.FirmwareStatusNotificationConfirmation), nil } func (cp *CP) StatusNotification(request *core.StatusNotificationRequest) (*core.StatusNotificationConfirmation, error) { - if request != nil && request.ConnectorId == cp.connector { - cp.mu.Lock() - defer cp.mu.Unlock() - - if cp.status == nil { - cp.status = request - close(cp.statusC) // signal initial status received - } else if request.Timestamp == nil || cp.timestampValid(request.Timestamp.Time) { - cp.status = request - } else { - cp.log.TRACE.Printf("ignoring status: %s < %s", request.Timestamp.Time, cp.status.Timestamp) - } + if request == nil { + return nil, ErrInvalidRequest } - return new(core.StatusNotificationConfirmation), nil + conn := cp.connectorByID(request.ConnectorId) + if conn == nil { + return nil, ErrInvalidConnector + } + + return conn.StatusNotification(request) } func (cp *CP) DataTransfer(request *core.DataTransferRequest) (*core.DataTransferConfirmation, error) { @@ -78,99 +72,47 @@ func (cp *CP) DataTransfer(request *core.DataTransferRequest) (*core.DataTransfe func (cp *CP) Heartbeat(request *core.HeartbeatRequest) (*core.HeartbeatConfirmation, error) { res := &core.HeartbeatConfirmation{ - CurrentTime: types.NewDateTime(cp.clock.Now()), + CurrentTime: types.NewDateTime(time.Now()), } return res, nil } func (cp *CP) MeterValues(request *core.MeterValuesRequest) (*core.MeterValuesConfirmation, error) { - if request != nil && request.ConnectorId == cp.connector { - cp.mu.Lock() - defer cp.mu.Unlock() - - if request.TransactionId != nil && cp.txnId == 0 { - cp.log.DEBUG.Printf("hijacking transaction: %d", *request.TransactionId) - cp.txnId = *request.TransactionId - } - - for _, meterValue := range request.MeterValue { - // ignore old meter value requests - if meterValue.Timestamp.Time.After(cp.meterUpdated) { - for _, sample := range meterValue.SampledValue { - cp.measurements[getSampleKey(sample)] = sample - cp.meterUpdated = cp.clock.Now() - } - } - } + if request == nil { + return nil, ErrInvalidRequest } - return new(core.MeterValuesConfirmation), nil -} - -func getSampleKey(s types.SampledValue) string { - if s.Phase != "" { - return string(s.Measurand) + "@" + string(s.Phase) + conn := cp.connectorByID(request.ConnectorId) + if conn == nil { + return nil, ErrInvalidConnector } - return string(s.Measurand) + return conn.MeterValues(request) } func (cp *CP) StartTransaction(request *core.StartTransactionRequest) (*core.StartTransactionConfirmation, error) { - if request == nil || request.ConnectorId != cp.connector { - return new(core.StartTransactionConfirmation), nil + if request == nil { + return nil, ErrInvalidRequest } - cp.mu.Lock() - defer cp.mu.Unlock() - - res := &core.StartTransactionConfirmation{ - IdTagInfo: &types.IdTagInfo{ - Status: types.AuthorizationStatusAccepted, // accept - }, - TransactionId: 1, // default + conn := cp.connectorByID(request.ConnectorId) + if conn == nil { + return nil, ErrInvalidConnector } - // create new transaction - if request != nil && time.Since(request.Timestamp.Time) < transactionExpiry { // only respect transactions in the last hour - cp.txnCount++ - res.TransactionId = cp.txnCount - } - - cp.txnId = res.TransactionId - - return res, nil + return conn.StartTransaction(request) } func (cp *CP) StopTransaction(request *core.StopTransactionRequest) (*core.StopTransactionConfirmation, error) { - if request != nil { - cp.mu.Lock() - defer cp.mu.Unlock() - - // reset transaction - if time.Since(request.Timestamp.Time) < transactionExpiry { // only respect transactions in the last hour - // log mismatching id but close transaction anyway - if request.TransactionId != cp.txnId { - cp.log.ERROR.Printf("stop transaction: invalid id %d", request.TransactionId) - } - - cp.txnId = 0 - } + if request == nil { + return nil, ErrInvalidRequest } - res := &core.StopTransactionConfirmation{ - IdTagInfo: &types.IdTagInfo{ - Status: types.AuthorizationStatusAccepted, // accept - }, + conn := cp.connectorByTransactionID(request.TransactionId) + if conn == nil { + return nil, ErrInvalidTransaction } - return res, nil -} - -func (cp *CP) DiagnosticStatusNotification(request *firmware.DiagnosticsStatusNotificationRequest) (*firmware.DiagnosticsStatusNotificationConfirmation, error) { - return &firmware.DiagnosticsStatusNotificationConfirmation{}, nil -} - -func (cp *CP) FirmwareStatusNotification(request *firmware.FirmwareStatusNotificationRequest) (*firmware.FirmwareStatusNotificationConfirmation, error) { - return &firmware.FirmwareStatusNotificationConfirmation{}, nil + return conn.StopTransaction(request) } diff --git a/charger/ocpp/cs.go b/charger/ocpp/cs.go index 06d2e5bec..3f5929fd9 100644 --- a/charger/ocpp/cs.go +++ b/charger/ocpp/cs.go @@ -43,8 +43,7 @@ func (cs *CS) errorHandler(errC <-chan error) { } } -// chargepointByID returns a configured charge point identified by id. -func (cs *CS) chargepointByID(id string) (*CP, error) { +func (cs *CS) ChargepointByID(id string) (*CP, error) { cp, ok := cs.cps[id] if !ok { return nil, fmt.Errorf("unknown charge point: %s", id) @@ -103,7 +102,7 @@ func (cs *CS) ChargePointDisconnected(chargePoint ocpp16.ChargePointConnection) cs.log.DEBUG.Printf("charge point disconnected: %s", chargePoint.ID()) - if cp, err := cs.chargepointByID(chargePoint.ID()); err != nil { + if cp, err := cs.ChargepointByID(chargePoint.ID()); err != nil { cp.connect(false) } } diff --git a/charger/ocpp/cs_core.go b/charger/ocpp/cs_core.go index 1b70dca6f..ac833ff40 100644 --- a/charger/ocpp/cs_core.go +++ b/charger/ocpp/cs_core.go @@ -44,19 +44,13 @@ func (cs *CS) TriggerMessageRequest(id string, requestedMessage remotetrigger.Me } } -func (cs *CS) TriggerMeterValuesRequest(id string, connector int) { - cs.TriggerMessageRequest(id, core.MeterValuesFeatureName, func(request *remotetrigger.TriggerMessageRequest) { - request.ConnectorId = &connector - }) -} - // cp actions func (cs *CS) OnAuthorize(id string, request *core.AuthorizeRequest) (*core.AuthorizeConfirmation, error) { cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -68,7 +62,7 @@ func (cs *CS) OnBootNotification(id string, request *core.BootNotificationReques cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -80,7 +74,7 @@ func (cs *CS) OnDataTransfer(id string, request *core.DataTransferRequest) (*cor cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -92,7 +86,7 @@ func (cs *CS) OnHeartbeat(id string, request *core.HeartbeatRequest) (*core.Hear cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -104,7 +98,7 @@ func (cs *CS) OnMeterValues(id string, request *core.MeterValuesRequest) (*core. cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -116,7 +110,7 @@ func (cs *CS) OnStatusNotification(id string, request *core.StatusNotificationRe cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -128,7 +122,7 @@ func (cs *CS) OnStartTransaction(id string, request *core.StartTransactionReques cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -140,7 +134,7 @@ func (cs *CS) OnStopTransaction(id string, request *core.StopTransactionRequest) cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -152,7 +146,7 @@ func (cs *CS) OnDiagnosticsStatusNotification(id string, request *firmware.Diagn cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } @@ -164,7 +158,7 @@ func (cs *CS) OnFirmwareStatusNotification(id string, request *firmware.Firmware cs.mu.Lock() defer cs.mu.Unlock() - cp, err := cs.chargepointByID(id) + cp, err := cs.ChargepointByID(id) if err != nil { return nil, err } diff --git a/charger/ocpp_test.go b/charger/ocpp_test.go index f4ea887a2..ae2a47261 100644 --- a/charger/ocpp_test.go +++ b/charger/ocpp_test.go @@ -17,7 +17,6 @@ const ( ocppTestUrl = "ws://localhost:8887" ocppTestConnectTimeout = 10 * time.Second ocppTestTimeout = 3 * time.Second - ocppTestConnector = 1 ) func TestOcpp(t *testing.T) { @@ -34,7 +33,7 @@ func (suite *ocppTestSuite) SetupSuite() { suite.NotNil(ocpp.Instance()) } -func (suite *ocppTestSuite) startChargePoint(id string) ocpp16.ChargePoint { +func (suite *ocppTestSuite) startChargePoint(id string, connectorId int) ocpp16.ChargePoint { // set a handler for all callback functions handler := &ChargePointHandler{ triggerC: make(chan remotetrigger.MessageTrigger, 1), @@ -48,14 +47,14 @@ func (suite *ocppTestSuite) startChargePoint(id string) ocpp16.ChargePoint { // let cs handle the trigger messages go func() { for msg := range handler.triggerC { - suite.handleTrigger(cp, msg) + suite.handleTrigger(cp, connectorId, msg) } }() return cp } -func (suite *ocppTestSuite) handleTrigger(cp ocpp16.ChargePoint, msg remotetrigger.MessageTrigger) { +func (suite *ocppTestSuite) handleTrigger(cp ocpp16.ChargePoint, connectorId int, msg remotetrigger.MessageTrigger) { switch msg { case core.BootNotificationFeatureName: if res, err := cp.BootNotification("demo", "evcc"); err != nil { @@ -65,14 +64,14 @@ func (suite *ocppTestSuite) handleTrigger(cp ocpp16.ChargePoint, msg remotetrigg } case core.StatusNotificationFeatureName: - if res, err := cp.StatusNotification(ocppTestConnector, core.NoError, core.ChargePointStatusAvailable); err != nil { + if res, err := cp.StatusNotification(connectorId, core.NoError, core.ChargePointStatusAvailable); err != nil { suite.T().Log("StatusNotification:", err) } else { suite.T().Log("StatusNotification:", res) } case core.MeterValuesFeatureName: - if res, err := cp.MeterValues(1, []types.MeterValue{ + if res, err := cp.MeterValues(connectorId, []types.MeterValue{ { Timestamp: types.NewDateTime(suite.clock.Now()), SampledValue: []types.SampledValue{ @@ -92,41 +91,54 @@ func (suite *ocppTestSuite) handleTrigger(cp ocpp16.ChargePoint, msg remotetrigg } func (suite *ocppTestSuite) TestConnect() { - // start cp client - cp := suite.startChargePoint("test") - suite.NoError(cp.Start(ocppTestUrl)) - suite.True(cp.IsConnected()) + // 1st charge point- remote + cp1 := suite.startChargePoint("test-1", 1) + suite.NoError(cp1.Start(ocppTestUrl)) + suite.True(cp1.IsConnected()) - // start cp server - c, err := NewOCPP("test", ocppTestConnector, "", "", 0, false, false, ocppTestConnectTimeout, ocppTestTimeout, "A") - if err != nil { + // 1st charge point- local + c1, err := NewOCPP("test-1", 1, "", "", 0, false, false, ocppTestConnectTimeout, ocppTestTimeout, "A") + suite.Require().NoError(err) + + { + suite.clock.Add(ocppTestTimeout) + c1.conn.TestClock(suite.clock) + + // status + _, err = c1.Status() suite.NoError(err) - return + + // power + f, err := c1.currentPower() + suite.NoError(err) + suite.Equal(1e3, f) + + // energy + f, err = c1.totalEnergy() + suite.NoError(err) + suite.Equal(1.2, f) } - suite.clock.Add(ocppTestTimeout) - c.cp.TestClock(suite.clock) - - // status - _, err = c.Status() - suite.NoError(err) - - // power - f, err := c.currentPower() - suite.NoError(err) - suite.Equal(1e3, f) - - // energy - f, err = c.totalEnergy() - suite.NoError(err) - suite.Equal(1.2, f) - - // 2nd charge point - cp2 := suite.startChargePoint("test2") - suite.NoError(cp2.Start(ocppTestUrl)) + // 2nd charge point - remote + cp2 := suite.startChargePoint("test-2", 1) + suite.Require().NoError(cp2.Start(ocppTestUrl)) suite.True(cp2.IsConnected()) + // 2nd charge point - local + c2, err := NewOCPP("test-2", 1, "", "", 0, false, false, ocppTestConnectTimeout, ocppTestTimeout, "A") + suite.Require().NoError(err) + + { + suite.clock.Add(ocppTestTimeout) + c2.conn.TestClock(suite.clock) + + // status + _, err = c2.Status() + suite.NoError(err) + } + // error on unconfigured 2nd charge point - _, err = cp2.BootNotification("demo", "evcc") - suite.Error(err) + cp3 := suite.startChargePoint("unconfigured", 1) + _, err = cp3.BootNotification("model", "vendor") + suite.Require().Error(err) }