chore: use context (#17600)

This commit is contained in:
andig 2024-12-05 18:11:06 +01:00 • committed by GitHub
parent a85fcc425b
commit 7b437602b4
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
17 changed files with 234 additions and 127 deletions

View file

@ -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 {

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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 {

View file

@ -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 {

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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 {

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}
}
}

View file

@ -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)
}

View file

@ -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)
}

View file

@ -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)
}