diff --git a/charger/keba.go b/charger/keba.go index f3e7379a3..552da8e67 100644 --- a/charger/keba.go +++ b/charger/keba.go @@ -4,10 +4,7 @@ import ( "encoding/json" "errors" "fmt" - "io" - "net" "reflect" - "strings" "time" "github.com/andig/evcc/api" @@ -19,7 +16,6 @@ import ( const ( udpTimeout = time.Second - kebaPort = 7090 ) // RFID contains access credentials @@ -34,6 +30,7 @@ type Keba struct { rfid RFID timeout time.Duration recv chan keba.UDPMsg + sender *keba.Sender } func init() { @@ -58,19 +55,20 @@ func NewKebaFromConfig(other map[string]interface{}) (api.Charger, error) { } // NewKeba creates a new charger -func NewKeba(conn, serial string, rfid RFID, timeout time.Duration) (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 { - keba.Instance, err = keba.New(log, fmt.Sprintf(":%d", kebaPort)) + keba.Instance, err = keba.New(log) if err != nil { return nil, err } } // add default port - conn = util.DefaultPort(conn, kebaPort) + conn := util.DefaultPort(uri, keba.Port) + sender, err := keba.NewSender(uri) c := &Keba{ log: log, @@ -78,6 +76,7 @@ func NewKeba(conn, serial string, rfid RFID, timeout time.Duration) (api.Charger rfid: rfid, timeout: timeout, recv: make(chan keba.UDPMsg), + sender: sender, } // use serial to subscribe if defined for docker scenarios @@ -85,24 +84,9 @@ func NewKeba(conn, serial string, rfid RFID, timeout time.Duration) (api.Charger serial = conn } - return c, keba.Instance.Subscribe(serial, c.recv) -} + keba.Instance.Subscribe(serial, c.recv) -func (c *Keba) send(msg string) error { - raddr, err := net.ResolveUDPAddr("udp", c.conn) - if err != nil { - return err - } - - conn, err := net.DialUDP("udp", nil, raddr) - if err != nil { - return err - } - - defer conn.Close() - - _, err = io.Copy(conn, strings.NewReader(msg)) - return err + return c, err } func (c *Keba) receive(report int, resC chan<- keba.UDPMsg, errC chan<- error, closeC <-chan struct{}) { @@ -140,7 +124,7 @@ func (c *Keba) roundtrip(msg string, report int, res interface{}) error { go c.receive(report, resC, errC, closeC) - if err := c.send(msg); err != nil { + if err := c.sender.Send(msg); err != nil { return err } diff --git a/charger/keba/listener.go b/charger/keba/listener.go index bb2c2bdd2..0c78b8775 100644 --- a/charger/keba/listener.go +++ b/charger/keba/listener.go @@ -13,8 +13,14 @@ import ( const ( udpBufferSize = 1024 + // Port is the KEBA UDP port + Port = 7090 + // OK is the KEBA confirmation message OK = "TCH-OK :done" + + // Any subscriber receives all messages + Any = "" ) // Instance is the KEBA listener instance @@ -37,8 +43,8 @@ type Listener struct { } // New creates a UDP listener that clients can subscribe to -func New(log *util.Logger, addr string) (*Listener, error) { - laddr, err := net.ResolveUDPAddr("udp", addr) +func New(log *util.Logger) (*Listener, error) { + laddr, err := net.ResolveUDPAddr("udp", fmt.Sprintf(":%d", Port)) if err != nil { return nil, err } @@ -60,16 +66,11 @@ func New(log *util.Logger, addr string) (*Listener, error) { } // Subscribe adds a client address and message channel -func (l *Listener) Subscribe(addr string, c chan<- UDPMsg) error { +func (l *Listener) Subscribe(addr string, c chan<- UDPMsg) { l.mux.Lock() defer l.mux.Unlock() - if _, exists := l.clients[addr]; exists { - return fmt.Errorf("duplicate subscription: %s", addr) - } - l.clients[addr] = c - return nil } func (l *Listener) listen() { @@ -78,7 +79,7 @@ func (l *Listener) listen() { for { read, addr, err := l.conn.ReadFrom(b) if err != nil { - l.log.ERROR.Printf("listener: %v", err) + l.log.TRACE.Printf("listener: %v", err) continue } @@ -107,6 +108,8 @@ 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 == Any: + return true case addr == msg.Addr: return true case msg.Report != nil && addr == msg.Report.Serial: @@ -125,7 +128,7 @@ func (l *Listener) send(msg UDPMsg) { select { case client <- msg: default: - l.log.TRACE.Println("listener: recv blocked") + l.log.TRACE.Println("recv: listener blocked") } break } diff --git a/charger/keba/sender.go b/charger/keba/sender.go new file mode 100644 index 000000000..c951e65e3 --- /dev/null +++ b/charger/keba/sender.go @@ -0,0 +1,37 @@ +package keba + +import ( + "io" + "net" + "strings" + + "github.com/andig/evcc/util" +) + +// Sender is a KEBA UDP sender +type Sender struct { + conn *net.UDPConn +} + +// NewSender creates KEBA UDP sender +func NewSender(addr string) (*Sender, error) { + addr = util.DefaultPort(addr, Port) + raddr, err := net.ResolveUDPAddr("udp", addr) + + var conn *net.UDPConn + if err == nil { + conn, err = net.DialUDP("udp", nil, raddr) + } + + c := &Sender{ + conn: conn, + } + + return c, err +} + +// Send msg to receiver +func (c *Sender) Send(msg string) error { + _, err := io.Copy(c.conn, strings.NewReader(msg)) + return err +} diff --git a/cmd/detect.go b/cmd/detect.go new file mode 100644 index 000000000..34499c598 --- /dev/null +++ b/cmd/detect.go @@ -0,0 +1,142 @@ +package cmd + +import ( + "fmt" + "net" + "os" + "strings" + + "github.com/andig/evcc/detect" + "github.com/andig/evcc/util" + "github.com/korylprince/ipnetgen" + "github.com/olekukonko/tablewriter" + "github.com/spf13/cobra" + "github.com/spf13/viper" +) + +// detectCmd represents the vehicle command +var detectCmd = &cobra.Command{ + Use: "detect [host ...] [subnet ...]", + Short: "Auto-detect compatible hardware", + Long: `Automatic discovery using detect scans the local network for available devices. +Scanning focuses on devices that are commonly used that are detectable with reasonable efforts. + +On successful detection, suggestions for EVCC configuration can be made. The suggestions should simplify +configuring EVCC but are probably not sufficient for fully automatic configuration.`, + Run: runDetect, +} + +func init() { + rootCmd.AddCommand(detectCmd) +} + +// IPsFromSubnet creates a list of ip addresses for given subnet +func IPsFromSubnet(arg string) (res []string) { + gen, err := ipnetgen.New(arg) + if err != nil { + log.FATAL.Fatal("could not create iterator") + } + + for ip := gen.Next(); ip != nil; ip = gen.Next() { + res = append(res, ip.String()) + } + + return res +} + +// ParseHostIPNet converts host or cidr into a host list +func ParseHostIPNet(arg string) (res []string) { + if ip := net.ParseIP(arg); ip != nil { + return []string{ip.String()} + } + + _, ipnet, err := net.ParseCIDR(arg) + + // simple host + if err != nil { + return []string{arg} + } + + // check subnet size + if bits, _ := ipnet.Mask.Size(); bits < 24 { + log.INFO.Println("skipping large subnet:", ipnet) + return + } + + return IPsFromSubnet(arg) +} + +func display(res []detect.Result) { + table := tablewriter.NewWriter(os.Stdout) + table.SetHeader([]string{"IP", "Hostname", "Task", "Details"}) + table.SetAutoMergeCells(true) + table.SetRowLine(true) + + for _, hit := range res { + switch hit.ID { + case detect.TaskPing, detect.TaskTCP80, detect.TaskTCP502: + continue + + default: + host := "" + hosts, err := net.LookupAddr(hit.Host) + if err == nil && len(hosts) > 0 { + host = strings.TrimSuffix(hosts[0], ".") + } + + details := "" + if hit.Details != nil { + details = fmt.Sprintf("%+v", hit.Details) + } + + // fmt.Printf("%-16s %-20s %-16s %s\n", hit.Host, host, hit.ID, details) + table.Append([]string{hit.Host, host, hit.ID, details}) + } + } + + fmt.Println("") + table.Render() + + fmt.Println(` +Please open https://github.com/andig/evcc/issues/new in your browser and copy the +results above into a new issue. Please tell us: + + 1. Is the scan result correct? + 2. If not correct: please describe your hardware setup.`) +} + +func runDetect(cmd *cobra.Command, args []string) { + util.LogLevel(viper.GetString("log"), nil) + + println(viper.GetString("log")) + fmt.Println(` +Auto detection will now start to scan the network for available devices. +Scanning focuses on devices that are commonly used that are detectable with reasonable efforts. +On successful detection, suggestions for EVCC configuration can be made. The suggestions should simplify +configuring EVCC but are probably not sufficient for fully automatic configuration.`) + fmt.Println() + + // args + var hosts []string + for _, arg := range args { + hosts = append(hosts, ParseHostIPNet(arg)...) + } + + // autodetect + if len(hosts) == 0 { + ips := util.LocalIPs() + if len(ips) == 0 { + log.FATAL.Fatal("could not find ip") + } + + myIP := ips[0] + log.INFO.Println("my ip:", myIP.IP) + + hosts = append(hosts, "127.0.0.1") + hosts = append(hosts, IPsFromSubnet(myIP.String())...) + } + + // magic happens here + res := detect.Work(log, 50, hosts) + display(res) +} diff --git a/detect/analyze.go b/detect/analyze.go new file mode 100644 index 000000000..f990f0dcb --- /dev/null +++ b/detect/analyze.go @@ -0,0 +1,95 @@ +package detect + +type Criteria map[string]interface{} + +func filter(list []Result, criteria []Criteria) (match []Result) { + for _, res := range list { + for _, criterium := range criteria { + ok := true + + for matchKey, matchVal := range criterium { + if foundVal, found := res.Attributes[matchKey]; !found || foundVal != matchVal { + ok = false + break + } + } + + if ok { + match = append(match, res) + } + } + } + + return match +} + +type TypeSummary struct { + Results []Result + Found, Unique bool +} + +type Summary struct { + Charger, Grid, PV, Charge, Battery, Meter TypeSummary +} + +func summarize(res []Result) TypeSummary { + return TypeSummary{ + Results: res, + Found: len(res) > 0, + Unique: len(res) == 1, + } +} + +const ( + tid = "task.id" + smaHttp = "details.http" +) + +func Consolidate(res []Result) Summary { + grid := filter(res, []Criteria{ + {tid: taskOpenwb}, + {tid: taskSMA, smaHttp: false}, + {tid: taskE3DC}, + {tid: taskSonnen}, + }) + + pv := filter(res, []Criteria{ + {tid: taskOpenwb}, + {tid: taskInverter}, + }) + + battery := filter(res, []Criteria{ + {tid: taskOpenwb}, + {tid: taskE3DC}, + {tid: taskSonnen}, + {tid: taskBattery}, + }) + + charger := filter(res, []Criteria{ + {tid: taskOpenwb}, + {tid: taskWallbe}, + {tid: taskPhoenixEMCP}, + {tid: taskEVSEWifi}, + {tid: taskGoE}, + {tid: taskKEBA}, + }) + + charge := filter(res, []Criteria{ + {tid: taskOpenwb}, + {tid: taskKEBA}, + }) + + meter := filter(res, []Criteria{ + {tid: taskSMA, smaHttp: true}, + {tid: taskMeter}, + }) + + return Summary{ + Grid: summarize(grid), + PV: summarize(pv), + Battery: summarize(battery), + Charger: summarize(charger), + Charge: summarize(charge), + Meter: summarize(meter), + } +} diff --git a/detect/definitions.go b/detect/definitions.go new file mode 100644 index 000000000..242d59684 --- /dev/null +++ b/detect/definitions.go @@ -0,0 +1,228 @@ +package detect + +import "time" + +var ( + taskList = &TaskList{} + + sunspecIDs = []int{1, 2, 3, 71, 126} // modbus ids + chargeStatus = []int{0x41, 0x42, 0x43} // status values A..C +) + +const timeout = 200 * time.Millisecond + +// public task ids +const ( + TaskPing = "ping" + TaskTCP80 = "tcp_80" + TaskTCP502 = "tcp_502" + TaskSunspec = "sunspec" +) + +// private task ids +const ( + taskOpenwb = "openwb" + taskSMA = "sma" + taskKEBA = "KEBA" + taskE3DC = "e3dc_simple" + taskSonnen = "sonnen" + taskPowerwall = "powerwall" + taskWallbe = "wallbe" + taskPhoenixEMCP = "em-cp" + taskEVSEWifi = "evsewifi" + taskGoE = "go-e" + taskInverter = "inverter" + taskBattery = "battery" + taskMeter = "meter" +) + +func init() { + taskList.Add(Task{ + ID: taskSMA, + Type: "sma", + }) + + taskList.Add(Task{ + ID: taskKEBA, + Type: "keba", + }) + + taskList.Add(Task{ + ID: TaskPing, + Type: "ping", + }) + + taskList.Add(Task{ + ID: TaskTCP502, + Type: "tcp", + Depends: TaskPing, + Config: map[string]interface{}{ + "port": 502, + }, + }) + + taskList.Add(Task{ + ID: TaskSunspec, + Type: "modbus", + Depends: TaskTCP502, + Config: map[string]interface{}{ + "ids": sunspecIDs, + "models": []int{1}, + "point": "Mn", + }, + }) + + taskList.Add(Task{ + ID: taskInverter, + Type: "modbus", + Depends: TaskSunspec, + Config: map[string]interface{}{ + "ids": sunspecIDs, + "models": []int{101, 103}, + "point": "W", + "invalid": []int{0xFFFF}, + }, + }) + + taskList.Add(Task{ + ID: taskBattery, + Type: "modbus", + Depends: TaskSunspec, + Config: map[string]interface{}{ + "ids": sunspecIDs, + "models": []int{124}, + "point": "ChaSt", + "invalid": []int{0xFFFF}, + }, + }) + + taskList.Add(Task{ + ID: taskMeter, + Type: "modbus", + Depends: TaskSunspec, + Config: map[string]interface{}{ + "ids": sunspecIDs, + "models": []int{201, 203}, + "point": "W", + }, + }) + + taskList.Add(Task{ + ID: taskE3DC, + Type: "modbus", + Depends: TaskTCP502, + Config: map[string]interface{}{ + "ids": []int{1, 2, 3, 4, 5, 6}, + "address": 40000, + "type": "holding", + "decode": "uint16", + "values": []int{0xE3DC}, + }, + }) + + taskList.Add(Task{ + ID: taskWallbe, + Type: "modbus", + Depends: TaskTCP502, + Config: map[string]interface{}{ + "ids": []int{255}, + "address": 100, + "type": "input", + "decode": "uint16", + "values": chargeStatus, + }, + }) + + taskList.Add(Task{ + ID: taskPhoenixEMCP, + Type: "modbus", + Depends: TaskTCP502, + Config: map[string]interface{}{ + "ids": []int{180}, + "address": 100, + "type": "input", + "decode": "uint16", + "values": chargeStatus, + }, + }) + + taskList.Add(Task{ + ID: taskOpenwb, + Type: "mqtt", + Depends: TaskPing, + Config: map[string]interface{}{ + "topic": "openWB", + }, + }) + + taskList.Add(Task{ + ID: TaskTCP80, + Type: "tcp", + Depends: TaskPing, + Config: map[string]interface{}{ + "port": 80, + }, + }) + + taskList.Add(Task{ + ID: taskGoE, + Type: "http", + Depends: TaskTCP80, + Config: map[string]interface{}{ + "path": "/status", + "jq": ".car", + }, + }) + + taskList.Add(Task{ + ID: taskEVSEWifi, + Type: "http", + Depends: TaskTCP80, + Config: map[string]interface{}{ + "path": "/getParameters", + "jq": ".type", + }, + }) + + taskList.Add(Task{ + ID: taskSonnen, + Type: "http", + Depends: TaskPing, + Config: map[string]interface{}{ + "port": 8080, + "path": "/api/v1/status", + "jq": ".GridFeedIn_W", + }, + }) + + taskList.Add(Task{ + ID: taskPowerwall, + Type: "http", + Depends: TaskTCP80, + Config: map[string]interface{}{ + "path": "/api/meters/aggregates", + "jq": ".load", + }, + }) + + // // see https://github.com/andig/evcc-config/pull/5/files + // taskList.Add(Task{ + // ID: "fronius", + // Type: "http", + // Depends: TaskTCP80, + // Config: map[string]interface{}{ + // "path": "/solar_api/v1/GetPowerFlowRealtimeData.fcgi", + // "jq": ".Body.Data.Site.P_Grid", + // }, + // }) + + // taskList.Add(Task{ + // ID: "volkszähler", + // Type: "http", + // Depends: TaskTCP80, + // Config: map[string]interface{}{ + // "path": "/middleware.php/entity.json", + // "timeout": 500 * time.Millisecond, + // }, + // }) +} diff --git a/detect/http.go b/detect/http.go new file mode 100644 index 000000000..ce08c7d8a --- /dev/null +++ b/detect/http.go @@ -0,0 +1,117 @@ +package detect + +import ( + "fmt" + "io/ioutil" + "net/http" + "strings" + "time" + + "github.com/andig/evcc/util" + "github.com/andig/evcc/util/jq" + "github.com/itchyny/gojq" +) + +func init() { + registry.Add("http", HttpHandlerFactory) +} + +type HttpResult struct { + Jq interface{} +} + +func HttpHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := HttpHandler{ + Schema: "http", + Port: 80, + Method: "GET", + Codes: []int{200}, + Header: map[string]string{ + "Content-type": "application/json", + }, + Timeout: 3 * timeout, + } + + err := util.DecodeOther(conf, &handler) + + if !(handler.Schema == "http" && handler.Port == 80 || + handler.Schema == "https" && handler.Port == 443) { + handler.optionalPort = fmt.Sprintf(":%d", handler.Port) + } + + if handler.Jq != "" { + query, err := gojq.Parse(handler.Jq) + if err != nil { + return nil, fmt.Errorf("invalid jq query: %s (%s)", handler.Jq, err) + } + + handler.query = query + } + + return &handler, err +} + +type HttpHandler struct { + query *gojq.Query + Port int + optionalPort string + Schema, Method, Path string + Codes []int + Header map[string]string + Jq string + Timeout time.Duration +} + +func (h *HttpHandler) Test(log *util.Logger, ip string) []interface{} { + uri := fmt.Sprintf("%s://%s%s/%s", h.Schema, ip, h.optionalPort, strings.TrimLeft(h.Path, "/")) + req, err := http.NewRequest(strings.ToUpper(h.Method), uri, nil) + if err != nil { + return nil + } + + client := http.Client{ + Timeout: h.Timeout, + } + + resp, err := client.Do(req) + if err != nil { + return nil + } + + defer resp.Body.Close() + + if len(h.Codes) > 0 { + var status bool + for _, code := range h.Codes { + if resp.StatusCode == code { + status = true + break + } + } + + if !status { + return nil + } + } + + body, err := ioutil.ReadAll(resp.Body) + if err != nil { + return nil + } + + var res HttpResult + if h.query != nil { + val, err := jq.Query(h.query, body) + res.Jq = val + + if val == nil || err != nil { + return nil + } + } + + if err == nil { + return []interface{}{res} + } + + return nil +} diff --git a/detect/keba.go b/detect/keba.go new file mode 100644 index 000000000..ba7c52cc7 --- /dev/null +++ b/detect/keba.go @@ -0,0 +1,88 @@ +package detect + +import ( + "sync" + "time" + + "github.com/andig/evcc/charger/keba" + "github.com/andig/evcc/util" +) + +type KebaResult struct { + Addr, Serial string +} + +func init() { + registry.Add("keba", KEBAHandlerFactory) +} + +func KEBAHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := KEBAHandler{ + Timeout: 5 * timeout, + } + + err := util.DecodeOther(conf, &handler) + + return &handler, err +} + +type KEBAHandler struct { + mux sync.Mutex + listener *keba.Listener + Timeout time.Duration +} + +var instance *keba.Listener + +func (h *KEBAHandler) Test(log *util.Logger, ip string) []interface{} { + h.mux.Lock() + + if h.listener == nil { + var err error + if h.listener, err = keba.New(log); err != nil { + log.ERROR.Println("keba:", err) + } + + h.mux.Unlock() + return nil + } + + h.mux.Unlock() + + resC := make(chan keba.UDPMsg) + h.listener.Subscribe(ip, resC) + + sender, err := keba.NewSender(ip) + if err != nil { + log.ERROR.Println("keba:", err) + return nil + } + + timer := time.NewTimer(h.Timeout) +WAIT: + for { + go func() { + _ = sender.Send("report 1") + }() + + select { + case t := <-resC: + log.INFO.Println(t) + if t.Report == nil { + continue + } + + r := KebaResult{ + Addr: t.Addr, + Serial: t.Report.Serial, + } + + return []interface{}{r} + + case <-timer.C: + break WAIT + } + } + + return nil +} diff --git a/detect/modbus.go b/detect/modbus.go new file mode 100644 index 000000000..ef3fcfb37 --- /dev/null +++ b/detect/modbus.go @@ -0,0 +1,202 @@ +package detect + +import ( + "encoding/binary" + "errors" + "fmt" + "time" + + "github.com/andig/evcc/util" + "github.com/andig/evcc/util/modbus" + gridx "github.com/grid-x/modbus" + "github.com/volkszaehler/mbmd/meters" + "github.com/volkszaehler/mbmd/meters/rs485" + "github.com/volkszaehler/mbmd/meters/sunspec" +) + +func init() { + registry.Add("modbus", ModbusHandlerFactory) +} + +type ModbusResult struct { + SlaveID uint8 + Model int + Point string + Value interface{} +} + +func (r *ModbusResult) Configuration(handler TaskHandler, res Result) map[string]interface{} { + port := handler.(*ModbusHandler).Port + cc := map[string]interface{}{ + "uri": fmt.Sprintf("%s:%d", res.Host, port), + "model": "sunspec", + "id": r.SlaveID, + } + + return cc +} + +func ModbusHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := ModbusHandler{ + Port: 502, + IDs: []uint8{1}, + Models: []int{1}, + Point: "Md", // Model + Timeout: 10 * timeout, + } + + err := util.DecodeOther(conf, &handler) + + if err == nil && len(handler.IDs) == 0 { + err = errors.New("missing slave IDs") + } + + if handler.Register.Address > 0 { + handler.op, err = modbus.RegisterOperation(handler.Register) + } + + return &handler, err +} + +type ModbusHandler struct { + Port int + IDs []uint8 + Models []int + Point string + Register modbus.Register `mapstructure:",squash"` + Values []int + Invalid []int + op rs485.Operation + Timeout time.Duration +} + +func (h *ModbusHandler) testRegister(log *util.Logger, conn gridx.Client) bool { + var bytes []byte + var err error + + switch h.op.FuncCode { + case rs485.ReadHoldingReg: + bytes, err = conn.ReadHoldingRegisters(h.op.OpCode, h.op.ReadLen) + case rs485.ReadInputReg: + bytes, err = conn.ReadInputRegisters(h.op.OpCode, h.op.ReadLen) + } + + if err != nil { + return false + } + + if len(h.Values) == 0 { + return true + } + + var u uint64 + switch h.op.ReadLen { + case 1: + u = uint64(binary.BigEndian.Uint16(bytes)) + case 2: + u = uint64(binary.BigEndian.Uint32(bytes)) + case 4: + u = binary.BigEndian.Uint64(bytes) + } + + for _, val := range h.Values { + if u == uint64(val) { + return true + } + } + + return false +} + +func (h *ModbusHandler) testSunSpec(log *util.Logger, conn meters.Connection, dev *sunspec.SunSpec, mr *ModbusResult) bool { + err := dev.Initialize(conn.ModbusClient()) + if errors.Is(err, meters.ErrPartiallyOpened) { + err = nil + } + if err != nil { + return false + } + + if len(h.Models) == 0 { + return true + } + + for _, model := range h.Models { + _, res, err := dev.QueryPointAny( + conn.ModbusClient(), + model, + 0, + h.Point, + ) + + if err == nil { + mr.Model = model + mr.Point = h.Point + mr.Value = res.Value() + + log.TRACE.Printf("model %d point %s: %v", model, mr.Point, mr.Value) + + if len(h.Invalid) == 0 { + return true + } + + var val int + switch typ := res.Type(); typ { + case "int16": + val = int(res.Int16()) + case "uint16": + val = int(res.Uint16()) + case "enum16": + val = int(res.Enum16()) + default: + panic("invalid point type: " + typ) + } + + for _, inv := range h.Invalid { + if val != inv { + return true + } + } + } else { + log.DEBUG.Printf("model %d: %v", model, err) + } + } + + return false +} + +func (h *ModbusHandler) Test(log *util.Logger, ip string) (res []interface{}) { + addr := fmt.Sprintf("%s:%d", ip, h.Port) + conn := meters.NewTCP(addr) + dev := sunspec.NewDevice("sunspec") + + defer conn.Close() + + conn.Logger(log.TRACE) + conn.Timeout(h.Timeout) + + for _, slaveID := range h.IDs { + // grace period for id switch + conn.Slave(slaveID) + time.Sleep(100 * time.Millisecond) + + mr := ModbusResult{ + SlaveID: slaveID, + } + + var ok bool + if h.op.OpCode > 0 { + log.DEBUG.Printf("slave id: %d op: %v", slaveID, h.op) + ok = h.testRegister(log, conn.ModbusClient()) + } else { + log.DEBUG.Printf("slave id: %d models: %v", slaveID, h.Models) + ok = h.testSunSpec(log, conn, dev, &mr) + } + + if ok { + res = append(res, mr) + } + } + + return res +} diff --git a/detect/mqtt.go b/detect/mqtt.go new file mode 100644 index 000000000..21e35794c --- /dev/null +++ b/detect/mqtt.go @@ -0,0 +1,75 @@ +package detect + +import ( + "errors" + "fmt" + "time" + + "github.com/andig/evcc/util" + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +func init() { + registry.Add("mqtt", MqttHandlerFactory) +} + +func MqttHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := MqttHandler{ + Port: 1883, + Timeout: timeout, + } + err := util.DecodeOther(conf, &handler) + + if err == nil && handler.Port == 0 { + err = errors.New("missing port") + } + + return &handler, err +} + +type MqttHandler struct { + Port int + Topic string + Timeout time.Duration +} + +func (h *MqttHandler) Test(log *util.Logger, ip string) []interface{} { + broker := fmt.Sprintf("%s:%d", ip, h.Port) + + opt := mqtt.NewClientOptions() + opt.AddBroker(broker) + opt.SetConnectTimeout(timeout) + + client := mqtt.NewClient(opt) + + var ok bool + token := client.Connect() + if token.Wait() { + ok = token.Error() == nil + } + + if ok && h.Topic != "" { + recv := make(chan bool, 1) + _ = client.Subscribe(h.Topic, 1, func(mqtt.Client, mqtt.Message) { + recv <- true + }) + + timer := time.NewTimer(timeout) + WAIT: + for { + select { + case <-recv: + break WAIT + case <-timer.C: + ok = false + break WAIT + } + } + } + + if ok { + return []interface{}{nil} + } + + return nil +} diff --git a/detect/ping.go b/detect/ping.go new file mode 100644 index 000000000..e9198e066 --- /dev/null +++ b/detect/ping.go @@ -0,0 +1,55 @@ +package detect + +import ( + "time" + + "github.com/andig/evcc/util" + "github.com/go-ping/ping" +) + +func init() { + registry.Add("ping", PingHandlerFactory) +} + +func PingHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := PingHandler{ + Count: 1, + Timeout: timeout, + } + + err := util.DecodeOther(conf, &handler) + + return &handler, err +} + +type PingHandler struct { + Count int + Timeout time.Duration +} + +func (h *PingHandler) Test(log *util.Logger, ip string) (res []interface{}) { + pinger, err := ping.NewPinger(ip) + if err != nil { + panic(err) + } + + pinger.Count = h.Count + pinger.Timeout = h.Timeout + + if err = pinger.Run(); err != nil { + log.FATAL.Println("ping:", err) + log.FATAL.Println("") + log.FATAL.Println("In order to run evcc in discovery mode, make sure to allow ping:") + log.FATAL.Println("") + log.FATAL.Println(" sudo sysctl -w net.ipv4.ping_group_range=\"0 2147483647\"") + log.FATAL.Fatalln("") + } + + stat := pinger.Statistics() + + if stat.PacketsRecv == 0 { + return nil + } + + return []interface{}{nil} +} diff --git a/detect/registry.go b/detect/registry.go new file mode 100644 index 000000000..7bfc65dfe --- /dev/null +++ b/detect/registry.go @@ -0,0 +1,30 @@ +package detect + +import ( + "fmt" + + "github.com/andig/evcc/util" +) + +type TaskHandler interface { + Test(log *util.Logger, ip string) []interface{} +} + +type TaskHandlerRegistry map[string]func(map[string]interface{}) (TaskHandler, error) + +var registry TaskHandlerRegistry = make(map[string]func(map[string]interface{}) (TaskHandler, error)) + +func (r TaskHandlerRegistry) Add(name string, factory func(map[string]interface{}) (TaskHandler, error)) { + if _, exists := r[name]; exists { + panic(fmt.Sprintf("cannot register duplicate charger type: %s", name)) + } + r[name] = factory +} + +func (r TaskHandlerRegistry) Get(name string) (func(map[string]interface{}) (TaskHandler, error), error) { + factory, exists := r[name] + if !exists { + return nil, fmt.Errorf("charger type not registered: %s", name) + } + return factory, nil +} diff --git a/detect/sma.go b/detect/sma.go new file mode 100644 index 000000000..0148d95ac --- /dev/null +++ b/detect/sma.go @@ -0,0 +1,102 @@ +package detect + +import ( + "crypto/tls" + "fmt" + "net/http" + "sync" + "time" + + "github.com/andig/evcc/meter/sma" + "github.com/andig/evcc/util" +) + +type SmaResult struct { + Addr, Serial string + Http bool +} + +func init() { + registry.Add("sma", SMAHandlerFactory) +} + +func SMAHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := SMAHandler{ + Timeout: 5 * time.Second, + } + + err := util.DecodeOther(conf, &handler) + + return &handler, err +} + +type SMAHandler struct { + mux sync.Mutex + listener *sma.Listener + Timeout time.Duration +} + +func (h *SMAHandler) httpAvailable(ip string) bool { + uri := fmt.Sprintf("https://%s", ip) + + client := http.Client{ + Timeout: time.Second, + Transport: &http.Transport{ + TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, + }, + } + + resp, err := client.Get(uri) + if err != nil { + return false + } + + resp.Body.Close() + return true +} + +func (h *SMAHandler) Test(log *util.Logger, ip string) (res []interface{}) { + h.mux.Lock() + + if h.listener != nil { + h.mux.Unlock() + return nil + } + + var err error + if h.listener, err = sma.New(log); err != nil { + log.ERROR.Println("sma:", err) + return nil + } + h.mux.Unlock() + + resC := make(chan sma.Telegram) + h.listener.Subscribe(sma.Any, resC) + + timer := time.NewTimer(h.Timeout) +WAIT: + for { + select { + case t := <-resC: + // eliminate duplicates + for _, r := range res { + if r.(SmaResult).Serial == t.Serial { + continue WAIT + } + } + + r := SmaResult{ + Addr: t.Addr, + Serial: t.Serial, + Http: h.httpAvailable(t.Addr), + } + + res = append(res, r) + + case <-timer.C: + break WAIT + } + } + + return res +} diff --git a/detect/tasklist.go b/detect/tasklist.go new file mode 100644 index 000000000..44233ec8b --- /dev/null +++ b/detect/tasklist.go @@ -0,0 +1,125 @@ +package detect + +import ( + "fmt" + "sync" + + "github.com/andig/evcc/util" + "github.com/thoas/go-funk" +) + +type Task struct { + ID, Type string + Depends string + Config map[string]interface{} +} + +type TaskList struct { + tasks []Task + handlers []TaskHandler + once sync.Once +} + +func (l *TaskList) Add(task Task) { + l.tasks = append(l.tasks, task) +} + +func (l *TaskList) Count() int { + return len(l.tasks) +} + +func (l *TaskList) delete(i int) { + if len(l.tasks) == 1 { + l.tasks = nil + } + + res := l.tasks[:i] + if i < len(l.tasks)-1 { + res = append(res, l.tasks[i+1:]...) + } + + l.tasks = res +} + +func (l *TaskList) sort() { + var res []Task + + for len(l.tasks) > 0 { + last := len(l.tasks) + + NEXT: + for i, task := range l.tasks { + if task.Depends == "" { + res = append(res, task) + l.delete(i) + break NEXT + } + + for _, sortedTask := range res { + if task.Depends == sortedTask.ID { + res = append(res, task) + l.delete(i) + break NEXT + } + } + } + + if last == len(l.tasks) { + panic("tasks with unmatched dependencies: " + fmt.Sprintf("%v", l)) + } + } + + l.tasks = res +} + +func (l *TaskList) createHandlers() { + for _, task := range l.tasks { + factory, err := registry.Get(task.Type) + if err != nil { + panic("invalid task type " + task.Type) + } + + handler, err := factory(task.Config) + if err != nil { + panic("invalid config: " + err.Error()) + } + + l.handlers = append(l.handlers, handler) + } +} + +func (l *TaskList) Test(log *util.Logger, ip string) (res []Result) { + l.once.Do(func() { + l.sort() + l.createHandlers() + }) + + failed := make([]string, 0) + +HANDLERS: + for id, handler := range l.handlers { + task := l.tasks[id] + + if funk.ContainsString(failed, task.Depends) { + continue HANDLERS + } + + results := handler.Test(log, ip) + if len(results) > 0 { + log.INFO.Printf("ip: %s task: %s ok", ip, task.ID) + + for _, detail := range results { + res = append(res, Result{ + Task: task, + Host: ip, + Details: detail, + }) + } + } else { + // log.INFO.Printf("ip: %s task: %s nok", ip, task.ID) + failed = append(failed, task.ID) + } + } + + return res +} diff --git a/detect/tcp.go b/detect/tcp.go new file mode 100644 index 000000000..e1d53913d --- /dev/null +++ b/detect/tcp.go @@ -0,0 +1,48 @@ +package detect + +import ( + "errors" + "fmt" + "net" + "time" + + "github.com/andig/evcc/util" +) + +func init() { + registry.Add("tcp", TcpHandlerFactory) +} + +func TcpHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { + handler := TcpHandler{ + Timeout: timeout, + } + err := util.DecodeOther(conf, &handler) + + if err == nil && handler.Port == 0 { + err = errors.New("missing port") + } + + handler.dialer = net.Dialer{Timeout: handler.Timeout} + return &handler, err +} + +type TcpHandler struct { + Port int + Timeout time.Duration + dialer net.Dialer +} + +func (h *TcpHandler) Test(log *util.Logger, ip string) []interface{} { + addr := fmt.Sprintf("%s:%d", ip, h.Port) + conn, err := h.dialer.Dial("tcp", addr) + if err == nil { + defer conn.Close() + } + + if err == nil { + return []interface{}{nil} + } + + return nil +} diff --git a/detect/work.go b/detect/work.go new file mode 100644 index 000000000..73531cebd --- /dev/null +++ b/detect/work.go @@ -0,0 +1,92 @@ +package detect + +import ( + "sort" + "strings" + "sync" + + "github.com/andig/evcc/util" + "github.com/fatih/structs" + "github.com/jeremywohl/flatten" +) + +type Result struct { + Task + Host string + Details interface{} + Attributes map[string]interface{} +} + +func workers(log *util.Logger, num int, tasks <-chan string, hits chan<- []Result) *sync.WaitGroup { + var wg sync.WaitGroup + for i := 0; i < num; i++ { + wg.Add(1) + go func() { + workunit(log, tasks, hits) + wg.Done() + }() + } + + return &wg +} + +func workunit(log *util.Logger, tasks <-chan string, hits chan<- []Result) { + for ip := range tasks { + res := taskList.Test(log, ip) + hits <- res + } +} + +func Work(log *util.Logger, num int, hosts []string) []Result { + tasks := make(chan string) + hits := make(chan []Result) + done := make(chan struct{}) + + wg := workers(log, num, tasks, hits) + + var res []Result + go func() { + for hits := range hits { + res = append(res, hits...) + } + done <- struct{}{} + }() + + for _, host := range hosts { + tasks <- host + } + + close(tasks) + wg.Wait() + + close(hits) + <-done + + return postProcess(res) +} + +func postProcess(res []Result) []Result { + for idx, hit := range res { + if sma, ok := hit.Details.(SmaResult); ok { + hit.Host = sma.Addr + } + + hit.Attributes = make(map[string]interface{}) + flat, _ := flatten.Flatten(structs.Map(hit), "", flatten.DotStyle) + for k, v := range flat { + hit.Attributes[strings.ToLower(k)] = v + } + + res[idx] = hit + } + + // sort by host + sort.Slice(res, func(i, j int) bool { + if res[i].Host == res[j].Host { + return res[i].Type < res[j].Type + } + return res[i].Host < res[j].Host + }) + + return res +} diff --git a/go.mod b/go.mod index d19041504..b7f009fab 100644 --- a/go.mod +++ b/go.mod @@ -11,6 +11,8 @@ require ( github.com/containrrr/shoutrrr v0.0.0-20200721140131-bafc331a1968 github.com/denisbrodbeck/machineid v1.0.1 github.com/eclipse/paho.mqtt.golang v1.2.0 + github.com/fatih/structs v1.1.0 + github.com/go-ping/ping v0.0.0-20201022122018-3977ed72668a github.com/go-telegram-bot-api/telegram-bot-api v4.6.4+incompatible github.com/godbus/dbus/v5 v5.0.3 github.com/golang/mock v1.4.3 @@ -20,15 +22,17 @@ require ( github.com/gorilla/mux v1.8.0 github.com/gorilla/websocket v1.4.2 github.com/gregdel/pushover v0.0.0-20200416074932-c8ad547caed4 - github.com/grid-x/modbus v0.0.0-20200831145459-cb26bc3b5d3d // indirect + github.com/grid-x/modbus v0.0.0-20200831145459-cb26bc3b5d3d github.com/hashicorp/go-version v1.2.1 github.com/imdario/mergo v0.3.11 github.com/influxdata/influxdb-client-go/v2 v2.1.0 github.com/itchyny/gojq v0.11.0 + github.com/jeremywohl/flatten v1.0.1 github.com/joeshaw/carwings v0.0.0-20191118152321-61b46581307a github.com/jsgoecke/tesla v0.0.0-20200530171421-e02ebd220e5a github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51 github.com/koron/go-ssdp v0.0.0-20191105050749-2e1c40ed0b5d + github.com/korylprince/ipnetgen v0.0.0-20160712025547-34a2abb431f4 github.com/kr/pretty v0.2.0 // indirect github.com/lorenzodonini/ocpp-go v0.12.0 github.com/lunixbochs/struc v0.0.0-20200707160740-784aaebc1d40 @@ -37,15 +41,17 @@ require ( github.com/muka/go-bluetooth v0.0.0-20200619025933-f6113f7141c5 github.com/mxschmitt/golang-combinations v1.0.0 github.com/nirasan/go-oauth-pkce-code-verifier v0.0.0-20170819232839-0fbfe93532da + github.com/olekukonko/tablewriter v0.0.4 github.com/spf13/cobra v1.0.0 github.com/spf13/jwalterweatherman v1.1.0 github.com/spf13/pflag v1.0.5 github.com/spf13/viper v1.7.1 github.com/stretchr/testify v1.6.1 // indirect github.com/tcnksm/go-latest v0.0.0-20170313132115-e3007ae9052e + github.com/thoas/go-funk v0.7.0 github.com/tv42/httpunix v0.0.0-20191220191345-2ba4b9c3382c - github.com/volkszaehler/mbmd v0.0.0-20200831092453-b235d6a65b21 - golang.org/x/net v0.0.0-20200707034311-ab3426394381 + github.com/volkszaehler/mbmd v0.0.0-20201115202927-ff826598e117 + golang.org/x/net v0.0.0-20200904194848-62affa334b73 gopkg.in/ini.v1 v1.57.0 gopkg.in/yaml.v3 v3.0.0-20200605160147-a5ece683394c ) diff --git a/go.sum b/go.sum index ee02a9be0..03c9dd9c6 100644 --- a/go.sum +++ b/go.sum @@ -19,14 +19,13 @@ github.com/PuerkitoBio/goquery v1.5.1 h1:PSPBGne8NIUWw+/7vFBV+kG2J/5MOjbzc7154Oa github.com/PuerkitoBio/goquery v1.5.1/go.mod h1:GsLWisAFVj4WgDibEWF4pvYnkVQBpKBKeU+7zCJoLcc= github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= +github.com/alvaroloes/enumer v1.1.2 h1:5khqHB33TZy1GWCO/lZwcroBFh7u+0j40T83VUbfAMY= github.com/alvaroloes/enumer v1.1.2/go.mod h1:FxrjvuXoDAx9isTJrv4c+T410zFi0DtXIT0m65DJ+Wo= github.com/andig/evcc v0.0.0-20200727161511-d58eb15f2dc9/go.mod h1:8HONEC6cC2s4k0u3QL7GIjrYOZYTOKiiXybw0FIJL0A= github.com/andig/evcc-config v0.0.0-20201030211256-8664d1cdae95 h1:GsXyGiwc3gf1XB5pj80llnbxsB9fqjFC2JKdooQSrs0= github.com/andig/evcc-config v0.0.0-20201030211256-8664d1cdae95/go.mod h1:N0hIjIy+5E2AR1fF7Tg2IzBlblBrnFvCCaDGAaHzbWk= github.com/andig/gosunspec v0.0.0-20200429133549-3cf6a82fed9c h1:AMtX56iHlNYVxMID7fe9efuVtaxgtdjyMeolg7q87IE= github.com/andig/gosunspec v0.0.0-20200429133549-3cf6a82fed9c/go.mod h1:YkshK8WMzYn1iXAZzHUO75gIqhMSan2ctgBVtBkRIyA= -github.com/andig/ocpp-go v0.12.1-0.20201110090118-4fdb491db96c h1:2pVN47ILzwG8KkDMw+inI2VVLi+q2o4j1DEOhHIvGtA= -github.com/andig/ocpp-go v0.12.1-0.20201110090118-4fdb491db96c/go.mod h1:mN2Tv8bNqQlFMECj3IqcjZL7fOi3PmNnyVuEjICdxgU= github.com/andig/ocpp-go v0.12.1-0.20201110113243-43b1af9c1480 h1:5yKucDmOcJ6vRNeAn7PCRqYwU3/UrPY4dhwBUoNR+ug= github.com/andig/ocpp-go v0.12.1-0.20201110113243-43b1af9c1480/go.mod h1:mN2Tv8bNqQlFMECj3IqcjZL7fOi3PmNnyVuEjICdxgU= github.com/andybalholm/cascadia v1.1.0 h1:BuuO6sSfQNFRu1LppgbD25Hr2vLYW25JvxHs5zzsLTo= @@ -56,6 +55,7 @@ github.com/coreos/go-semver v0.2.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3Ee github.com/coreos/go-semver v0.3.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3EedlOD2RNk= github.com/coreos/go-systemd v0.0.0-20190321100706-95778dfbb74e/go.mod h1:F5haX7vjVVG0kc13fIWeqUViNPyEJxv/OmvnBo0Yme4= github.com/coreos/pkg v0.0.0-20180928190104-399ea9e2e55f/go.mod h1:E3G3o1h8I7cfcXa63jLwjI0eiQQMgzzUDFVpN/nH/eA= +github.com/cpuguy83/go-md2man/v2 v2.0.0 h1:EoUDS0afbrsXAZ9YQ9jdu/mZ2sXgT1/2yyNng4PGlyM= github.com/cpuguy83/go-md2man/v2 v2.0.0/go.mod h1:maD7wRr/U5Z6m/iR4s+kqSMx2CaBsrgA7czyZG/E6dU= github.com/cyberdelia/templates v0.0.0-20141128023046-ca7fffd4298c/go.mod h1:GyV+0YP4qX0UQ7r2MoYZ+AvYDp12OF5yg4q8rGnyNh4= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -91,6 +91,8 @@ github.com/go-gl/glfw v0.0.0-20190409004039-e6da0acd62b1/go.mod h1:vR7hzQXu2zJy9 github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE= github.com/go-logfmt/logfmt v0.4.0/go.mod h1:3RMwSq7FuexP4Kalkev3ejPJsZTpXXBr9+V4qmtdjCk= +github.com/go-ping/ping v0.0.0-20201022122018-3977ed72668a h1:O9xspHB2yrvKfMQ1m6OQhqe37i5yvg0dXAYMuAjugmM= +github.com/go-ping/ping v0.0.0-20201022122018-3977ed72668a/go.mod h1:35JbSyV/BYqHwwRA6Zr1uVDm1637YlNOU61wI797NPI= github.com/go-playground/locales v0.12.1 h1:2FITxuFt/xuCNP1Acdhv62OzaCiviiE4kotfhkmOqEc= github.com/go-playground/locales v0.12.1/go.mod h1:IUMDtCfWo/w/mtMfIE/IG2K+Ey3ygWanZIBtBW0W2TM= github.com/go-playground/universal-translator v0.16.0 h1:X++omBR/4cE2MNg91AoC3rmGrCjJ8eAeUP/K/EKx4DM= @@ -209,6 +211,8 @@ github.com/jarcoal/httpmock v1.0.4 h1:jp+dy/+nonJE4g4xbVtl9QdrUNbn6/3hDT5R4nDIZn github.com/jarcoal/httpmock v1.0.4/go.mod h1:ATjnClrvW/3tijVmpL/va5Z3aAyGvqU3gCT8nX0Txik= github.com/jehiah/go-strftime v0.0.0-20171201141054-1d33003b3869 h1:IPJ3dvxmJ4uczJe5YQdrYB16oTJlGSC/OyZDqUk9xX4= github.com/jehiah/go-strftime v0.0.0-20171201141054-1d33003b3869/go.mod h1:cJ6Cj7dQo+O6GJNiMx+Pa94qKj+TG8ONdKHgMNIyyag= +github.com/jeremywohl/flatten v1.0.1 h1:LrsxmB3hfwJuE+ptGOijix1PIfOoKLJ3Uee/mzbgtrs= +github.com/jeremywohl/flatten v1.0.1/go.mod h1:4AmD/VxjWcI5SRB0n6szE2A6s2fsNHDLO0nAlMHgfLQ= github.com/joeshaw/carwings v0.0.0-20191118152321-61b46581307a h1:wOlcxK/k5rPhtMyEFVayqVh0x3lWvnyEaK8Aw96zx4M= github.com/joeshaw/carwings v0.0.0-20191118152321-61b46581307a/go.mod h1:tB0OlpicmRVTL1Vksc5XRiYo+wkK2kl/GI7eGMIl6Rs= github.com/jonboulle/clockwork v0.1.0/go.mod h1:Ii8DK3G1RaLaWxj9trq07+26W01tbo22gdxWY5EU2bo= @@ -230,6 +234,8 @@ github.com/konsorten/go-windows-terminal-sequences v1.0.3 h1:CE8S1cTafDpPvMhIxNJ github.com/konsorten/go-windows-terminal-sequences v1.0.3/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/koron/go-ssdp v0.0.0-20191105050749-2e1c40ed0b5d h1:68u9r4wEvL3gYg2jvAOgROwZ3H+Y3hIDk4tbbmIjcYQ= github.com/koron/go-ssdp v0.0.0-20191105050749-2e1c40ed0b5d/go.mod h1:5Ky9EC2xfoUKUor0Hjgi2BJhCSXJfMOFlmyYrVKGQMk= +github.com/korylprince/ipnetgen v0.0.0-20160712025547-34a2abb431f4 h1:yMBxUzMGyS40obSiIkcj6d8HllfujIcwftQDPDK64Hg= +github.com/korylprince/ipnetgen v0.0.0-20160712025547-34a2abb431f4/go.mod h1:WUy1jS5kZzPfL2gWyJq3LAV5XzK71v9G/NMcT0TAYVo= github.com/kr/fs v0.1.0/go.mod h1:FFnZGqtBN9Gxj7eW1uZ42v5BccTP0vu6NEaFoC2HwRg= github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= @@ -264,6 +270,8 @@ github.com/mattn/go-isatty v0.0.10/go.mod h1:qgIWMr58cqv1PHHyhnkY9lrL7etaEgOFcME github.com/mattn/go-isatty v0.0.11/go.mod h1:PhnuNfih5lzO57/f3n+odYbM4JtupLOxQOAqxQCu2WE= github.com/mattn/go-isatty v0.0.12 h1:wuysRhFDzyxgEmMf5xjvJ2M9dZoWAXNNr5LSBS7uHXY= github.com/mattn/go-isatty v0.0.12/go.mod h1:cbi8OIDigv2wuxKPP5vlRcQ1OAZbq2CE4Kysco4FUpU= +github.com/mattn/go-runewidth v0.0.7/go.mod h1:H031xJmbD/WCDINGzjvQ9THkh0rPKHF+m2gUSrubnMI= +github.com/mattn/go-runewidth v0.0.9 h1:Lm995f3rfxdpd6TSmuVCHVb/QhupuXlYr8sCI/QdE+0= github.com/mattn/go-runewidth v0.0.9/go.mod h1:H031xJmbD/WCDINGzjvQ9THkh0rPKHF+m2gUSrubnMI= github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0= github.com/miekg/dns v1.0.14/go.mod h1:W1PPwlIAgtquWBMBEV9nkV9Cazfe8ScdGz/Lj7v3Nrg= @@ -292,12 +300,15 @@ github.com/mxschmitt/golang-combinations v1.0.0/go.mod h1:RbMhWvfCelHR6WROvT2bVf github.com/nirasan/go-oauth-pkce-code-verifier v0.0.0-20170819232839-0fbfe93532da h1:qiPWuGGr+1GQE6s9NPSK8iggR/6x/V+0snIoOPYsBgc= github.com/nirasan/go-oauth-pkce-code-verifier v0.0.0-20170819232839-0fbfe93532da/go.mod h1:DvuJJ/w1Y59rG8UTDxsMk5U+UJXJwuvUgbiJSm9yhX8= github.com/oklog/ulid v1.3.1/go.mod h1:CirwcVhetQ6Lv90oh/F+FBtV6XMibvdAFo93nm5qn4U= +github.com/olekukonko/tablewriter v0.0.4 h1:vHD/YYe1Wolo78koG299f7V/VAS08c6IpCLn+Ejf/w8= +github.com/olekukonko/tablewriter v0.0.4/go.mod h1:zq6QwlOf5SlnkVbMSr5EoBv3636FWnp+qbPhuoO21uA= github.com/onsi/ginkgo v1.6.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/ginkgo v1.8.0 h1:VkHVNpR4iVnU8XQR6DBm8BqYjN7CRzw+xKUbVVbbW9w= github.com/onsi/ginkgo v1.8.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/gomega v1.5.0 h1:izbySO9zDPmjJ8rDjLvkA2zJHIo+HkYXHnf7eN7SSyo= github.com/onsi/gomega v1.5.0/go.mod h1:ex+gbHU/CVuBBDIJjb2X0qEXbFg53c61hWP/1CpauHY= github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= +github.com/pascaldekloe/name v0.0.0-20180628100202-0fd16699aae1 h1:/I3lTljEEDNYLho3/FUB7iD/oc2cEFgVmbHzV+O0PtU= github.com/pascaldekloe/name v0.0.0-20180628100202-0fd16699aae1/go.mod h1:eD5JxqMiuNYyFNmyY9rkJ/slN8y59oEu4Ei7F8OoKWQ= github.com/paypal/gatt v0.0.0-20151011220935-4ae819d591cf/go.mod h1:+AwQL2mK3Pd3S+TUwg0tYQjid0q1txyNUJuuSmz8Kdk= github.com/pbnjay/strptime v0.0.0-20140226051138-5c05b0d668c9 h1:4lfz0keanz7/gAlvJ7lAe9zmE08HXxifBZJC0AdeGKo= @@ -326,9 +337,11 @@ github.com/prometheus/procfs v0.0.0-20190507164030-5867b95ac084/go.mod h1:TjEm7z github.com/prometheus/tsdb v0.7.1/go.mod h1:qhTCs0VvXwvX/y3TZrWD7rabWM+ijKTux40TwIPHuXU= github.com/rogpeppe/fastuuid v0.0.0-20150106093220-6724a57986af/go.mod h1:XWv6SoW27p1b0cqNHllgS5HIMJraePCO15w5zCzIWYg= github.com/rogpeppe/go-internal v1.3.0/go.mod h1:M8bDsm7K2OlrFYOpmOWEs/qY81heoFRclV5y23lUDJ4= +github.com/russross/blackfriday/v2 v2.0.1 h1:lPqVAte+HuHNfhJ/0LC98ESWRz8afy9tM/0RK8m9o+Q= github.com/russross/blackfriday/v2 v2.0.1/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/ryanuber/columnize v0.0.0-20160712163229-9b3edd62028f/go.mod h1:sm1tb6uqfes/u+d4ooFouqFdy9/2g9QGwK3SQygK0Ts= github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529/go.mod h1:DxrIzT+xaE7yg65j358z/aeFdxmN0P9QXhEzd20vsDc= +github.com/shurcooL/sanitized_anchor_name v1.0.0 h1:PdmoCO6wvbs+7yrJyMORt4/BmY5IYyJwS/kOiWx8mHo= github.com/shurcooL/sanitized_anchor_name v1.0.0/go.mod h1:1NzhyTcUVG4SuEtjjoZeVRXNmyL/1OwPU0+IJeTBvfc= github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo= github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= @@ -382,6 +395,8 @@ github.com/tebeka/strftime v0.1.3 h1:5HQXOqWKYRFfNyBMNVc9z5+QzuBtIXy03psIhtdJYto github.com/tebeka/strftime v0.1.3/go.mod h1:7wJm3dZlpr4l/oVK0t1HYIc4rMzQ2XJlOMIUJUJH6XQ= github.com/technoweenie/multipartstreamer v1.0.1 h1:XRztA5MXiR1TIRHxH2uNxXxaIkKQDeX7m2XsSOlQEnM= github.com/technoweenie/multipartstreamer v1.0.1/go.mod h1:jNVxdtShOxzAsukZwTSw6MDx5eUJoiEBsSvzDU9uzog= +github.com/thoas/go-funk v0.7.0 h1:GmirKrs6j6zJbhJIficOsz2aAI7700KsU/5YrdHRM1Y= +github.com/thoas/go-funk v0.7.0/go.mod h1:+IWnUfUmFO1+WVYQWQtIJHeRRdaIyyYglZN7xzUPe4Q= github.com/tmc/grpc-websocket-proxy v0.0.0-20190109142713-0ad062ec5ee5/go.mod h1:ncp9v5uamzpCO7NfCPTXjqaC+bZgJeR0sMTm6dMHP7U= github.com/tv42/httpunix v0.0.0-20191220191345-2ba4b9c3382c h1:u6SKchux2yDvFQnDHS3lPnIRmfVJ5Sxy3ao2SIdysLQ= github.com/tv42/httpunix v0.0.0-20191220191345-2ba4b9c3382c/go.mod h1:hzIxponao9Kjc7aWznkXaL4U4TWaDSs8zcsY4Ka08nM= @@ -390,8 +405,8 @@ github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyC github.com/valyala/fasttemplate v1.0.1/go.mod h1:UQGH1tvbgY+Nz5t2n7tXsz52dQxojPUpymEIMZ47gx8= github.com/valyala/fasttemplate v1.1.0/go.mod h1:UQGH1tvbgY+Nz5t2n7tXsz52dQxojPUpymEIMZ47gx8= github.com/volkszaehler/mbmd v0.0.0-20200717102329-c4d965bd1eac/go.mod h1:sldLyJCKVO9GQkit55U5WPNr9U4KmUn1SCdqFN9Gjb4= -github.com/volkszaehler/mbmd v0.0.0-20200831092453-b235d6a65b21 h1:KiwFeqxuXkdnv67iuR2qOiyevfPrO4Ma00Ch+J3Tr5I= -github.com/volkszaehler/mbmd v0.0.0-20200831092453-b235d6a65b21/go.mod h1:PQGaIeLLHhwwqQ56gnHJn9xSlwPoI1OeIXyGc/eLP0Y= +github.com/volkszaehler/mbmd v0.0.0-20201115202927-ff826598e117 h1:jKhfYg79as+otwlxGbZcGT8tzfi12oRddmsRQrFtSxw= +github.com/volkszaehler/mbmd v0.0.0-20201115202927-ff826598e117/go.mod h1:rq6/XBs3AGX2d7qgpTgpVrbwN56o19RBQl5FopJMIZ4= github.com/xiang90/probing v0.0.0-20190116061207-43a291ad63a2/go.mod h1:UETIi67q53MR2AWcXfiuqkDkRtnGDLqkBTpCHuJHxtU= github.com/xordataexchange/crypt v0.0.3-0.20170626215501-b2862e3d0a77/go.mod h1:aYKd//L2LvnjZzWKhF00oedf4jCCReLcmhLdhm1A27Q= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= @@ -457,8 +472,9 @@ golang.org/x/net v0.0.0-20200202094626-16171245cfb2/go.mod h1:z5CRVTTTmAJ677TzLL golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200324143707-d3edc9973b7e/go.mod h1:qpuaurCH72eLCgpAm/N6yyVIVM9cpaDIP3A8BGJEC5A= golang.org/x/net v0.0.0-20200625001655-4c5254603344/go.mod h1:/O7V0waA8r7cgGh81Ro3o1hOxt32SMVPicZroKQ2sZA= -golang.org/x/net v0.0.0-20200707034311-ab3426394381 h1:VXak5I6aEWmAXeQjA+QSZzlgNrpq9mjcfDemuexIKsU= golang.org/x/net v0.0.0-20200707034311-ab3426394381/go.mod h1:/O7V0waA8r7cgGh81Ro3o1hOxt32SMVPicZroKQ2sZA= +golang.org/x/net v0.0.0-20200904194848-62affa334b73 h1:MXfv8rhZWmFeqX3GNZRsd6vOLoaCHjYEX3qkRo3YBUA= +golang.org/x/net v0.0.0-20200904194848-62affa334b73/go.mod h1:/O7V0waA8r7cgGh81Ro3o1hOxt32SMVPicZroKQ2sZA= golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= golang.org/x/oauth2 v0.0.0-20190604053449-0f29369cfe45/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= diff --git a/hems/semp/helper.go b/hems/semp/helper.go deleted file mode 100644 index a501ee0a3..000000000 --- a/hems/semp/helper.go +++ /dev/null @@ -1,30 +0,0 @@ -package semp - -import "net" - -// LocalIPs returns a slice of local IPv4 addresses -func LocalIPs() []net.IP { - ips := make([]net.IP, 0) - - ifaces, err := net.Interfaces() - if err != nil { - panic(err) - } - - for _, i := range ifaces { - addrs, err := i.Addrs() - if err != nil { - panic(err) - } - - for _, addr := range addrs { - if ip, ok := addr.(*net.IPNet); ok { - if !ip.IP.IsLoopback() && ip.IP.To4() != nil { - ips = append(ips, ip.IP) - } - } - } - } - - return ips -} diff --git a/hems/semp/semp.go b/hems/semp/semp.go index 387363aa3..9454af85a 100644 --- a/hems/semp/semp.go +++ b/hems/semp/semp.go @@ -142,7 +142,7 @@ func (s *SEMP) callbackURI() string { } ip := "localhost" - ips := LocalIPs() + ips := util.LocalIPs() if len(ips) > 0 { ip = ips[0].String() } else { diff --git a/icon.png b/icon.png new file mode 100644 index 000000000..1a137b77e Binary files /dev/null and b/icon.png differ diff --git a/meter/sma/listener.go b/meter/sma/listener.go index 108f97bfc..3da9a21c0 100644 --- a/meter/sma/listener.go +++ b/meter/sma/listener.go @@ -19,8 +19,8 @@ const ( msgPreamble = 28 // preamble size in bytes msgCodeLength = 4 // length in bytes - // All subscriber receives all messages - All = "" + // Any subscriber receives all messages + Any = "" ) // Obis defines an Obis code as understood my the EMETER protocol @@ -203,11 +203,11 @@ func (l *Listener) send(msg Telegram) { defer l.mux.Unlock() for identifier, client := range l.clients { - if identifier == msg.Addr || identifier == msg.Serial || identifier == All { + if identifier == msg.Addr || identifier == msg.Serial || identifier == Any { select { case client <- msg: default: - l.log.TRACE.Println("listener: recv blocked") + l.log.TRACE.Println("recv: listener blocked") } break } diff --git a/util/net.go b/util/net.go index a4e4bbd27..66589d9d9 100644 --- a/util/net.go +++ b/util/net.go @@ -13,3 +13,29 @@ func DefaultPort(conn string, port int) string { return conn } + +// LocalIPs returns a slice of local IPv4 addresses +func LocalIPs() (ips []net.IPNet) { + ifaces, err := net.Interfaces() + if err != nil { + panic(err) + } + + for _, i := range ifaces { + addrs, err := i.Addrs() + if err != nil { + panic(err) + } + + for _, addr := range addrs { + // fmt.Println(addr) + if ip, ok := addr.(*net.IPNet); ok { + if !ip.IP.IsLoopback() && ip.IP.To4() != nil { + ips = append(ips, *ip) + } + } + } + } + + return ips +}