From 95cfe455b8856dc79123e4beb75c7375754cf71a Mon Sep 17 00:00:00 2001 From: andig Date: Wed, 23 Dec 2020 17:10:55 +0100 Subject: [PATCH] Add serial to address mappings cache to Keba listener for making simple messages routable via serial (#546) #429 introduced a bug when KEBA was configured using serial number. The TCH :OK messages received from KEBA could not be routed to the subscriber since they don't contain serial numbers. This PR adds a serial<->address cache to the listener for maintaining this mapping and hence make the OK responses routable. --- charger/keba.go | 4 ++-- charger/keba/listener.go | 20 ++++++++++++++++---- charger/keba/sender.go | 7 ++++++- detect/keba.go | 2 +- 4 files changed, 25 insertions(+), 8 deletions(-) diff --git a/charger/keba.go b/charger/keba.go index f6d8625ca..b380fd959 100644 --- a/charger/keba.go +++ b/charger/keba.go @@ -58,8 +58,8 @@ func NewKebaFromConfig(other map[string]interface{}) (api.Charger, error) { func NewKeba(uri, serial string, rfid RFID, timeout time.Duration) (api.Charger, error) { log := util.NewLogger("keba") - var err error if keba.Instance == nil { + var err error keba.Instance, err = keba.New(log) if err != nil { return nil, err @@ -68,7 +68,7 @@ func NewKeba(uri, serial string, rfid RFID, timeout time.Duration) (api.Charger, // add default port conn := util.DefaultPort(uri, keba.Port) - sender, err := keba.NewSender(uri) + sender, err := keba.NewSender(log, conn) c := &Keba{ log: log, diff --git a/charger/keba/listener.go b/charger/keba/listener.go index 0c78b8775..9d0928165 100644 --- a/charger/keba/listener.go +++ b/charger/keba/listener.go @@ -40,6 +40,7 @@ type Listener struct { log *util.Logger conn *net.UDPConn clients map[string]chan<- UDPMsg + cache map[string]string } // New creates a UDP listener that clients can subscribe to @@ -58,6 +59,7 @@ func New(log *util.Logger) (*Listener, error) { log: log, conn: conn, clients: make(map[string]chan<- UDPMsg), + cache: make(map[string]string), } go l.listen() @@ -65,7 +67,7 @@ func New(log *util.Logger) (*Listener, error) { return l, nil } -// Subscribe adds a client address and message channel +// Subscribe adds a client address or serial and message channel to the list of subscribers func (l *Listener) Subscribe(addr string, c chan<- UDPMsg) { l.mux.Lock() defer l.mux.Unlock() @@ -94,7 +96,7 @@ func (l *Listener) listen() { if body != OK { var report Report if err := json.Unmarshal([]byte(body), &report); err != nil { - l.log.WARN.Printf("listener: %v", err) + l.log.WARN.Printf("recv: invalid message: %v", err) continue } @@ -105,15 +107,25 @@ func (l *Listener) listen() { } } -// addrMatches checks if either message sender or serial matched given addr +// addrMatches checks if either message sender or serial matches given addr func (l *Listener) addrMatches(addr string, msg UDPMsg) bool { switch { case addr == Any: return true + case addr == msg.Addr: return true - case msg.Report != nil && addr == msg.Report.Serial: + + // simple response like TCH :OK where cached serial for sender address matches + case l.cache[addr] == msg.Addr: return true + + // report response with matching serial + case msg.Report != nil && addr == msg.Report.Serial: + // cache address for serial to make simple TCH :OK messages routable using serial + l.cache[msg.Report.Serial] = msg.Addr + return true + default: return false } diff --git a/charger/keba/sender.go b/charger/keba/sender.go index c951e65e3..72afb5293 100644 --- a/charger/keba/sender.go +++ b/charger/keba/sender.go @@ -10,11 +10,13 @@ import ( // Sender is a KEBA UDP sender type Sender struct { + log *util.Logger + addr string conn *net.UDPConn } // NewSender creates KEBA UDP sender -func NewSender(addr string) (*Sender, error) { +func NewSender(log *util.Logger, addr string) (*Sender, error) { addr = util.DefaultPort(addr, Port) raddr, err := net.ResolveUDPAddr("udp", addr) @@ -24,6 +26,8 @@ func NewSender(addr string) (*Sender, error) { } c := &Sender{ + log: log, + addr: addr, conn: conn, } @@ -32,6 +36,7 @@ func NewSender(addr string) (*Sender, error) { // Send msg to receiver func (c *Sender) Send(msg string) error { + c.log.TRACE.Printf("send to %s %v", c.addr, msg) _, err := io.Copy(c.conn, strings.NewReader(msg)) return err } diff --git a/detect/keba.go b/detect/keba.go index ec69727b1..df9701d17 100644 --- a/detect/keba.go +++ b/detect/keba.go @@ -50,7 +50,7 @@ func (h *KEBAHandler) Test(log *util.Logger, ip string) []interface{} { resC := make(chan keba.UDPMsg) h.listener.Subscribe(ip, resC) - sender, err := keba.NewSender(ip) + sender, err := keba.NewSender(log, ip) if err != nil { log.ERROR.Println("keba:", err) return nil