diff --git a/core/loadpoint.go b/core/loadpoint.go index af7f5447a..75322276a 100644 --- a/core/loadpoint.go +++ b/core/loadpoint.go @@ -23,6 +23,8 @@ import ( "github.com/avast/retry-go/v3" "github.com/benbjohnson/clock" "github.com/cjrd/allocate" + "github.com/emirpasic/gods/queues" + aq "github.com/emirpasic/gods/queues/arrayqueue" ) const ( @@ -150,6 +152,8 @@ type LoadPoint struct { chargeRemainingDuration time.Duration // Remaining charge duration chargeRemainingEnergy float64 // Remaining charge energy in Wh progress *Progress // Step-wise progress indicator + + tasks queues.Queue // tasks to be executed } // NewLoadPointFromConfig creates a new loadpoint @@ -270,6 +274,7 @@ func NewLoadPoint(log *util.Logger) *LoadPoint { Disable: ThresholdConfig{Delay: 3 * time.Minute, Threshold: 0}, // t, W GuardDuration: 5 * time.Minute, progress: NewProgress(0, 10), // soc progress indicator + tasks: aq.New(), } return lp @@ -817,16 +822,7 @@ func (lp *LoadPoint) setActiveVehicle(vehicle api.Vehicle) { // release lock to unblock api lp.Unlock() - // publish odometer once - if vs, ok := lp.vehicle.(api.VehicleOdometer); ok { - if odo, err := vs.Odometer(); err == nil { - lp.log.DEBUG.Printf("vehicle odometer: %.0fkm", odo) - lp.publish("vehicleOdometer", odo) - } else { - lp.log.ERROR.Printf("vehicle odometer: %v", err) - } - } - + lp.addTask(lp.vehicleOdometer) lp.applyAction(vehicle.OnIdentified()) // re-apply lock to match defer above @@ -923,6 +919,18 @@ func (lp *LoadPoint) identifyVehicleByStatus() { } } +// vehicleOdometer updates odometer +func (lp *LoadPoint) vehicleOdometer() { + if vs, ok := lp.vehicle.(api.VehicleOdometer); ok { + if odo, err := vs.Odometer(); err == nil { + lp.log.DEBUG.Printf("vehicle odometer: %.0fkm", odo) + lp.publish("vehicleOdometer", odo) + } else { + lp.log.ERROR.Printf("vehicle odometer: %v", err) + } + } +} + // updateChargerStatus updates charger status and detects car connected/disconnected events func (lp *LoadPoint) updateChargerStatus() error { status, err := lp.charger.Status() @@ -1429,8 +1437,24 @@ func (lp *LoadPoint) publishSoCAndRange() { } } +// addTask adds a single task to the queue +func (lp *LoadPoint) addTask(task func()) { + lp.tasks.Enqueue(task) +} + +// processTasks executes a single task from the queue +func (lp *LoadPoint) processTasks() { + if lp.tasks != nil { + if task, ok := lp.tasks.Dequeue(); ok { + task.(func())() + } + } +} + // Update is the main control function. It reevaluates meters and charger state func (lp *LoadPoint) Update(sitePower float64, cheap, batteryBuffered bool) { + lp.processTasks() + mode := lp.GetMode() lp.publish("mode", mode) diff --git a/go.mod b/go.mod index 8dd527690..71a2dfbe9 100644 --- a/go.mod +++ b/go.mod @@ -23,6 +23,7 @@ require ( github.com/dustin/go-humanize v1.0.0 github.com/dylanmei/iso8601 v0.1.0 github.com/eclipse/paho.mqtt.golang v1.3.5 + github.com/emirpasic/gods v1.18.1 github.com/evcc-io/eebus v0.0.0-20220628095038-be707d322cfc github.com/fatih/structs v1.1.0 github.com/foogod/go-powerwall v0.2.0 diff --git a/go.sum b/go.sum index 9eda5c6e4..d299db26b 100644 --- a/go.sum +++ b/go.sum @@ -216,6 +216,8 @@ github.com/eapache/queue v1.1.0/go.mod h1:6eCeP0CKFpHLu8blIFXhExK/dRa7WDZfr6jVFP github.com/eclipse/paho.mqtt.golang v1.3.5 h1:sWtmgNxYM9P2sP+xEItMozsR3w0cqZFlqnNN1bdl41Y= github.com/eclipse/paho.mqtt.golang v1.3.5/go.mod h1:eTzb4gxwwyWpqBUHGQZ4ABAV7+Jgm1PklsYT/eo8Hcc= github.com/edsrzf/mmap-go v1.0.0/go.mod h1:YO35OhQPt3KJa3ryjFM5Bs14WD66h8eGKpfaBNrHW5M= +github.com/emirpasic/gods v1.18.1 h1:FXtiHYKDGKCW2KzwZKx0iC0PQmdlorYgdFG9jPXJ1Bc= +github.com/emirpasic/gods v1.18.1/go.mod h1:8tpGGwCnJ5H4r6BWwaV6OrWmMoPhUl5jm/FMNAnJvWQ= github.com/envoyproxy/go-control-plane v0.6.9/go.mod h1:SBwIajubJHhxtWwsL9s8ss4safvEdbitLhGGK48rN6g= github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= github.com/envoyproxy/go-control-plane v0.9.1-0.20191026205805-5f8ba28d4473/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4=