From 7b437602b40c840b612b62232cd8073cf6c0faaf Mon Sep 17 00:00:00 2001 From: andig Date: Thu, 5 Dec 2024 18:11:06 +0100 Subject: [PATCH] chore: use context (#17600) --- charger/alfen.go | 21 ++++++++++++++------- charger/amperfied.go | 21 ++++++++++++++------- charger/dadapower.go | 21 ++++++++++++++------- charger/daheimladen-mb.go | 21 ++++++++++++++------- charger/delta.go | 21 ++++++++++++++------- charger/hardybarth-salia.go | 21 ++++++++++++++------- charger/heidelberg-ec.go | 21 ++++++++++++++------- charger/keba-modbus.go | 21 ++++++++++++++------- charger/mennekes-compact.go | 22 ++++++++++++++-------- charger/mypv-elwa2.go | 21 ++++++++++++++------- charger/obo.go | 21 ++++++++++++++------- charger/openwb-pro.go | 21 ++++++++++++++------- charger/pulsares.go | 21 ++++++++++++++------- charger/pulsatrix.go | 24 ++++++++++-------------- charger/schneider-v3.go | 21 ++++++++++++++------- charger/vestel.go | 21 ++++++++++++++------- charger/webasto-next.go | 21 ++++++++++++++------- 17 files changed, 234 insertions(+), 127 deletions(-) diff --git a/charger/alfen.go b/charger/alfen.go index 068f62699..dc3ff4a6f 100644 --- a/charger/alfen.go +++ b/charger/alfen.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "math" "sync" @@ -52,13 +53,13 @@ const ( ) func init() { - registry.Add("alfen", NewAlfenFromConfig) + registry.AddCtx("alfen", NewAlfenFromConfig) } //go:generate go run ../cmd/tools/decorate.go -f decorateAlfen -b *Alfen -r api.Charger -t "api.PhaseSwitcher,Phases1p3p,func(int) error" -t "api.PhaseGetter,GetPhases,func() (int, error)" // NewAlfenFromConfig creates a Alfen charger from generic config -func NewAlfenFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewAlfenFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{ ID: 1, } @@ -67,11 +68,11 @@ func NewAlfenFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewAlfen(cc.URI, cc.ID) + return NewAlfen(ctx, cc.URI, cc.ID) } // NewAlfen creates Alfen charger -func NewAlfen(uri string, slaveID uint8) (api.Charger, error) { +func NewAlfen(ctx context.Context, uri string, slaveID uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, slaveID) if err != nil { return nil, err @@ -89,7 +90,7 @@ func NewAlfen(uri string, slaveID uint8) (api.Charger, error) { conn: conn, } - go wb.heartbeat() + go wb.heartbeat(ctx) _, v2, v3, err := wb.Voltages() @@ -108,8 +109,14 @@ func NewAlfen(uri string, slaveID uint8) (api.Charger, error) { return decorateAlfen(wb, phasesS, phasesG), err } -func (wb *Alfen) heartbeat() { - for range time.Tick(25 * time.Second) { +func (wb *Alfen) heartbeat(ctx context.Context) { + for tick := time.Tick(25 * time.Second); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + wb.mu.Lock() var curr float64 if wb.enabled { diff --git a/charger/amperfied.go b/charger/amperfied.go index 6dd367c39..854d6270a 100644 --- a/charger/amperfied.go +++ b/charger/amperfied.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "encoding/hex" "fmt" @@ -54,13 +55,13 @@ const ( ) func init() { - registry.Add("amperfied", NewAmperfiedFromConfig) + registry.AddCtx("amperfied", NewAmperfiedFromConfig) } //go:generate go run ../cmd/tools/decorate.go -f decorateAmperfied -b *Amperfied -r api.Charger -t "api.PhaseSwitcher,Phases1p3p,func(int) error" -t "api.PhaseGetter,GetPhases,func() (int, error)" // NewAmperfiedFromConfig creates a Amperfied charger from generic config -func NewAmperfiedFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewAmperfiedFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := struct { modbus.TcpSettings `mapstructure:",squash"` Phases1p3p bool @@ -74,11 +75,11 @@ func NewAmperfiedFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewAmperfied(cc.URI, cc.ID, cc.Phases1p3p) + return NewAmperfied(ctx, cc.URI, cc.ID, cc.Phases1p3p) } // NewAmperfied creates Amperfied charger -func NewAmperfied(uri string, slaveID uint8, phases bool) (api.Charger, error) { +func NewAmperfied(ctx context.Context, uri string, slaveID uint8, phases bool) (api.Charger, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, slaveID) if err != nil { return nil, err @@ -103,7 +104,7 @@ func NewAmperfied(uri string, slaveID uint8, phases bool) (api.Charger, error) { return nil, fmt.Errorf("failsafe timeout: %w", err) } if u := binary.BigEndian.Uint16(b); u > 0 { - go wb.heartbeat(time.Duration(u) * time.Millisecond / 2) + go wb.heartbeat(ctx, time.Duration(u)*time.Millisecond/2) } var phases1p3p func(int) error @@ -116,8 +117,14 @@ func NewAmperfied(uri string, slaveID uint8, phases bool) (api.Charger, error) { return decorateAmperfied(wb, phases1p3p, phasesG), nil } -func (wb *Amperfied) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Amperfied) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.Status(); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/dadapower.go b/charger/dadapower.go index 2196ab47a..26751c633 100644 --- a/charger/dadapower.go +++ b/charger/dadapower.go @@ -1,6 +1,7 @@ package charger import ( + "context" "encoding/binary" "errors" "fmt" @@ -37,22 +38,22 @@ type Dadapower struct { } func init() { - registry.Add("dadapower", NewDadapowerFromConfig) + registry.AddCtx("dadapower", NewDadapowerFromConfig) } // NewDadapowerFromConfig creates a Dadapower charger from generic config -func NewDadapowerFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewDadapowerFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{} if err := util.DecodeOther(other, &cc); err != nil { return nil, err } - return NewDadapower(cc.URI, cc.ID) + return NewDadapower(ctx, cc.URI, cc.ID) } // NewDadapower creates a Dadapower charger -func NewDadapower(uri string, id uint8) (*Dadapower, error) { +func NewDadapower(ctx context.Context, uri string, id uint8) (*Dadapower, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, id) if err != nil { return nil, err @@ -80,13 +81,19 @@ func NewDadapower(uri string, id uint8) (*Dadapower, error) { wb.regOffset = (uint16(id) - 1) * 1000 } - go wb.heartbeat() + go wb.heartbeat(ctx) return wb, nil } -func (wb *Dadapower) heartbeat() { - for range time.Tick(time.Minute) { +func (wb *Dadapower) heartbeat(ctx context.Context) { + for tick := time.Tick(time.Minute); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.ReadInputRegisters(dadapowerRegFailsafeTimeout, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/daheimladen-mb.go b/charger/daheimladen-mb.go index 3890e09b4..4f5397c17 100644 --- a/charger/daheimladen-mb.go +++ b/charger/daheimladen-mb.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "fmt" "time" @@ -57,11 +58,11 @@ const ( ) func init() { - registry.Add("daheimladen-mb", NewDaheimLadenMBFromConfig) + registry.AddCtx("daheimladen-mb", NewDaheimLadenMBFromConfig) } // NewDaheimLadenMBFromConfig creates a DaheimLadenMB charger from generic config -func NewDaheimLadenMBFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewDaheimLadenMBFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{ ID: 255, } @@ -70,11 +71,11 @@ func NewDaheimLadenMBFromConfig(other map[string]interface{}) (api.Charger, erro return nil, err } - return NewDaheimLadenMB(cc.URI, cc.ID) + return NewDaheimLadenMB(ctx, cc.URI, cc.ID) } // NewDaheimLadenMB creates DaheimLadenMB charger -func NewDaheimLadenMB(uri string, id uint8) (api.Charger, error) { +func NewDaheimLadenMB(ctx context.Context, uri string, id uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, id) if err != nil { return nil, err @@ -104,14 +105,20 @@ func NewDaheimLadenMB(uri string, id uint8) (api.Charger, error) { return nil, fmt.Errorf("failsafe timeout: %w", err) } if u := binary.BigEndian.Uint16(b); u > 0 { - go wb.heartbeat(time.Duration(u) * time.Second / 2) + go wb.heartbeat(ctx, time.Duration(u)*time.Second/2) } return wb, err } -func (wb *DaheimLadenMB) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *DaheimLadenMB) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.ReadHoldingRegisters(dlRegSafeCurrent, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/delta.go b/charger/delta.go index e81c9467b..e1cd38077 100644 --- a/charger/delta.go +++ b/charger/delta.go @@ -1,6 +1,7 @@ package charger import ( + "context" "fmt" "math" "sync" @@ -59,11 +60,11 @@ const ( ) func init() { - registry.Add("delta", NewDeltaFromConfig) + registry.AddCtx("delta", NewDeltaFromConfig) } // NewDeltaFromConfig creates a Delta charger from generic config -func NewDeltaFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewDeltaFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := struct { Connector uint16 modbus.Settings `mapstructure:",squash"` @@ -78,11 +79,11 @@ func NewDeltaFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewDelta(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Connector) + return NewDelta(ctx, cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Connector) } // NewDelta creates Delta charger -func NewDelta(uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8, connector uint16) (api.Charger, error) { +func NewDelta(ctx context.Context, uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8, connector uint16) (api.Charger, error) { conn, err := modbus.NewConnection(uri, device, comset, baudrate, proto, slaveID) if err != nil { return nil, err @@ -124,15 +125,21 @@ func NewDelta(uri, device, comset string, baudrate int, proto modbus.Protocol, s return nil, fmt.Errorf("failsafe timeout: %w", err) } if u := encoding.Uint16(b); u > 0 { - go wb.heartbeat(time.Duration(u) * time.Second / 2) + go wb.heartbeat(ctx, time.Duration(u)*time.Second/2) } } return wb, nil } -func (wb *Delta) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Delta) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + wb.mu.Lock() var curr float64 if wb.enabled { diff --git a/charger/hardybarth-salia.go b/charger/hardybarth-salia.go index f5a8ff471..99defeaef 100644 --- a/charger/hardybarth-salia.go +++ b/charger/hardybarth-salia.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "errors" "fmt" "net/http" @@ -49,13 +50,13 @@ type Salia struct { } func init() { - registry.Add("hardybarth-salia", NewSaliaFromConfig) + registry.AddCtx("hardybarth-salia", NewSaliaFromConfig) } //go:generate go run ../cmd/tools/decorate.go -f decorateSalia -b *Salia -r api.Charger -t "api.Meter,CurrentPower,func() (float64, error)" -t "api.MeterEnergy,TotalEnergy,func() (float64, error)" -t "api.PhaseCurrents,Currents,func() (float64, float64, float64, error)" -t "api.PhaseSwitcher,Phases1p3p,func(int) error" -t "api.PhaseGetter,GetPhases,func() (int, error)" // NewSaliaFromConfig creates a Salia cPH2 charger from generic config -func NewSaliaFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewSaliaFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := struct { URI string Cache time.Duration @@ -67,11 +68,11 @@ func NewSaliaFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewSalia(cc.URI, cc.Cache) + return NewSalia(ctx, cc.URI, cc.Cache) } // NewSalia creates Hardy Barth charger with Salia controller -func NewSalia(uri string, cache time.Duration) (api.Charger, error) { +func NewSalia(ctx context.Context, uri string, cache time.Duration) (api.Charger, error) { log := util.NewLogger("salia") uri = strings.TrimSuffix(uri, "/") + "/api" @@ -122,7 +123,7 @@ func NewSalia(uri string, cache time.Duration) (api.Charger, error) { return nil, err } - go wb.heartbeat() + go wb.heartbeat(ctx) wb.pause(false) @@ -148,12 +149,18 @@ func NewSalia(uri string, cache time.Duration) (api.Charger, error) { return decorateSalia(wb, currentPower, totalEnergy, currents, phasesS, phasesG), nil } -func (wb *Salia) heartbeat() { +func (wb *Salia) heartbeat(ctx context.Context) { bo := backoff.NewExponentialBackOff( backoff.WithInitialInterval(5*time.Second), backoff.WithMaxInterval(time.Minute)) - for range time.Tick(30 * time.Second) { + for tick := time.Tick(30 * time.Second); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if err := backoff.Retry(func() error { return wb.post(salia.HeartBeat, "alive") }, bo); err != nil { diff --git a/charger/heidelberg-ec.go b/charger/heidelberg-ec.go index 99de7e588..b70a438a7 100644 --- a/charger/heidelberg-ec.go +++ b/charger/heidelberg-ec.go @@ -19,6 +19,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "fmt" "time" @@ -54,14 +55,14 @@ const ( ) func init() { - registry.Add("heidelberg", NewHeidelbergECFromConfig) + registry.AddCtx("heidelberg", NewHeidelbergECFromConfig) } // https://wallbox.heidelberg.com/wp-content/uploads/2021/05/EC_ModBus_register_table_20210222.pdf (newer) // https://cdn.shopify.com/s/files/1/0101/2409/9669/files/heidelberg-energy-control-modbus.pdf (older) // NewHeidelbergECFromConfig creates a HeidelbergEC charger from generic config -func NewHeidelbergECFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewHeidelbergECFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.Settings{ ID: 1, } @@ -70,11 +71,11 @@ func NewHeidelbergECFromConfig(other map[string]interface{}) (api.Charger, error return nil, err } - return NewHeidelbergEC(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID) + return NewHeidelbergEC(ctx, cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID) } // NewHeidelbergEC creates HeidelbergEC charger -func NewHeidelbergEC(uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8) (api.Charger, error) { +func NewHeidelbergEC(ctx context.Context, uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, device, comset, baudrate, proto, slaveID) if err != nil { return nil, err @@ -107,14 +108,20 @@ func NewHeidelbergEC(uri, device, comset string, baudrate int, proto modbus.Prot return nil, fmt.Errorf("failsafe timeout: %w", err) } if u := binary.BigEndian.Uint16(b) / 4; u > 0 { - go wb.heartbeat(time.Duration(u) * time.Millisecond) + go wb.heartbeat(ctx, time.Duration(u)*time.Millisecond) } return wb, nil } -func (wb *HeidelbergEC) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *HeidelbergEC) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.Status(); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/keba-modbus.go b/charger/keba-modbus.go index bdbb7d932..3f37da720 100644 --- a/charger/keba-modbus.go +++ b/charger/keba-modbus.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "encoding/hex" "fmt" @@ -62,13 +63,13 @@ const ( ) func init() { - registry.Add("keba-modbus", NewKebaFromConfig) + registry.AddCtx("keba-modbus", NewKebaFromConfig) } //go:generate go run ../cmd/tools/decorate.go -f decorateKeba -b *Keba -r api.Charger -t "api.Meter,CurrentPower,func() (float64, error)" -t "api.MeterEnergy,TotalEnergy,func() (float64, error)" -t "api.PhaseCurrents,Currents,func() (float64, float64, float64, error)" -t "api.Identifier,Identify,func() (string, error)" -t "api.StatusReasoner,StatusReason,func() (api.Reason, error)" -t "api.PhaseSwitcher,Phases1p3p,func(int) error" -t "api.PhaseGetter,GetPhases,func() (int, error)" // NewKebaFromConfig creates a new Keba ModbusTCP charger -func NewKebaFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewKebaFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := struct { embed `mapstructure:",squash"` modbus.TcpSettings `mapstructure:",squash"` @@ -82,7 +83,7 @@ func NewKebaFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - wb, err := NewKeba(cc.embed, cc.URI, cc.ID) + wb, err := NewKeba(ctx, cc.embed, cc.URI, cc.ID) if err != nil { return nil, err } @@ -130,14 +131,14 @@ func NewKebaFromConfig(other map[string]interface{}) (api.Charger, error) { } if u := binary.BigEndian.Uint32(b); u > 0 { - go wb.heartbeat(time.Duration(u) * time.Second / 2) + go wb.heartbeat(ctx, time.Duration(u)*time.Second/2) } return decorateKeba(wb, currentPower, totalEnergy, currents, identify, reason, phasesS, phasesG), nil } // NewKeba creates a new charger -func NewKeba(embed embed, uri string, slaveID uint8) (*Keba, error) { +func NewKeba(ctx context.Context, embed embed, uri string, slaveID uint8) (*Keba, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, slaveID) if err != nil { return nil, err @@ -159,8 +160,14 @@ func NewKeba(embed embed, uri string, slaveID uint8) (*Keba, error) { return wb, err } -func (wb *Keba) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Keba) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.Enabled(); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/mennekes-compact.go b/charger/mennekes-compact.go index 577e22e56..4b4c9b556 100644 --- a/charger/mennekes-compact.go +++ b/charger/mennekes-compact.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "fmt" "math" @@ -63,11 +64,11 @@ const ( ) func init() { - registry.Add("mennekes-compact", NewMennekesCompactFromConfig) + registry.AddCtx("mennekes-compact", NewMennekesCompactFromConfig) } // NewMennekesCompactFromConfig creates a new Mennekes ModbusTCP charger -func NewMennekesCompactFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewMennekesCompactFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := struct { modbus.Settings `mapstructure:",squash"` Timeout time.Duration @@ -83,11 +84,11 @@ func NewMennekesCompactFromConfig(other map[string]interface{}) (api.Charger, er return nil, err } - return NewMennekesCompact(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Timeout) + return NewMennekesCompact(ctx, cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Settings.Protocol(), cc.ID, cc.Timeout) } // NewMennekesCompact creates Mennekes charger -func NewMennekesCompact(uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8, timeout time.Duration) (api.Charger, error) { +func NewMennekesCompact(ctx context.Context, uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8, timeout time.Duration) (api.Charger, error) { conn, err := modbus.NewConnection(uri, device, comset, baudrate, proto, slaveID) if err != nil { return nil, err @@ -110,14 +111,19 @@ func NewMennekesCompact(uri, device, comset string, baudrate int, proto modbus.P } // failsafe - go wb.heartbeat(mennekesHeartbeatInterval) + go wb.heartbeat(ctx, mennekesHeartbeatInterval) return wb, err } -func (wb *MennekesCompact) heartbeat(timeout time.Duration) { - tick := time.NewTicker(timeout) - for ; true; <-tick.C { +func (wb *MennekesCompact) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.WriteSingleRegister(mennekesRegHeartbeat, mennekesHeartbeatToken); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/mypv-elwa2.go b/charger/mypv-elwa2.go index bc74c12c0..151e2f6af 100644 --- a/charger/mypv-elwa2.go +++ b/charger/mypv-elwa2.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "sync/atomic" "time" @@ -44,13 +45,13 @@ const ( ) func init() { - registry.Add("ac-elwa-2", NewMyPvElwa2FromConfig) + registry.AddCtx("ac-elwa-2", NewMyPvElwa2FromConfig) } // https://github.com/evcc-io/evcc/discussions/12761 // NewMyPvElwa2FromConfig creates a MyPvElwa2 charger from generic config -func NewMyPvElwa2FromConfig(other map[string]interface{}) (api.Charger, error) { +func NewMyPvElwa2FromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{ ID: 1, } @@ -59,11 +60,11 @@ func NewMyPvElwa2FromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewMyPvElwa2(cc.URI, cc.ID) + return NewMyPvElwa2(ctx, cc.URI, cc.ID) } // NewMyPvElwa2 creates myPV AC Elwa 2 charger -func NewMyPvElwa2(uri string, slaveID uint8) (api.Charger, error) { +func NewMyPvElwa2(ctx context.Context, uri string, slaveID uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, slaveID) if err != nil { return nil, err @@ -81,7 +82,7 @@ func NewMyPvElwa2(uri string, slaveID uint8) (api.Charger, error) { conn: conn, } - go wb.heartbeat(30 * time.Second) + go wb.heartbeat(ctx, 30*time.Second) return wb, nil } @@ -100,8 +101,14 @@ func (wb *MyPvElwa2) Features() []api.Feature { return []api.Feature{api.IntegratedDevice, api.Heating} } -func (wb *MyPvElwa2) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *MyPvElwa2) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if power := uint16(atomic.LoadUint32(&wb.power)); power > 0 { enabled, err := wb.Enabled() if err == nil && enabled { diff --git a/charger/obo.go b/charger/obo.go index 314adda87..3dc85067d 100644 --- a/charger/obo.go +++ b/charger/obo.go @@ -1,6 +1,7 @@ package charger import ( + "context" "encoding/binary" "fmt" "time" @@ -24,11 +25,11 @@ const ( ) func init() { - registry.Add("obo", NewOboFromConfig) + registry.AddCtx("obo", NewOboFromConfig) } // NewOboFromConfig creates a OBO Bettermann charger from generic config -func NewOboFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewOboFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.Settings{ Baudrate: 19200, Comset: "8E1", @@ -39,11 +40,11 @@ func NewOboFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewObo(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID) + return NewObo(ctx, cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID) } // NewObo creates OBO Bettermann charger -func NewObo(uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8) (api.Charger, error) { +func NewObo(ctx context.Context, uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, device, comset, baudrate, proto, slaveID) if err != nil { return nil, err @@ -63,7 +64,7 @@ func NewObo(uri, device, comset string, baudrate int, proto modbus.Protocol, sla return nil, fmt.Errorf("failsafe timeout: %w", err) } if u := binary.BigEndian.Uint16(b); u > 0 { - go wb.heartbeat(time.Duration(u) * time.Millisecond / 2) + go wb.heartbeat(ctx, time.Duration(u)*time.Millisecond/2) } // lightshow @@ -82,8 +83,14 @@ func NewObo(uri, device, comset string, baudrate int, proto modbus.Protocol, sla return wb, nil } -func (wb *Obo) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Obo) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.ReadHoldingRegisters(dlRegSafeCurrent, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/openwb-pro.go b/charger/openwb-pro.go index 872985987..1feec1ca2 100644 --- a/charger/openwb-pro.go +++ b/charger/openwb-pro.go @@ -1,6 +1,7 @@ package charger import ( + "context" "fmt" "strings" "time" @@ -13,7 +14,7 @@ import ( ) func init() { - registry.Add("openwbpro", NewOpenWBProFromConfig) + registry.AddCtx("openwbpro", NewOpenWBProFromConfig) } // https://openwb.de/main/?page_id=771 @@ -27,7 +28,7 @@ type OpenWBPro struct { } // NewOpenWBProFromConfig creates a OpenWBPro charger from generic config -func NewOpenWBProFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewOpenWBProFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := struct { URI string Cache time.Duration @@ -39,11 +40,11 @@ func NewOpenWBProFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewOpenWBPro(util.DefaultScheme(cc.URI, "http"), cc.Cache) + return NewOpenWBPro(ctx, util.DefaultScheme(cc.URI, "http"), cc.Cache) } // NewOpenWBPro creates OpenWBPro charger -func NewOpenWBPro(uri string, cache time.Duration) (*OpenWBPro, error) { +func NewOpenWBPro(ctx context.Context, uri string, cache time.Duration) (*OpenWBPro, error) { log := util.NewLogger("owbpro") wb := &OpenWBPro{ @@ -59,13 +60,19 @@ func NewOpenWBPro(uri string, cache time.Duration) (*OpenWBPro, error) { return res, err }, cache) - go wb.heartbeat(log) + go wb.heartbeat(ctx, log) return wb, nil } -func (wb *OpenWBPro) heartbeat(log *util.Logger) { - for range time.Tick(30 * time.Second) { +func (wb *OpenWBPro) heartbeat(ctx context.Context, log *util.Logger) { + for tick := time.Tick(30 * time.Second); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.statusG.Get(); err != nil { log.ERROR.Printf("heartbeat: %v", err) } diff --git a/charger/pulsares.go b/charger/pulsares.go index 73122ecaa..71c54e348 100644 --- a/charger/pulsares.go +++ b/charger/pulsares.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "fmt" "time" @@ -47,13 +48,13 @@ const ( ) func init() { - registry.Add("pulsares", NewPulsaresFromConfig) + registry.AddCtx("pulsares", NewPulsaresFromConfig) } //go:generate go run ../cmd/tools/decorate.go -f decoratePulsares -b *Pulsares -r api.Charger -t "api.PhaseSwitcher,Phases1p3p,func(int) error" // NewPulsaresFromConfig creates a Pulsares charger from generic config -func NewPulsaresFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewPulsaresFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.Settings{ ID: 1, } @@ -62,7 +63,7 @@ func NewPulsaresFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - wb, err := NewPulsares(cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID) + wb, err := NewPulsares(ctx, cc.URI, cc.Device, cc.Comset, cc.Baudrate, cc.Protocol(), cc.ID) if err != nil { return nil, err } @@ -76,7 +77,7 @@ func NewPulsaresFromConfig(other map[string]interface{}) (api.Charger, error) { } // NewPulsares creates Pulsares charger -func NewPulsares(uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8) (*Pulsares, error) { +func NewPulsares(ctx context.Context, uri, device, comset string, baudrate int, proto modbus.Protocol, slaveID uint8) (*Pulsares, error) { conn, err := modbus.NewConnection(uri, device, comset, baudrate, proto, slaveID) if err != nil { return nil, err @@ -120,14 +121,20 @@ func NewPulsares(uri, device, comset string, baudrate int, proto modbus.Protocol } if t > 0 { - go wb.heartbeat(t / 2) + go wb.heartbeat(ctx, t/2) } return wb, err } -func (wb *Pulsares) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Pulsares) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.ReadHoldingRegisters(pulsaresRegBackup, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/pulsatrix.go b/charger/pulsatrix.go index 87c69a548..76512a5e4 100644 --- a/charger/pulsatrix.go +++ b/charger/pulsatrix.go @@ -49,7 +49,6 @@ type Pulsatrix struct { conn *websocket.Conn uri string enabled bool - quit chan struct{} data *util.Monitor[pulsatrixData] } @@ -115,9 +114,9 @@ func (c *Pulsatrix) connectWs() error { if err := c.Enable(false); err != nil { c.log.ERROR.Println(err) } - c.quit = make(chan struct{}) + go c.wsReader() - go c.heartbeat() + go c.heartbeat(ctx) return nil } @@ -151,7 +150,6 @@ func (c *Pulsatrix) wsReader() { c.mu.Lock() c.conn.Close(websocket.StatusNormalClosure, "Reconnecting") c.conn = nil - close(c.quit) c.mu.Unlock() c.reconnectWs() @@ -195,19 +193,17 @@ func (c *Pulsatrix) parseWsMessage(messageType websocket.MessageType, message [] } // Heartbeat sends a heartbeat to the pulsatrix SECC -func (c *Pulsatrix) heartbeat() { - ticker := time.NewTicker(3 * time.Minute) - defer ticker.Stop() - - for { +func (c *Pulsatrix) heartbeat(ctx context.Context) { + for tick := time.Tick(3 * time.Minute); ; { select { - case <-ticker.C: - if err := c.Enable(c.enabled); err != nil { - c.log.ERROR.Println(err) - } - case <-c.quit: + case <-tick: + case <-ctx.Done(): return } + + if err := c.Enable(c.enabled); err != nil { + c.log.ERROR.Println(err) + } } } diff --git a/charger/schneider-v3.go b/charger/schneider-v3.go index 569a8b5a7..66c63c3da 100644 --- a/charger/schneider-v3.go +++ b/charger/schneider-v3.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "fmt" "time" @@ -55,13 +56,13 @@ const ( ) func init() { - registry.Add("schneider-v3", NewSchneiderV3FromConfig) + registry.AddCtx("schneider-v3", NewSchneiderV3FromConfig) } // https://download.schneider-electric.com/files?p_enDocType=Other+technical+guide&p_File_Name=GEX1969300-04.pdf&p_Doc_Ref=GEX1969300 // NewSchneiderV3FromConfig creates a Schneider charger from generic config -func NewSchneiderV3FromConfig(other map[string]interface{}) (api.Charger, error) { +func NewSchneiderV3FromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{ ID: 255, } @@ -70,11 +71,11 @@ func NewSchneiderV3FromConfig(other map[string]interface{}) (api.Charger, error) return nil, err } - return NewSchneiderV3(cc.URI, cc.ID) + return NewSchneiderV3(ctx, cc.URI, cc.ID) } // NewSchneiderV3 creates Schneider charger -func NewSchneiderV3(uri string, id uint8) (api.Charger, error) { +func NewSchneiderV3(ctx context.Context, uri string, id uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, id) if err != nil { return nil, err @@ -108,14 +109,20 @@ func NewSchneiderV3(uri string, id uint8) (api.Charger, error) { return nil, fmt.Errorf("heartbeat timeout: %w", err) } if u := encoding.Uint16(b); u != 2 { - go wb.heartbeat(2 * time.Second) + go wb.heartbeat(ctx, 2*time.Second) } return wb, nil } -func (wb *Schneider) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Schneider) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.WriteSingleRegister(schneiderRegLifebit, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/vestel.go b/charger/vestel.go index 55d1f0936..e38824efe 100644 --- a/charger/vestel.go +++ b/charger/vestel.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "fmt" "time" @@ -60,11 +61,11 @@ type Vestel struct { } func init() { - registry.Add("vestel", NewVestelFromConfig) + registry.AddCtx("vestel", NewVestelFromConfig) } // NewVestelFromConfig creates a Vestel charger from generic config -func NewVestelFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewVestelFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{ ID: 255, } @@ -73,11 +74,11 @@ func NewVestelFromConfig(other map[string]interface{}) (api.Charger, error) { return nil, err } - return NewVestel(cc.URI, cc.ID) + return NewVestel(ctx, cc.URI, cc.ID) } // NewVestel creates a Vestel charger -func NewVestel(uri string, id uint8) (*Vestel, error) { +func NewVestel(ctx context.Context, uri string, id uint8) (*Vestel, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, id) if err != nil { return nil, err @@ -108,13 +109,19 @@ func NewVestel(uri string, id uint8) (*Vestel, error) { if timeout < time.Second { timeout = time.Second } - go wb.heartbeat(timeout) + go wb.heartbeat(ctx, timeout) return wb, nil } -func (wb *Vestel) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *Vestel) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.WriteSingleRegister(vestelRegAlive, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) } diff --git a/charger/webasto-next.go b/charger/webasto-next.go index 4f009df33..deb7c0885 100644 --- a/charger/webasto-next.go +++ b/charger/webasto-next.go @@ -18,6 +18,7 @@ package charger // SOFTWARE. import ( + "context" "encoding/binary" "fmt" "time" @@ -51,11 +52,11 @@ const ( ) func init() { - registry.Add("webasto-next", NewWebastoNextFromConfig) + registry.AddCtx("webasto-next", NewWebastoNextFromConfig) } // NewWebastoNextFromConfig creates a WebastoNext charger from generic config -func NewWebastoNextFromConfig(other map[string]interface{}) (api.Charger, error) { +func NewWebastoNextFromConfig(ctx context.Context, other map[string]interface{}) (api.Charger, error) { cc := modbus.TcpSettings{ ID: 255, } @@ -64,11 +65,11 @@ func NewWebastoNextFromConfig(other map[string]interface{}) (api.Charger, error) return nil, err } - return NewWebastoNext(cc.URI, cc.ID) + return NewWebastoNext(ctx, cc.URI, cc.ID) } // NewWebastoNext creates WebastoNext charger -func NewWebastoNext(uri string, id uint8) (api.Charger, error) { +func NewWebastoNext(ctx context.Context, uri string, id uint8) (api.Charger, error) { conn, err := modbus.NewConnection(uri, "", "", 0, modbus.Tcp, id) if err != nil { return nil, err @@ -98,14 +99,20 @@ func NewWebastoNext(uri string, id uint8) (api.Charger, error) { return nil, fmt.Errorf("failsafe timeout: %w", err) } if u := binary.BigEndian.Uint16(b); u > 0 { - go wb.heartbeat(time.Duration(u) * time.Second / 2) + go wb.heartbeat(ctx, time.Duration(u)*time.Second/2) } return wb, err } -func (wb *WebastoNext) heartbeat(timeout time.Duration) { - for range time.Tick(timeout) { +func (wb *WebastoNext) heartbeat(ctx context.Context, timeout time.Duration) { + for tick := time.Tick(timeout); ; { + select { + case <-tick: + case <-ctx.Done(): + return + } + if _, err := wb.conn.WriteSingleRegister(tqRegLifeBit, 1); err != nil { wb.log.ERROR.Println("heartbeat:", err) }