chore: simplify handling backoffs
This commit is contained in:
parent
7bca010e8b
commit
ab570ebf55
24 changed files with 40 additions and 59 deletions
|
|
@ -208,8 +208,7 @@ func (c *Easee) chargerSite(charger string) (easee.Site, error) {
|
|||
|
||||
// connect creates an HTTP connection to the signalR hub
|
||||
func (c *Easee) connect(ts oauth2.TokenSource) func() (signalr.Connection, error) {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.MaxInterval = time.Minute
|
||||
bo := backoff.NewExponentialBackOff(backoff.WithMaxInterval(time.Minute))
|
||||
|
||||
return func() (conn signalr.Connection, err error) {
|
||||
defer func() {
|
||||
|
|
|
|||
|
|
@ -132,9 +132,9 @@ func NewSalia(uri string, cache time.Duration) (api.Charger, error) {
|
|||
}
|
||||
|
||||
func (wb *Salia) heartbeat() {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.InitialInterval = 5 * time.Second
|
||||
bo.MaxInterval = time.Minute
|
||||
bo := backoff.NewExponentialBackOff(
|
||||
backoff.WithInitialInterval(5*time.Second),
|
||||
backoff.WithMaxInterval(time.Minute))
|
||||
|
||||
for range time.Tick(30 * time.Second) {
|
||||
if err := backoff.Retry(func() error {
|
||||
|
|
|
|||
|
|
@ -124,10 +124,10 @@ func (c *Pulsatrix) connectWs() error {
|
|||
|
||||
// ReconnectWs reconnects to a pulsatrix SECC websocket
|
||||
func (c *Pulsatrix) reconnectWs() {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.InitialInterval = time.Second
|
||||
bo.MaxInterval = 1 * time.Minute
|
||||
bo.MaxElapsedTime = 0 * time.Second // retry forever; default is 15 min
|
||||
bo := backoff.NewExponentialBackOff(
|
||||
backoff.WithInitialInterval(time.Second),
|
||||
backoff.WithMaxInterval(time.Minute),
|
||||
backoff.WithMaxElapsedTime(0)) // retry forever; default is 15 min
|
||||
if err := backoff.Retry(c.connectWs, bo); err != nil {
|
||||
c.log.ERROR.Println(err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -18,9 +18,7 @@ var (
|
|||
|
||||
// bo returns an exponential backoff for reading meter power quickly
|
||||
func bo() *backoff.ExponentialBackOff {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.MaxElapsedTime = time.Second
|
||||
return bo
|
||||
return backoff.NewExponentialBackOff(backoff.WithMaxElapsedTime(time.Second))
|
||||
}
|
||||
|
||||
// powerToCurrent is a helper function to convert power to per-phase current
|
||||
|
|
|
|||
|
|
@ -1369,10 +1369,7 @@ func (lp *Loadpoint) pvMaxCurrent(mode api.ChargeMode, sitePower float64, batter
|
|||
|
||||
// UpdateChargePowerAndCurrents updates charge meter power and currents for load management
|
||||
func (lp *Loadpoint) UpdateChargePowerAndCurrents() {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.MaxElapsedTime = time.Second
|
||||
|
||||
if power, err := backoff.RetryWithData(lp.chargeMeter.CurrentPower, bo); err == nil {
|
||||
if power, err := backoff.RetryWithData(lp.chargeMeter.CurrentPower, bo()); err == nil {
|
||||
lp.Lock()
|
||||
lp.chargePower = power // update value if no error
|
||||
lp.Unlock()
|
||||
|
|
@ -1409,7 +1406,7 @@ func (lp *Loadpoint) UpdateChargePowerAndCurrents() {
|
|||
lp.publish(keys.ChargeCurrents, lp.chargeCurrents)
|
||||
|
||||
return nil
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
lp.log.ERROR.Printf("charge currents: %v", err)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -139,8 +139,7 @@ func NewDsmr(uri, energy string, timeout time.Duration) (api.Meter, error) {
|
|||
// based on https://github.com/basvdlei/gotsmart/blob/master/gotsmart.go
|
||||
func (m *Dsmr) run(conn net.Conn, done chan struct{}) {
|
||||
log := util.NewLogger("dsmr")
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.MaxInterval = 5 * time.Minute
|
||||
bo := backoff.NewExponentialBackOff(backoff.WithMaxInterval(5 * time.Minute))
|
||||
|
||||
handle := func(op string, err error) {
|
||||
log.ERROR.Printf("%s: %v", op, err)
|
||||
|
|
|
|||
|
|
@ -60,9 +60,9 @@ func (m *Server) GetInverter(ip string) *util.Monitor[Inverter] {
|
|||
}
|
||||
|
||||
func (m *Server) readData() {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.MaxInterval = time.Second
|
||||
bo.MaxElapsedTime = 10 * time.Second
|
||||
bo := backoff.NewExponentialBackOff(
|
||||
backoff.WithMaxInterval(time.Second),
|
||||
backoff.WithMaxElapsedTime(10*time.Second))
|
||||
|
||||
for {
|
||||
mu.RLock()
|
||||
|
|
|
|||
|
|
@ -87,9 +87,9 @@ func NewRCT(uri, usage string, cache time.Duration, capacity func() float64) (ap
|
|||
return nil, err
|
||||
}
|
||||
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.InitialInterval = 10 * time.Millisecond
|
||||
bo.MaxElapsedTime = time.Second
|
||||
bo := backoff.NewExponentialBackOff(
|
||||
backoff.WithInitialInterval(10*time.Millisecond),
|
||||
backoff.WithMaxElapsedTime(time.Second))
|
||||
|
||||
m := &RCT{
|
||||
usage: strings.ToLower(usage),
|
||||
|
|
|
|||
|
|
@ -82,7 +82,6 @@ func NewAmberFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
|
||||
func (t *Amber) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Minute)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -91,7 +90,7 @@ func (t *Amber) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(t.GetJSON(uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -55,7 +55,7 @@ func NewAwattarFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
|
||||
func (t *Awattar) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
client := request.NewHelper(t.log)
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
|
|
@ -64,7 +64,7 @@ func (t *Awattar) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(client.GetJSON(t.uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -102,7 +102,6 @@ func (t *EdfTempo) RefreshToken(_ *oauth2.Token) (*oauth2.Token, error) {
|
|||
|
||||
func (t *EdfTempo) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -125,7 +124,7 @@ func (t *EdfTempo) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(t.GetJSON(uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -79,7 +79,7 @@ func NewElectricityMapsFromConfig(other map[string]interface{}) (api.Tariff, err
|
|||
|
||||
func (t *ElectricityMaps) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
uri := fmt.Sprintf("%s/carbon-intensity/forecast?zone=%s", t.uri, t.zone)
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
|
|
@ -88,7 +88,7 @@ func (t *ElectricityMaps) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(t.GetJSON(uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
if res.Error != "" {
|
||||
err = errors.New(res.Error)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -60,7 +60,6 @@ func NewEleringFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
func (t *Elering) run(done chan error) {
|
||||
var once sync.Once
|
||||
client := request.NewHelper(t.log)
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -73,7 +72,7 @@ func (t *Elering) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(client.GetJSON(uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -59,7 +59,6 @@ func NewEnerginetFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
func (t *Energinet) run(done chan error) {
|
||||
var once sync.Once
|
||||
client := request.NewHelper(t.log)
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -73,7 +72,7 @@ func (t *Energinet) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(client.GetJSON(uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -85,8 +85,6 @@ func NewEntsoeFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
func (t *Entsoe) run(done chan error) {
|
||||
var once sync.Once
|
||||
|
||||
bo := newBackoff()
|
||||
|
||||
// Data updated by ESO every half hour, but we only need data every hour to stay current.
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -127,7 +125,7 @@ func (t *Entsoe) run(done chan error) {
|
|||
default:
|
||||
return backoff.Permanent(errors.New("invalid document name: " + doc.XMLName.Local))
|
||||
}
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -39,7 +39,7 @@ func NewGroupeEFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
|
||||
func (t *GroupeE) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
client := request.NewHelper(t.log)
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
|
|
@ -55,7 +55,7 @@ func (t *GroupeE) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(client.GetJSON(uri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -89,7 +89,7 @@ func NewGrünStromIndexFromConfig(other map[string]interface{}) (api.Tariff, err
|
|||
func (t *GrünStromIndex) run(done chan error) {
|
||||
var once sync.Once
|
||||
client := request.NewHelper(t.log)
|
||||
bo := newBackoff()
|
||||
|
||||
uri := fmt.Sprintf("https://api.corrently.io/v2.0/gsi/prediction?zip=%s", t.zip)
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
|
|
|
|||
|
|
@ -11,11 +11,11 @@ import (
|
|||
"github.com/evcc-io/evcc/util/request"
|
||||
)
|
||||
|
||||
func newBackoff() backoff.BackOff {
|
||||
bo := backoff.NewExponentialBackOff()
|
||||
bo.InitialInterval = time.Second
|
||||
bo.MaxElapsedTime = time.Minute
|
||||
return bo
|
||||
func bo() backoff.BackOff {
|
||||
return backoff.NewExponentialBackOff(
|
||||
backoff.WithInitialInterval(time.Second),
|
||||
backoff.WithMaxElapsedTime(time.Minute),
|
||||
)
|
||||
}
|
||||
|
||||
// backoffPermanentError returns a permanent error in case of HTTP 400
|
||||
|
|
|
|||
|
|
@ -57,7 +57,6 @@ func NewNgesoFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
func (t *Ngeso) run(done chan error) {
|
||||
var once sync.Once
|
||||
client := request.NewHelper(t.log)
|
||||
bo := newBackoff()
|
||||
|
||||
// Use national results by default.
|
||||
var tReq ngeso.CarbonForecastRequest
|
||||
|
|
|
|||
|
|
@ -84,7 +84,6 @@ func NewOctopusFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
func (t *Octopus) run(done chan error) {
|
||||
var once sync.Once
|
||||
client := request.NewHelper(t.log)
|
||||
bo := newBackoff()
|
||||
|
||||
var restQueryUri string
|
||||
|
||||
|
|
@ -115,7 +114,7 @@ func (t *Octopus) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(client.GetJSON(restQueryUri, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -70,7 +70,6 @@ func NewPunFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
|
||||
func (t *Pun) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -81,7 +80,7 @@ func (t *Pun) run(done chan error) {
|
|||
today, err = t.getData(time.Now())
|
||||
|
||||
return err
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -49,7 +49,6 @@ func NewSmartEnergyFromConfig(other map[string]interface{}) (api.Tariff, error)
|
|||
func (t *SmartEnergy) run(done chan error) {
|
||||
var once sync.Once
|
||||
client := request.NewHelper(t.log)
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -57,7 +56,7 @@ func (t *SmartEnergy) run(done chan error) {
|
|||
|
||||
if err := backoff.Retry(func() error {
|
||||
return backoffPermanentError(client.GetJSON(smartenergy.URI, &res))
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -80,7 +80,6 @@ func NewConfigurableFromConfig(other map[string]interface{}) (api.Tariff, error)
|
|||
|
||||
func (t *Tariff) run(forecastG func() (string, error), done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
tick := time.NewTicker(time.Hour)
|
||||
for ; true; <-tick.C {
|
||||
|
|
@ -97,7 +96,7 @@ func (t *Tariff) run(forecastG func() (string, error), done chan error) {
|
|||
data[i].Price = t.totalPrice(r.Price)
|
||||
}
|
||||
return nil
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
|
|
@ -72,7 +72,6 @@ func NewTibberFromConfig(other map[string]interface{}) (api.Tariff, error) {
|
|||
|
||||
func (t *Tibber) run(done chan error) {
|
||||
var once sync.Once
|
||||
bo := newBackoff()
|
||||
|
||||
v := map[string]interface{}{
|
||||
"id": graphql.ID(t.homeID),
|
||||
|
|
@ -94,7 +93,7 @@ func (t *Tibber) run(done chan error) {
|
|||
ctx, cancel := context.WithTimeout(context.Background(), request.Timeout)
|
||||
defer cancel()
|
||||
return t.client.Query(ctx, &res, v)
|
||||
}, bo); err != nil {
|
||||
}, bo()); err != nil {
|
||||
once.Do(func() { done <- err })
|
||||
|
||||
t.log.ERROR.Println(err)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue