Tinkerforge Warp: fix websocket credentials (#27737)
This commit is contained in:
parent
fabf42918d
commit
d6cb77ef62
3 changed files with 104 additions and 46 deletions
|
|
@ -20,7 +20,7 @@ import (
|
|||
"github.com/evcc-io/evcc/charger/warp"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/evcc-io/evcc/util/request"
|
||||
"github.com/jpfielding/go-http-digest/pkg/digest"
|
||||
"github.com/icholy/digest"
|
||||
)
|
||||
|
||||
type WarpWS struct {
|
||||
|
|
@ -56,13 +56,6 @@ type WarpWS struct {
|
|||
pmLowLevelState warp.PmLowLevelState
|
||||
}
|
||||
|
||||
type warpEvent struct {
|
||||
Topic string `json:"topic"`
|
||||
Payload json.RawMessage `json:"payload"`
|
||||
}
|
||||
|
||||
var _ api.ChargerEx = (*WarpWS)(nil)
|
||||
|
||||
func init() {
|
||||
registry.AddCtx("warp-ws", NewWarpWSFromConfig)
|
||||
}
|
||||
|
|
@ -85,7 +78,7 @@ func NewWarpWSFromConfig(ctx context.Context, other map[string]any) (api.Charger
|
|||
return nil, err
|
||||
}
|
||||
|
||||
wb, err := NewWarpWS(ctx, cc.URI, cc.User, cc.Password, cc.EnergyMeterIndex)
|
||||
wb, err := NewWarpWS(ctx, cc.URI, cc.EnergyMeterIndex, cc.User, cc.Password)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -109,15 +102,15 @@ func NewWarpWSFromConfig(ctx context.Context, other map[string]any) (api.Charger
|
|||
wb.pmURI = wb.uri
|
||||
wb.pmHelper = wb.Helper
|
||||
} else if cc.EnergyManagerURI != "" { // fallback to Energy Manager
|
||||
wb.pmURI, err = parseURI(cc.EnergyManagerURI, false)
|
||||
if wb.pmURI == "" {
|
||||
return nil, err
|
||||
} else if err != nil {
|
||||
wb.log.DEBUG.Println(err)
|
||||
}
|
||||
wb.pmURI = util.DefaultScheme(strings.TrimRight(cc.EnergyManagerURI, "/"), "http")
|
||||
wb.pmHelper = request.NewHelper(wb.log)
|
||||
|
||||
if cc.EnergyManagerUser != "" {
|
||||
wb.pmHelper.Client.Transport = digest.NewTransport(cc.EnergyManagerUser, cc.EnergyManagerPassword, wb.pmHelper.Client.Transport)
|
||||
wb.pmHelper.Client.Transport = &digest.Transport{
|
||||
Username: cc.EnergyManagerUser,
|
||||
Password: cc.EnergyManagerPassword,
|
||||
Transport: wb.pmHelper.Client.Transport,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -151,70 +144,125 @@ func NewWarpWSFromConfig(ctx context.Context, other map[string]any) (api.Charger
|
|||
return decorateWarpWS(wb, currentPower, totalEnergy, currents, voltages, identify, phases, getPhases), nil
|
||||
}
|
||||
|
||||
func NewWarpWS(ctx context.Context, uri, user, password string, meterIndex uint) (*WarpWS, error) {
|
||||
func NewWarpWS(ctx context.Context, uri string, meterIndex uint, user, password string) (*WarpWS, error) {
|
||||
log := util.NewLogger("warp-ws")
|
||||
|
||||
wsURI, err := parseURI(uri)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
client := request.NewHelper(log)
|
||||
|
||||
if user != "" {
|
||||
client.Client.Transport = digest.NewTransport(user, password, client.Client.Transport)
|
||||
client.Client.Transport = &digest.Transport{
|
||||
Username: user,
|
||||
Password: password,
|
||||
Transport: client.Client.Transport,
|
||||
}
|
||||
}
|
||||
|
||||
w := &WarpWS{
|
||||
Helper: client, log: log,
|
||||
uri: util.DefaultScheme(uri, "http"),
|
||||
uri: uri,
|
||||
meterIndex: meterIndex,
|
||||
meterMap: map[int]int{},
|
||||
metersValueIDsTopic: fmt.Sprintf("meters/%d/value_ids", meterIndex),
|
||||
metersValuesTopic: fmt.Sprintf("meters/%d/values", meterIndex),
|
||||
}
|
||||
|
||||
go w.run(ctx)
|
||||
go w.run(ctx, digest.Options{
|
||||
URI: wsURI,
|
||||
Username: user,
|
||||
Password: password,
|
||||
})
|
||||
|
||||
return w, nil
|
||||
}
|
||||
|
||||
func (w *WarpWS) run(ctx context.Context) {
|
||||
uri, err := parseURI(w.uri, true)
|
||||
if err != nil {
|
||||
w.log.DEBUG.Println(err)
|
||||
if uri == "" {
|
||||
return
|
||||
}
|
||||
}
|
||||
w.log.TRACE.Printf("connecting to %s …", uri)
|
||||
func (w *WarpWS) run(ctx context.Context, options digest.Options) {
|
||||
bo := backoff.NewExponentialBackOff(
|
||||
backoff.WithMaxElapsedTime(0),
|
||||
backoff.WithMaxInterval(30*time.Second),
|
||||
)
|
||||
|
||||
bo := backoff.NewExponentialBackOff(backoff.WithMaxElapsedTime(0))
|
||||
for ctx.Err() == nil {
|
||||
conn, _, err := websocket.Dial(ctx, uri, nil)
|
||||
w.log.DEBUG.Println("websocket: connecting")
|
||||
|
||||
conn, err := dialWebsocket(ctx, options)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
if !errors.Is(context.DeadlineExceeded, err) {
|
||||
w.log.ERROR.Printf("websocket: %v", err)
|
||||
}
|
||||
time.Sleep(bo.NextBackOff())
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(bo.NextBackOff()):
|
||||
}
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
bo.Reset()
|
||||
|
||||
if err := w.handleConnection(ctx, conn); err != nil {
|
||||
w.log.ERROR.Println(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func parseURI(uri string, toWS bool) (string, error) {
|
||||
func dialWebsocket(ctx context.Context, options digest.Options) (*websocket.Conn, error) {
|
||||
// err will be non nil if auth is needed
|
||||
conn, resp, err := websocket.Dial(ctx, options.URI, nil)
|
||||
if err == nil {
|
||||
return conn, nil
|
||||
}
|
||||
|
||||
if resp == nil || resp.StatusCode != http.StatusUnauthorized {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if options.Username == "" {
|
||||
return nil, errors.New("websocket: missing credentials")
|
||||
}
|
||||
|
||||
// extract challenge from response
|
||||
challenge, err := digest.ParseChallenge(resp.Header.Get("WWW-Authenticate"))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("websocket: %w", err)
|
||||
}
|
||||
|
||||
options.Method = "GET"
|
||||
options.Count = 1
|
||||
|
||||
cred, err := digest.Digest(challenge, options)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Dial with Digest Auth
|
||||
dialer := websocket.DialOptions{
|
||||
HTTPHeader: http.Header{
|
||||
"Authorization": []string{cred.String()},
|
||||
},
|
||||
}
|
||||
|
||||
conn, _, err = websocket.Dial(ctx, options.URI, &dialer)
|
||||
return conn, err
|
||||
}
|
||||
|
||||
// Returns parsed URI and hostname
|
||||
func parseURI(uri string) (string, error) {
|
||||
u, err := url.Parse(util.DefaultScheme(strings.TrimRight(uri, "/"), "http"))
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
if u.Scheme == "https" || u.Scheme == "wss" {
|
||||
u.Scheme = "http"
|
||||
err = fmt.Errorf("https or wss are not supported, using http/ws instead")
|
||||
}
|
||||
if toWS {
|
||||
u.Scheme = "ws"
|
||||
u.Path = path.Join(u.Path, "/ws")
|
||||
}
|
||||
return u.String(), err
|
||||
|
||||
u.Scheme = "ws"
|
||||
u.Path = path.Join(u.Path, "/ws")
|
||||
|
||||
return u.String(), nil
|
||||
}
|
||||
|
||||
func (w *WarpWS) handleConnection(ctx context.Context, conn *websocket.Conn) error {
|
||||
|
|
@ -230,7 +278,10 @@ func (w *WarpWS) handleConnection(ctx context.Context, conn *websocket.Conn) err
|
|||
|
||||
dec := json.NewDecoder(r)
|
||||
for {
|
||||
var event warpEvent
|
||||
var event struct {
|
||||
Topic string `json:"topic"`
|
||||
Payload json.RawMessage `json:"payload"`
|
||||
}
|
||||
if err := dec.Decode(&event); err != nil {
|
||||
if errors.Is(err, io.EOF) {
|
||||
break //next frame
|
||||
|
|
@ -238,7 +289,7 @@ func (w *WarpWS) handleConnection(ctx context.Context, conn *websocket.Conn) err
|
|||
return err
|
||||
}
|
||||
|
||||
w.log.TRACE.Printf("ws event %s: %s", event.Topic, event.Payload)
|
||||
w.log.TRACE.Printf("websocket: event %s: %s", event.Topic, event.Payload)
|
||||
if err := w.handleEvent(event.Topic, event.Payload); err != nil {
|
||||
w.log.ERROR.Printf("bad payload for topic %s: %v", event.Topic, err)
|
||||
}
|
||||
|
|
@ -361,6 +412,8 @@ func (w *WarpWS) MaxCurrent(current int64) error {
|
|||
return w.MaxCurrentMillis(float64(current))
|
||||
}
|
||||
|
||||
var _ api.ChargerEx = (*WarpWS)(nil)
|
||||
|
||||
// MaxCurrentMillis implements the api.ChargerEx interface
|
||||
func (w *WarpWS) MaxCurrentMillis(current float64) error {
|
||||
curr := int64(current * 1e3)
|
||||
|
|
|
|||
1
go.mod
1
go.mod
|
|
@ -55,6 +55,7 @@ require (
|
|||
github.com/hashicorp/go-version v1.8.0
|
||||
github.com/hasura/go-graphql-client v0.15.1
|
||||
github.com/holoplot/go-evdev v0.0.0-20250804134636-ab1d56a1fe83
|
||||
github.com/icholy/digest v1.1.0
|
||||
github.com/influxdata/influxdb-client-go/v2 v2.14.0
|
||||
github.com/insomniacslk/tapo v1.0.2
|
||||
github.com/itchyny/gojq v0.12.18
|
||||
|
|
|
|||
4
go.sum
4
go.sum
|
|
@ -399,6 +399,8 @@ github.com/hpcloud/tail v1.0.0/go.mod h1:ab1qPbhIpdTxEkNHXyeSf5vhxWSCs/tWer42PpO
|
|||
github.com/huandu/xstrings v1.5.0 h1:2ag3IFq9ZDANvthTwTiqSSZLjDc+BedvHPAp5tJy2TI=
|
||||
github.com/huandu/xstrings v1.5.0/go.mod h1:y5/lhBue+AyNmUVz9RLU9xbLR0o4KIIExikq4ovT0aE=
|
||||
github.com/hudl/fargo v1.3.0/go.mod h1:y3CKSmjA+wD2gak7sUSXTAoopbhU08POFhmITJgmKTg=
|
||||
github.com/icholy/digest v1.1.0 h1:HfGg9Irj7i+IX1o1QAmPfIBNu/Q5A5Tu3n/MED9k9H4=
|
||||
github.com/icholy/digest v1.1.0/go.mod h1:QNrsSGQ5v7v9cReDI0+eyjsXGUoRSUZQHeQ5C4XLa0Y=
|
||||
github.com/inconshreveable/mousetrap v1.0.0/go.mod h1:PxqpIevigyE2G7u3NXJIT2ANytuPF1OarO4DADm73n8=
|
||||
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
|
||||
github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw=
|
||||
|
|
@ -1066,6 +1068,8 @@ gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
|||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gorm.io/gorm v1.31.1 h1:7CA8FTFz/gRfgqgpeKIBcervUn3xSyPUmr6B2WXJ7kg=
|
||||
gorm.io/gorm v1.31.1/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs=
|
||||
gotest.tools/v3 v3.5.1 h1:EENdUnS3pdur5nybKYIh2Vfgc8IUNBjxDPSjtiJcOzU=
|
||||
gotest.tools/v3 v3.5.1/go.mod h1:isy3WKz7GK6uNw/sbHzfKBLvlvXwUyV06n6brMxxopU=
|
||||
honnef.co/go/tools v0.0.0-20180728063816-88497007e858/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4=
|
||||
honnef.co/go/tools v0.0.0-20190102054323-c2f93a96b099/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4=
|
||||
honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4=
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue