From 31f2df25dce5223f7c5a6880dd44669701bd1a4d Mon Sep 17 00:00:00 2001 From: andig Date: Thu, 3 Oct 2024 13:02:23 +0200 Subject: [PATCH] Site: parallelise meter reading (#15372) --- core/loadpoint.go | 40 ++++++------ core/prioritizer/prioritizer.go | 4 ++ core/site.go | 107 ++++++++++++++++++++------------ 3 files changed, 93 insertions(+), 58 deletions(-) diff --git a/core/loadpoint.go b/core/loadpoint.go index 9dd38b802..87b9e3382 100644 --- a/core/loadpoint.go +++ b/core/loadpoint.go @@ -1374,8 +1374,9 @@ func (lp *Loadpoint) pvMaxCurrent(mode api.ChargeMode, sitePower float64, batter } // UpdateChargePowerAndCurrents updates charge meter power and currents for load management -func (lp *Loadpoint) UpdateChargePowerAndCurrents() { - if power, err := backoff.RetryWithData(lp.chargeMeter.CurrentPower, bo()); err == nil { +func (lp *Loadpoint) UpdateChargePowerAndCurrents() float64 { + power, err := backoff.RetryWithData(lp.chargeMeter.CurrentPower, bo()) + if err == nil { lp.Lock() lp.chargePower = power // update value if no error lp.Unlock() @@ -1390,31 +1391,34 @@ func (lp *Loadpoint) UpdateChargePowerAndCurrents() { lp.log.WARN.Printf("charge power must not be negative: %.0f", power) } } else { + power = 0 lp.log.ERROR.Printf("charge power: %v", err) } // update charge currents lp.chargeCurrents = nil - phaseMeter, ok := lp.chargeMeter.(api.PhaseCurrents) - if !ok { - return // don't guess - } + if phaseMeter, ok := lp.chargeMeter.(api.PhaseCurrents); ok { + if err := backoff.Retry(func() error { + i1, i2, i3, err := phaseMeter.Currents() + if err != nil { + return err + } - if err := backoff.Retry(func() error { - i1, i2, i3, err := phaseMeter.Currents() - if err != nil { - return err + lp.Lock() + lp.chargeCurrents = []float64{i1, i2, i3} + lp.Unlock() + + lp.log.DEBUG.Printf("charge currents: %.3gA", lp.chargeCurrents) + lp.publish(keys.ChargeCurrents, lp.chargeCurrents) + + return nil + }, bo()); err != nil { + lp.log.ERROR.Printf("charge currents: %v", err) } - - lp.chargeCurrents = []float64{i1, i2, i3} - lp.log.DEBUG.Printf("charge currents: %.3gA", lp.chargeCurrents) - lp.publish(keys.ChargeCurrents, lp.chargeCurrents) - - return nil - }, bo()); err != nil { - lp.log.ERROR.Printf("charge currents: %v", err) } + + return power } // phasesFromChargeCurrents uses PhaseCurrents interface to count phases with current >=1A diff --git a/core/prioritizer/prioritizer.go b/core/prioritizer/prioritizer.go index fe4480751..39c90b8bf 100644 --- a/core/prioritizer/prioritizer.go +++ b/core/prioritizer/prioritizer.go @@ -2,12 +2,14 @@ package prioritizer import ( "fmt" + "sync" "github.com/evcc-io/evcc/core/loadpoint" "github.com/evcc-io/evcc/util" ) type Prioritizer struct { + mu sync.Mutex log *util.Logger demand map[loadpoint.API]float64 } @@ -21,7 +23,9 @@ func New(log *util.Logger) *Prioritizer { func (p *Prioritizer) UpdateChargePowerFlexibility(lp loadpoint.API) { if power := lp.GetChargePowerFlexibility(); power >= 0 { + p.mu.Lock() p.demand[lp] = power + p.mu.Unlock() } } diff --git a/core/site.go b/core/site.go index 961a58144..985c529a6 100644 --- a/core/site.go +++ b/core/site.go @@ -30,6 +30,7 @@ import ( "github.com/evcc-io/evcc/util/config" "github.com/evcc-io/evcc/util/telemetry" "github.com/smallnest/chanx" + "golang.org/x/sync/errgroup" ) const standbyPower = 10 // consider less than 10W as charger in standby @@ -101,6 +102,7 @@ type Site struct { // cached state gridPower float64 // Grid power pvPower float64 // PV power + auxPower float64 // Aux power batteryPower float64 // Battery charge power batterySoc float64 // Battery soc batteryMode api.BatteryMode // Battery mode (runtime only, not persisted) @@ -483,7 +485,30 @@ func (site *Site) updatePvMeters() { site.publish(keys.Pv, mm) } -// updateExtMeters updates ext meters. All measurements are optional. +// updateAuxMeters updates aux meters +func (site *Site) updateAuxMeters() { + if len(site.auxMeters) == 0 { + return + } + + mm := make([]meterMeasurement, len(site.auxMeters)) + + for i, meter := range site.auxMeters { + if power, err := meter.CurrentPower(); err == nil { + site.auxPower += power + mm[i].Power = power + site.log.DEBUG.Printf("aux power %d: %.0fW", i+1, power) + } else { + site.log.ERROR.Printf("aux meter %d: %v", i+1, err) + } + } + + site.log.DEBUG.Printf("aux power: %.0fW", site.auxPower) + site.publish(keys.AuxPower, site.auxPower) + site.publish(keys.Aux, mm) +} + +// updateExtMeters updates ext meters func (site *Site) updateExtMeters() { if len(site.extMeters) == 0 { return @@ -516,7 +541,7 @@ func (site *Site) updateExtMeters() { // Publishing will be done in separate PR } -// updateBatteryMeters updates battery meters. Power is retried, other measurements are optional. +// updateBatteryMeters updates battery meters func (site *Site) updateBatteryMeters() error { if len(site.batteryMeters) == 0 { return nil @@ -605,7 +630,7 @@ func (site *Site) updateBatteryMeters() error { return nil } -// updateGridMeter updates grid meter. Power is retried, other measurements are optional. +// updateGridMeter updates grid meter func (site *Site) updateGridMeter() error { if site.gridMeter == nil { return nil @@ -655,15 +680,17 @@ func (site *Site) updateGridMeter() error { return nil } -// updateMeter updates and publishes single meter func (site *Site) updateMeters() error { - // TODO parallelize once modbus supports that - site.updatePvMeters() - if err := site.updateBatteryMeters(); err != nil { - return err - } - site.updateExtMeters() - return site.updateGridMeter() + g, _ := errgroup.WithContext(context.Background()) + + g.Go(func() error { site.updatePvMeters(); return nil }) + g.Go(func() error { site.updateAuxMeters(); return nil }) + g.Go(func() error { site.updateExtMeters(); return nil }) + + g.Go(site.updateBatteryMeters) + g.Go(site.updateGridMeter) + + return g.Wait() } // sitePower returns @@ -715,27 +742,7 @@ func (site *Site) sitePower(totalChargePower, flexiblePower float64) (float64, b sitePower := sitePower(site.log, site.GetMaxGridSupplyWhileBatteryCharging(), site.gridPower, batteryPower, site.GetResidualPower()) // deduct smart loads - if len(site.auxMeters) > 0 { - var auxPower float64 - mm := make([]meterMeasurement, len(site.auxMeters)) - - for i, meter := range site.auxMeters { - if power, err := meter.CurrentPower(); err == nil { - auxPower += power - mm[i].Power = power - site.log.DEBUG.Printf("aux power %d: %.0fW", i+1, power) - } else { - site.log.ERROR.Printf("aux meter %d: %v", i+1, err) - } - } - - sitePower -= auxPower - - site.log.DEBUG.Printf("aux power: %.0fW", auxPower) - site.publish(keys.AuxPower, auxPower) - - site.publish(keys.Aux, mm) - } + sitePower -= site.auxPower // handle priority if flexiblePower > 0 { @@ -818,17 +825,37 @@ func (site *Site) publishTariffs(greenShareHome float64, greenShareLoadpoints fl } } +// updateLoadpoints updates all loadpoints' charge power +func (site *Site) updateLoadpoints() float64 { + var ( + wg sync.WaitGroup + mu sync.Mutex + sum float64 + ) + + wg.Add(len(site.loadpoints)) + for _, lp := range site.loadpoints { + go func() { + power := lp.UpdateChargePowerAndCurrents() + site.prioritizer.UpdateChargePowerFlexibility(lp) + + mu.Lock() + sum += power + mu.Unlock() + + wg.Done() + }() + } + wg.Wait() + + return sum +} + func (site *Site) update(lp updater) { site.log.DEBUG.Println("----") - // update all loadpoint's charge power - var totalChargePower float64 - for _, lp := range site.loadpoints { - lp.UpdateChargePowerAndCurrents() - totalChargePower += lp.GetChargePower() - - site.prioritizer.UpdateChargePowerFlexibility(lp) - } + // update loadpoints + totalChargePower := site.updateLoadpoints() // update all circuits' power and currents if site.circuit != nil {