diff --git a/.changeset/ocpp-last-power.md b/.changeset/ocpp-last-power.md new file mode 100644 index 00000000..3693c614 --- /dev/null +++ b/.changeset/ocpp-last-power.md @@ -0,0 +1,7 @@ +--- +"ftw": patch +--- + +OCPP charger power uses one accept rule for dispatch and forecast, so a +stale, per-phase, negative, or energy-only sample cannot publish a phantom EV load. +1.6 Available/unplug now zeros last power like 2.0.1. diff --git a/go/internal/ocpp/forecast_power_test.go b/go/internal/ocpp/forecast_power_test.go index 0f881c5d..c0f61cc3 100644 --- a/go/internal/ocpp/forecast_power_test.go +++ b/go/internal/ocpp/forecast_power_test.go @@ -2,6 +2,7 @@ package ocpp import ( "fmt" + "math" "strings" "testing" "time" @@ -16,12 +17,14 @@ import ( ) type forecastOCPPFixture struct { - h *Handler - tel *telemetry.Store - power func(float64, time.Time) - energy func() - status func() - stop func() + h *Handler + tel *telemetry.Store + power func(float64, time.Time) + energy func() + status func() + available func() + stop func() + phases func(l1, l2, l3 float64, at time.Time) } func newForecastOCPPFixture(t *testing.T, version string) forecastOCPPFixture { @@ -41,7 +44,17 @@ func newForecastOCPPFixture(t *testing.T, version string) forecastOCPPFixture { f.status = func() { h.OnStatusNotification("charger", &core.StatusNotificationRequest{ConnectorId: 1, Status: core.ChargePointStatusCharging}) } + f.available = func() { + h.OnStatusNotification("charger", &core.StatusNotificationRequest{ConnectorId: 1, Status: core.ChargePointStatusAvailable}) + } f.stop = func() { h.OnStopTransaction("charger", &core.StopTransactionRequest{}) } + f.phases = func(l1, l2, l3 float64, at time.Time) { + h.OnMeterValues("charger", &core.MeterValuesRequest{ConnectorId: 1, MeterValue: []types16.MeterValue{{Timestamp: types16.NewDateTime(at), SampledValue: []types16.SampledValue{ + {Value: fmt.Sprint(l1), Measurand: types16.MeasurandPowerActiveImport, Phase: types16.PhaseL1}, + {Value: fmt.Sprint(l2), Measurand: types16.MeasurandPowerActiveImport, Phase: types16.PhaseL2}, + {Value: fmt.Sprint(l3), Measurand: types16.MeasurandPowerActiveImport, Phase: types16.PhaseL3}, + }}}}) + } } else { v := &handlerV201{h} f.power = func(w float64, at time.Time) { @@ -53,9 +66,19 @@ func newForecastOCPPFixture(t *testing.T, version string) forecastOCPPFixture { f.status = func() { v.OnStatusNotification("charger", &availability.StatusNotificationRequest{EvseID: 1, ConnectorID: 1, ConnectorStatus: availability.ConnectorStatusOccupied}) } + f.available = func() { + v.OnStatusNotification("charger", &availability.StatusNotificationRequest{EvseID: 1, ConnectorID: 1, ConnectorStatus: availability.ConnectorStatusAvailable}) + } f.stop = func() { v.OnTransactionEvent("charger", &transactions.TransactionEventRequest{EventType: transactions.TransactionEventEnded}) } + f.phases = func(l1, l2, l3 float64, at time.Time) { + v.OnMeterValues("charger", &meter.MeterValuesRequest{EvseID: 1, MeterValue: []types201.MeterValue{{Timestamp: *types201.NewDateTime(at), SampledValue: []types201.SampledValue{ + {Value: l1, Measurand: types201.MeasurandPowerActiveImport, Phase: types201.PhaseL1}, + {Value: l2, Measurand: types201.MeasurandPowerActiveImport, Phase: types201.PhaseL2}, + {Value: l3, Measurand: types201.MeasurandPowerActiveImport, Phase: types201.PhaseL3}, + }}}}) + } } return f } @@ -143,8 +166,8 @@ func TestForecastOCPPRejectsFutureAndOutOfOrderSamples(t *testing.T) { if !r.Valid || r.EVW != 700 { t.Fatalf("out-of-order power replaced accepted measurement: %+v", r) } - if raw := f.tel.Get("charger", telemetry.DerEV); raw.RawW != 9000 { - t.Fatal("forecast qualification changed existing dispatch publication") + if raw := f.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 700 { + t.Fatalf("out-of-order power replaced dispatch lastPowerW: %+v", raw) } }) } @@ -157,8 +180,74 @@ func TestForecastOCPP201TransactionPowerAndMissingPower(t *testing.T) { if f.reading().Valid { t.Fatal("power-free transaction event became zero measurement") } + at := time.Now() + f.power(700, at) + v.OnTransactionEvent("charger", &transactions.TransactionEventRequest{EventType: transactions.TransactionEventUpdated, MeterValue: []types201.MeterValue{{Timestamp: *types201.NewDateTime(at), SampledValue: []types201.SampledValue{ + {Value: 16, Measurand: types201.MeasurandCurrentImport}, + {Value: 1234, Measurand: types201.MeasurandEnergyActiveImportRegister}, + }}}}) + if raw := f.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 700 { + t.Fatalf("energy-only transaction event wrote lastPowerW=%v", raw) + } + time.Sleep(2 * time.Millisecond) v.OnTransactionEvent("charger", &transactions.TransactionEventRequest{EventType: transactions.TransactionEventUpdated, MeterValue: []types201.MeterValue{{Timestamp: *types201.NewDateTime(time.Now()), SampledValue: []types201.SampledValue{{Value: 0, Measurand: types201.MeasurandPowerActiveImport}}}}}) if r := f.reading(); !r.Valid || r.EVW != 0 { t.Fatalf("real transaction zero rejected: %+v", r) } } + +func TestOCPPLastPowerSharedAcceptRule(t *testing.T) { + for _, version := range []string{"1.6", "2.0.1"} { + t.Run(version, func(t *testing.T) { + f := newForecastOCPPFixture(t, version) + at := time.Now() + f.power(700, at) + f.energy() + if raw := f.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 700 { + t.Fatalf("energy-only sample overwrote last power: %+v", raw) + } + f.power(-50, at.Add(time.Millisecond)) + if raw := f.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 700 { + t.Fatalf("negative import published: %+v", raw) + } + f.phases(300, 250, 150, at) + if raw := f.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 700 { + t.Fatalf("phased samples overwrote accepted total: %+v", raw) + } + f.available() + if raw := f.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 0 { + t.Fatalf("Available left lastPowerW set: %+v", raw) + } + if r := f.reading(); r.Valid { + t.Fatalf("Available published a measured zero: %+v", r) + } + + f2 := newForecastOCPPFixture(t, version) + f2.phases(300, 250, 150, time.Now()) + if raw := f2.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 700 { + t.Fatalf("phase-only samples not summed: %+v", raw) + } + + f3 := newForecastOCPPFixture(t, version) + f3.power(700, time.Now()) + f3.h.OnDisconnect("charger") + if raw := f3.tel.Get("charger", telemetry.DerEV); raw == nil || raw.RawW != 0 { + t.Fatalf("disconnect left lastPowerW set: %+v", raw) + } + }) + } +} + +func TestMeterPowerWPrefersUnphasedThenSum(t *testing.T) { + w, ok := meterPowerW([]powerSample{{w: 700}, {w: 300, phase: "L1"}}) + if !ok || w != 700 { + t.Fatalf("unphased+phase: got %v ok=%v", w, ok) + } + w, ok = meterPowerW([]powerSample{{w: 300, phase: "L1"}, {w: 250, phase: "L2"}, {w: 150, phase: "L3"}}) + if !ok || w != 700 { + t.Fatalf("phase sum: got %v ok=%v", w, ok) + } + if _, ok := meterPowerW([]powerSample{{w: -1}, {w: math.NaN(), phase: "L1"}}); ok { + t.Fatal("negative/NaN samples accepted") + } +} diff --git a/go/internal/ocpp/handlers.go b/go/internal/ocpp/handlers.go index 09f9b294..9a08e4dc 100644 --- a/go/internal/ocpp/handlers.go +++ b/go/internal/ocpp/handlers.go @@ -389,8 +389,7 @@ func (h *Handler) OnConnect(id string) { s.connectionGeneration++ s.connectedKnown = false s.charging = false - s.lastPowerW = 0 - s.forecastPower.Known = false + s.clearMeasuredPower() s.identityCurrent = false s.featureProfiles = "" s.steerable = nil @@ -416,8 +415,7 @@ func (h *Handler) OnDisconnect(id string) { h.cancelIdentityProbeLocked(s) s.connectedKnown = false s.charging = false - s.lastPowerW = 0 - s.forecastPower.Known = false + s.clearMeasuredPower() h.mu.Unlock() // Push a zero so the dispatch clamp releases — otherwise the last known // non-zero w would survive until staleness kicks in. @@ -488,6 +486,7 @@ func (h *Handler) OnStatusNotification(id string, req *core.StatusNotificationRe case core.ChargePointStatusAvailable, core.ChargePointStatusUnavailable: s.connected = false s.charging = false + s.clearMeasuredPower() case core.ChargePointStatusPreparing, core.ChargePointStatusFinishing, core.ChargePointStatusSuspendedEV, @@ -501,6 +500,7 @@ func (h *Handler) OnStatusNotification(id string, req *core.StatusNotificationRe case core.ChargePointStatusFaulted: s.connected = true s.charging = false + s.clearMeasuredPower() } h.mu.Unlock() @@ -521,6 +521,11 @@ func (h *Handler) OnMeterValues(id string, req *core.MeterValuesRequest) (*core. h.mu.Lock() received := time.Now() for _, mv := range req.MeterValue { + measured := received + if mv.Timestamp != nil { + measured = mv.Timestamp.Time + } + var samples []powerSample for _, sv := range mv.SampledValue { measurand := sv.Measurand // OCPP 1.6 default measurand if unspecified. @@ -533,17 +538,13 @@ func (h *Handler) OnMeterValues(id string, req *core.MeterValuesRequest) (*core. } switch measurand { case types.MeasurandPowerActiveImport: + if sv.Unit != "" && sv.Unit != types.UnitOfMeasureW && sv.Unit != types.UnitOfMeasureKW { + continue + } if sv.Unit == types.UnitOfMeasureKW { val *= 1000 } - s.lastPowerW = val - if sv.Phase == "" && (sv.Unit == "" || sv.Unit == types.UnitOfMeasureW || sv.Unit == types.UnitOfMeasureKW) { - measured := received - if mv.Timestamp != nil { - measured = mv.Timestamp.Time - } - s.recordForecastPower(val, measured, received) - } + samples = append(samples, powerSample{w: val, phase: string(sv.Phase)}) case types.MeasurandEnergyActiveImportRegister: if sv.Unit == types.UnitOfMeasureKWh { val *= 1000 @@ -553,6 +554,9 @@ func (h *Handler) OnMeterValues(id string, req *core.MeterValuesRequest) (*core. } } } + if w, ok := meterPowerW(samples); ok { + s.recordPower(w, measured, received) + } } h.mu.Unlock() @@ -593,8 +597,7 @@ func (h *Handler) OnStopTransaction(id string, req *core.StopTransactionRequest) sessionWh := float64(req.MeterStop) - s.sessionStartMeterWh s.transactionID = -1 s.charging = false - s.lastPowerW = 0 - s.forecastPower.Known = false + s.clearMeasuredPower() s.sessionMeterWh = sessionWh h.mu.Unlock() @@ -652,13 +655,54 @@ func (h *Handler) pushReading(id string, s *chargerState) { h.tel.Update(id, telemetry.DerEV, w, nil, blob) } -// recordForecastPower is separate from the existing status/dispatch power. -// Only a real aggregate power measurand may refresh it. Samples older than a -// connection or the accepted sample, and future/nonfinite values, cannot revive -// a stale or synthesized reading. The caller holds h.mu. -func (s *chargerState) recordForecastPower(w float64, measured, received time.Time) { +// powerSample is one Power.Active.Import after unit conversion. An empty phase +// is the charger total; anything else is one phase of that total. +type powerSample struct { + w float64 + phase string +} + +// meterPowerW picks one watts value for a MeterValue: the unphased total when +// present, otherwise the sum of finite, non-negative phase samples. +func meterPowerW(samples []powerSample) (float64, bool) { + var total, sum float64 + hasTotal, hasPhase := false, false + for _, s := range samples { + if math.IsNaN(s.w) || math.IsInf(s.w, 0) || s.w < 0 { + continue + } + if s.phase == "" { + total = s.w + hasTotal = true + } else { + sum += s.w + hasPhase = true + } + } + if hasTotal { + return total, true + } + if hasPhase { + return sum, true + } + return 0, false +} + +// recordPower is the shared accept rule for forecast and dispatch lastPowerW. +// Present, unphased-or-summed, finite, >= 0, and not older than the last +// accepted timestamp. Samples from before this socket or in the future cannot +// revive a stale or synthesized reading. The caller holds h.mu. +func (s *chargerState) recordPower(w float64, measured, received time.Time) bool { if math.IsNaN(w) || math.IsInf(w, 0) || w < 0 || measured.IsZero() || measured.After(received) || measured.Before(s.powerConnectedAt) || measured.UnixMilli() <= s.forecastPower.MeasuredAtMS { - return + return false } + s.lastPowerW = w s.forecastPower = telemetry.ForecastPowerSample{Version: 1, Known: true, Watts: w, MeasuredAtMS: measured.UnixMilli(), ReceivedAtMS: received.UnixMilli()} + return true +} + +// clearMeasuredPower zeros dispatch watts without recording a measured sample. +func (s *chargerState) clearMeasuredPower() { + s.lastPowerW = 0 + s.forecastPower.Known = false } diff --git a/go/internal/ocpp/handlers_v201.go b/go/internal/ocpp/handlers_v201.go index 41f46f24..9020eaaa 100644 --- a/go/internal/ocpp/handlers_v201.go +++ b/go/internal/ocpp/handlers_v201.go @@ -102,8 +102,7 @@ func (h *handlerV201) OnStatusNotification(id string, req *availability.StatusNo case availability.ConnectorStatusAvailable, availability.ConnectorStatusUnavailable: s.connected = false s.charging = false - s.lastPowerW = 0 - s.forecastPower.Known = false + s.clearMeasuredPower() case availability.ConnectorStatusOccupied, availability.ConnectorStatusReserved: s.connected = true s.connectedKnown = true @@ -113,8 +112,7 @@ func (h *handlerV201) OnStatusNotification(id string, req *availability.StatusNo s.connected = true s.connectedKnown = true s.charging = false - s.lastPowerW = 0 - s.forecastPower.Known = false + s.clearMeasuredPower() } faulted := req.ConnectorStatus == availability.ConnectorStatusFaulted h.mu.Unlock() @@ -140,9 +138,13 @@ func (h *handlerV201) OnTransactionEvent(id string, req *transactions.Transactio s := h.state(id) // Meter samples ride along with every event type. - powerW, energyWh, hasEnergy := sampledValuesV201(req.MeterValue) + energyWh, hasEnergy := sampledEnergyV201(req.MeterValue) h.mu.Lock() + acceptedPower := false + if req.EventType != transactions.TransactionEventEnded { + acceptedPower = s.recordPowerV201(req.MeterValue, time.Now()) + } switch req.EventType { case transactions.TransactionEventStarted: // 2.0.1 transaction ids are strings; the shared state keeps an int for @@ -164,8 +166,10 @@ func (h *handlerV201) OnTransactionEvent(id string, req *transactions.Transactio s.sessionMeterWh = energyWh - s.sessionStartMeterWh } // A zero power sample during a live transaction is a genuine pause, - // not a missing reading, so it is taken at face value. - s.charging = powerW > 0 + // not a missing reading. Energy/current without power keeps last watts. + if acceptedPower { + s.charging = s.lastPowerW > 0 + } case transactions.TransactionEventEnded: if hasEnergy { @@ -174,15 +178,10 @@ func (h *handlerV201) OnTransactionEvent(id string, req *transactions.Transactio s.transactionID = -1 s.transactionRef = "" s.charging = false - s.lastPowerW = 0 - s.forecastPower.Known = false - powerW = 0 + s.clearMeasuredPower() } - if req.EventType != transactions.TransactionEventEnded { - s.lastPowerW = powerW - } - s.recordForecastPowerV201(req.MeterValue, time.Now()) + powerW := s.lastPowerW sessionWh := s.sessionMeterWh ended := req.EventType == transactions.TransactionEventEnded h.mu.Unlock() @@ -210,11 +209,10 @@ func (h *handlerV201) OnTransactionEvent(id string, req *transactions.Transactio func (h *handlerV201) OnMeterValues(id string, req *meter.MeterValuesRequest) (*meter.MeterValuesResponse, error) { s := h.state(id) - powerW, energyWh, hasEnergy := sampledValuesV201(req.MeterValue) + energyWh, hasEnergy := sampledEnergyV201(req.MeterValue) h.mu.Lock() - s.lastPowerW = powerW - s.recordForecastPowerV201(req.MeterValue, time.Now()) + s.recordPowerV201(req.MeterValue, time.Now()) if hasEnergy && s.transactionID >= 0 { s.sessionMeterWh = energyWh - s.sessionStartMeterWh } @@ -237,33 +235,29 @@ func (h *handlerV201) OnAuthorize(id string, _ *authorization.AuthorizeRequest) }), nil } -// sampledValuesV201 pulls active-import power and energy out of a 2.0.1 meter -// value set, normalising kW/kWh to W/Wh. +// sampledEnergyV201 pulls active-import energy out of a 2.0.1 meter value set, +// normalising kWh to Wh. Power goes through recordPowerV201 so a missing +// measurand cannot write 0 W. // // 2.0.1 always states the measurand, so unlike 1.6 there is no default to // assume. hasEnergy distinguishes "no energy sample in this batch" from a // genuine zero reading, which matters because session energy is a difference // against the transaction's starting register. -func sampledValuesV201(values []types201.MeterValue) (powerW, energyWh float64, hasEnergy bool) { +func sampledEnergyV201(values []types201.MeterValue) (energyWh float64, hasEnergy bool) { for _, mv := range values { for _, sv := range mv.SampledValue { + if sv.Measurand != types201.MeasurandEnergyActiveImportRegister { + continue + } val := sv.Value - switch sv.Measurand { - case types201.MeasurandPowerActiveImport: - if unitIsKilo(sv.UnitOfMeasure) { - val *= 1000 - } - powerW = val - case types201.MeasurandEnergyActiveImportRegister: - if unitIsKilo(sv.UnitOfMeasure) { - val *= 1000 - } - energyWh = val - hasEnergy = true + if unitIsKilo(sv.UnitOfMeasure) { + val *= 1000 } + energyWh = val + hasEnergy = true } } - return powerW, energyWh, hasEnergy + return energyWh, hasEnergy } // unitIsKilo reports whether a sample is expressed in kW or kWh. An absent unit @@ -280,10 +274,12 @@ func unitIsKilo(u *types201.UnitOfMeasure) bool { } } -func (s *chargerState) recordForecastPowerV201(values []types201.MeterValue, received time.Time) { +func (s *chargerState) recordPowerV201(values []types201.MeterValue, received time.Time) bool { + accepted := false for _, mv := range values { + var samples []powerSample for _, sv := range mv.SampledValue { - if sv.Measurand != types201.MeasurandPowerActiveImport || sv.Phase != "" { + if sv.Measurand != types201.MeasurandPowerActiveImport { continue } w := sv.Value @@ -298,11 +294,19 @@ func (s *chargerState) recordForecastPowerV201(values []types201.MeterValue, rec w *= math.Pow10(*unit.Multiplier) } } - measured := mv.Timestamp.Time - if measured.IsZero() { - measured = received - } - s.recordForecastPower(w, measured, received) + samples = append(samples, powerSample{w: w, phase: string(sv.Phase)}) + } + w, ok := meterPowerW(samples) + if !ok { + continue + } + measured := mv.Timestamp.Time + if measured.IsZero() { + measured = received + } + if s.recordPower(w, measured, received) { + accepted = true } } + return accepted }