diff --git a/charger/keba.go b/charger/keba.go index f886c8d28..aaf265ebf 100644 --- a/charger/keba.go +++ b/charger/keba.go @@ -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 diff --git a/charger/keba/listener.go b/charger/keba/listener.go index ce11c2769..bb2c2bdd2 100644 --- a/charger/keba/listener.go +++ b/charger/keba/listener.go @@ -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: diff --git a/charger/keba_test.go b/charger/keba_test.go index 9816e5e4d..0588968d3 100644 --- a/charger/keba_test.go +++ b/charger/keba_test.go @@ -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) }