Refactor KEBA implementation for docker (#288)

* Handle errors creating keba listener
* Verify settings applied
* Fix tests
This commit is contained in:
andig 2020-08-15 23:08:01 +02:00 • committed by GitHub
parent 872eb58337
commit b2c2aa2008
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
3 changed files with 70 additions and 38 deletions

View file

@ -40,33 +40,36 @@ type Keba struct {
func NewKebaFromConfig(other map[string]interface{}) (api.Charger, error) {
cc := struct {
URI string
Serial string
Timeout time.Duration
RFID RFID
}{}
}{
Timeout: udpTimeout,
}
if err := util.DecodeOther(other, &cc); err != nil {
return nil, err
}
return NewKeba(cc.URI, cc.RFID, cc.Timeout)
return NewKeba(cc.URI, cc.Serial, cc.RFID, cc.Timeout)
}
// NewKeba creates a new charger
func NewKeba(conn string, rfid RFID, timeout time.Duration) (api.Charger, error) {
func NewKeba(conn, serial string, rfid RFID, timeout time.Duration) (api.Charger, error) {
log := util.NewLogger("keba")
var err error
if keba.Instance == nil {
keba.Instance = keba.New(log, fmt.Sprintf(":%s", kebaPort))
keba.Instance, err = keba.New(log, fmt.Sprintf(":%s", kebaPort))
if err != nil {
return nil, err
}
}
// add default port
if _, _, err := net.SplitHostPort(conn); err != nil {
if _, _, err = net.SplitHostPort(conn); err != nil {
conn = fmt.Sprintf("%s:%s", conn, kebaPort)
}
if timeout == 0 {
timeout = udpTimeout
}
c := &Keba{
log: log,
conn: conn,
@ -75,9 +78,12 @@ func NewKeba(conn string, rfid RFID, timeout time.Duration) (api.Charger, error)
recv: make(chan keba.UDPMsg),
}
keba.Instance.Subscribe(conn, c.recv)
// use serial to subscribe if defined for docker scenarios
if serial == "" {
serial = conn
}
return c, nil
return c, keba.Instance.Subscribe(serial, c.recv)
}
func (c *Keba) send(msg string) error {
@ -226,32 +232,43 @@ func (c *Keba) Enable(enable bool) error {
d = 1
}
// ignore result...
var resp string
err := c.roundtrip(fmt.Sprintf("ena %d", d), 0, &resp)
_ = c.roundtrip(fmt.Sprintf("ena %d", d), 0, &resp)
// ...and verify value
res, err := c.Enabled()
if err == nil && res != enable {
return fmt.Errorf("ena could not enable: %s", resp)
}
return err
}
// actualCurrent returns the actual current
func (c *Keba) actualCurrent() (int64, error) {
var kr keba.Report2
err := c.roundtrip("report 2", 2, &kr)
if err != nil {
return err
return 0, err
}
if string(resp) == keba.OK {
return nil
}
return fmt.Errorf("ena unexpected response: %s", resp)
return int64(kr.Curruser) / 1000, nil
}
// MaxCurrent implements the Charger.MaxCurrent interface
func (c *Keba) MaxCurrent(current int64) error {
// ignore result...
var resp string
err := c.roundtrip(fmt.Sprintf("curr %d", 1000*current), 0, &resp)
if err != nil {
return err
_ = c.roundtrip(fmt.Sprintf("curr %d", 1000*current), 0, &resp)
// ...and verify value
res, err := c.actualCurrent()
if err == nil && res != current {
return fmt.Errorf("curr could not set: %s", resp)
}
if resp == keba.OK {
return nil
}
return fmt.Errorf("curr unexpected response: %s", resp)
return err
}
// CurrentPower implements the Meter interface

View file

@ -2,6 +2,7 @@ package keba
import (
"encoding/json"
"fmt"
"net"
"strings"
"sync"
@ -36,37 +37,39 @@ type Listener struct {
}
// New creates a UDP listener that clients can subscribe to
func New(log *util.Logger, addr string) *Listener {
func New(log *util.Logger, addr string) (*Listener, error) {
laddr, err := net.ResolveUDPAddr("udp", addr)
if err != nil {
log.FATAL.Fatal(err)
return nil, err
}
conn, err := net.ListenUDP("udp", laddr)
if err != nil {
log.FATAL.Fatal(err)
return nil, err
}
l := &Listener{
log: log,
conn: conn,
log: log,
conn: conn,
clients: make(map[string]chan<- UDPMsg),
}
go l.listen()
return l
return l, nil
}
// Subscribe adds a client address and message channel
func (l *Listener) Subscribe(addr string, c chan<- UDPMsg) {
func (l *Listener) Subscribe(addr string, c chan<- UDPMsg) error {
l.mux.Lock()
defer l.mux.Unlock()
if l.clients == nil {
l.clients = make(map[string]chan<- UDPMsg)
if _, exists := l.clients[addr]; exists {
return fmt.Errorf("duplicate subscription: %s", addr)
}
l.clients[addr] = c
return nil
}
func (l *Listener) listen() {
@ -75,7 +78,7 @@ func (l *Listener) listen() {
for {
read, addr, err := l.conn.ReadFrom(b)
if err != nil {
l.log.WARN.Printf("listener: %v", err)
l.log.ERROR.Printf("listener: %v", err)
continue
}
@ -101,12 +104,24 @@ func (l *Listener) listen() {
}
}
// addrMatches checks if either message sender or serial matched given addr
func (l *Listener) addrMatches(addr string, msg UDPMsg) bool {
switch {
case addr == msg.Addr:
return true
case msg.Report != nil && addr == msg.Report.Serial:
return true
default:
return false
}
}
func (l *Listener) send(msg UDPMsg) {
l.mux.Lock()
defer l.mux.Unlock()
for addr, client := range l.clients {
if addr == msg.Addr {
if l.addrMatches(addr, msg) {
select {
case client <- msg:
default:

View file

@ -8,7 +8,7 @@ import (
func TestKeba(t *testing.T) {
var wb api.Charger
wb, err := NewKeba("foo", RFID{}, 0)
wb, err := NewKeba("foo", "bar", RFID{}, 0)
if err != nil {
t.Error(err)
}