From d6624ce05c731683c046d322331d2da75720dd4f Mon Sep 17 00:00:00 2001 From: andig Date: Tue, 11 May 2021 12:42:52 +0200 Subject: [PATCH] Improve detection by building a dependency tree (#986) --- cmd/detect.go | 17 ++-- detect/analyze.go | 10 +- detect/definitions.go | 158 +++++++++++++++--------------- detect/registry.go | 30 ------ detect/tasklist.go | 103 ++++++++++--------- detect/tasks/const.go | 5 + detect/{ => tasks}/http.go | 46 ++++++--- detect/{ => tasks}/keba.go | 29 +++--- detect/{ => tasks}/modbus.go | 33 +++++-- detect/{ => tasks}/mqtt.go | 23 +++-- detect/{ => tasks}/ping.go | 12 ++- detect/tasks/registry.go | 32 ++++++ detect/{ => tasks}/sma.go | 32 +++--- detect/tasks/tcp.go | 54 ++++++++++ detect/tasks/types.go | 43 ++++++++ detect/tcp.go | 48 --------- detect/work.go | 52 +++++----- go.mod | 3 +- go.sum | 8 +- internal/charger/keba/listener.go | 5 +- 20 files changed, 425 insertions(+), 318 deletions(-) delete mode 100644 detect/registry.go create mode 100644 detect/tasks/const.go rename detect/{ => tasks}/http.go (68%) rename detect/{ => tasks}/keba.go (75%) rename detect/{ => tasks}/modbus.go (85%) rename detect/{ => tasks}/mqtt.go (76%) rename detect/{ => tasks}/ping.go (80%) create mode 100644 detect/tasks/registry.go rename detect/{ => tasks}/sma.go (72%) create mode 100644 detect/tasks/tcp.go create mode 100644 detect/tasks/types.go delete mode 100644 detect/tcp.go diff --git a/cmd/detect.go b/cmd/detect.go index 16d60e615..e4d5312f1 100644 --- a/cmd/detect.go +++ b/cmd/detect.go @@ -1,12 +1,14 @@ package cmd import ( + "encoding/json" "fmt" "net" "os" "strings" "github.com/andig/evcc/detect" + "github.com/andig/evcc/detect/tasks" "github.com/andig/evcc/util" "github.com/korylprince/ipnetgen" "github.com/olekukonko/tablewriter" @@ -66,7 +68,7 @@ func ParseHostIPNet(arg string) (res []string) { return IPsFromSubnet(arg) } -func display(res []detect.Result) { +func display(res []tasks.Result) { table := tablewriter.NewWriter(os.Stdout) table.SetHeader([]string{"IP", "Hostname", "Task", "Details"}) table.SetAutoMergeCells(true) @@ -74,23 +76,20 @@ func display(res []detect.Result) { for _, hit := range res { switch hit.ID { - case detect.TaskPing, detect.TaskTCP80, detect.TaskTCP502: + case detect.TaskPing, detect.TaskHttp, detect.TaskModbus: continue default: host := "" - hosts, err := net.LookupAddr(hit.Host) + hosts, err := net.LookupAddr(hit.ResultDetails.IP) if err == nil && len(hosts) > 0 { host = strings.TrimSuffix(hosts[0], ".") } - details := "" - if hit.Details != nil { - details = fmt.Sprintf("%+v", hit.Details) - } + b, _ := json.Marshal(hit.ResultDetails) - // fmt.Printf("%-16s %-20s %-16s %s\n", hit.Host, host, hit.ID, details) - table.Append([]string{hit.Host, host, hit.ID, details}) + // fmt.Printf("%-16s %-20s %-16s %s\n", hit.ResultDetails.IP, host, hit.ID, details) + table.Append([]string{hit.ResultDetails.IP, host, hit.ID, string(b)}) } } diff --git a/detect/analyze.go b/detect/analyze.go index 4d4f16337..390d0bf22 100644 --- a/detect/analyze.go +++ b/detect/analyze.go @@ -1,8 +1,10 @@ package detect +import "github.com/andig/evcc/detect/tasks" + type Criteria map[string]interface{} -func filter(list []Result, criteria []Criteria) (match []Result) { +func filter(list []tasks.Result, criteria []Criteria) (match []tasks.Result) { for _, res := range list { for _, criterium := range criteria { ok := true @@ -24,7 +26,7 @@ func filter(list []Result, criteria []Criteria) (match []Result) { } type TypeSummary struct { - Results []Result + Results []tasks.Result Found, Unique bool } @@ -32,7 +34,7 @@ type Summary struct { Charger, Grid, PV, Charge, Battery, Meter TypeSummary } -func summarize(res []Result) TypeSummary { +func summarize(res []tasks.Result) TypeSummary { return TypeSummary{ Results: res, Found: len(res) > 0, @@ -45,7 +47,7 @@ const ( smaHttp = "details.http" ) -func Consolidate(res []Result) Summary { +func Consolidate(res []tasks.Result) Summary { grid := filter(res, []Criteria{ {tid: taskOpenwb}, {tid: taskSMA, smaHttp: false}, diff --git a/detect/definitions.go b/detect/definitions.go index 8bfdf1c7d..b0aa21487 100644 --- a/detect/definitions.go +++ b/detect/definitions.go @@ -1,6 +1,10 @@ package detect -import "time" +import ( + "time" + + "github.com/andig/evcc/detect/tasks" +) var ( taskList = &TaskList{} @@ -9,21 +13,19 @@ var ( 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" + TaskHttp = "tcp_http" + TaskModbus = "tcp_modbus" TaskSunspec = "sunspec" ) // private task ids const ( taskOpenwb = "openwb" - taskSMA = "sma" - taskKEBA = "KEBA" + taskSMA = "shm" + taskKEBA = "keba" taskE3DC = "e3dc_simple" taskSonnen = "sonnen" taskPowerwall = "powerwall" @@ -37,38 +39,38 @@ const ( taskMeter = "meter" taskFronius = "fronius" taskTasmota = "tasmota" - taskTPLink = "tplink" + // taskTPLink = "tplink" ) func init() { - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskSMA, - Type: "sma", + Type: tasks.Sma, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskKEBA, - Type: "keba", + Type: tasks.Keba, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: TaskPing, - Type: "ping", + Type: tasks.Ping, }) - taskList.Add(Task{ - ID: TaskTCP502, - Type: "tcp", + taskList.Add(tasks.Task{ + ID: TaskModbus, + Type: tasks.Tcp, Depends: TaskPing, Config: map[string]interface{}{ - "port": 502, + "ports": []int{502, 1502}, }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: TaskSunspec, - Type: "modbus", - Depends: TaskTCP502, + Type: tasks.Modbus, + Depends: TaskModbus, Config: map[string]interface{}{ "ids": sunspecIDs, "models": []int{1}, @@ -76,9 +78,9 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskInverter, - Type: "modbus", + Type: tasks.Modbus, Depends: TaskSunspec, Config: map[string]interface{}{ "ids": sunspecIDs, @@ -88,9 +90,9 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskBattery, - Type: "modbus", + Type: tasks.Modbus, Depends: TaskSunspec, Config: map[string]interface{}{ "ids": sunspecIDs, @@ -100,9 +102,9 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskMeter, - Type: "modbus", + Type: tasks.Modbus, Depends: TaskSunspec, Config: map[string]interface{}{ "ids": sunspecIDs, @@ -111,10 +113,10 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskE3DC, - Type: "modbus", - Depends: TaskTCP502, + Type: tasks.Modbus, + Depends: TaskModbus, Config: map[string]interface{}{ "ids": []int{1, 2, 3, 4, 5, 6}, "address": 40000, @@ -124,10 +126,10 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskWallbe, - Type: "modbus", - Depends: TaskTCP502, + Type: tasks.Modbus, + Depends: TaskModbus, Config: map[string]interface{}{ "ids": []int{255}, "address": 100, @@ -137,10 +139,10 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskPhoenixEMEth, - Type: "modbus", - Depends: TaskTCP502, + Type: tasks.Modbus, + Depends: TaskModbus, Config: map[string]interface{}{ "ids": []int{180}, "address": 100, @@ -150,10 +152,10 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskPhoenixEVEth, - Type: "modbus", - Depends: TaskTCP502, + Type: tasks.Modbus, + Depends: TaskModbus, Config: map[string]interface{}{ "ids": []int{255}, "address": 100, @@ -163,47 +165,47 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskOpenwb, - Type: "mqtt", + Type: tasks.Mqtt, Depends: TaskPing, Config: map[string]interface{}{ "topic": "openWB", }, }) - taskList.Add(Task{ - ID: TaskTCP80, - Type: "tcp", + taskList.Add(tasks.Task{ + ID: TaskHttp, + Type: tasks.Tcp, Depends: TaskPing, Config: map[string]interface{}{ - "port": 80, + "ports": []int{80, 443}, }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskGoE, - Type: "http", - Depends: TaskTCP80, + Type: tasks.Http, + Depends: TaskHttp, Config: map[string]interface{}{ "path": "/status", "jq": ".car", }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskEVSEWifi, - Type: "http", - Depends: TaskTCP80, + Type: tasks.Http, + Depends: TaskHttp, Config: map[string]interface{}{ "path": "/getParameters", "jq": ".type", }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskSonnen, - Type: "http", + Type: tasks.Http, Depends: TaskPing, Config: map[string]interface{}{ "port": 8080, @@ -212,52 +214,54 @@ func init() { }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskPowerwall, - Type: "http", - Depends: TaskTCP80, + Type: tasks.Http, + Depends: TaskHttp, Config: map[string]interface{}{ "path": "/api/meters/aggregates", "jq": ".load", }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskFronius, - Type: "http", - Depends: TaskTCP80, + Type: tasks.Http, + Depends: TaskHttp, Config: map[string]interface{}{ "path": "/solar_api/GetAPIVersion.cgi", "jq": ".BaseURL", }, }) - taskList.Add(Task{ + taskList.Add(tasks.Task{ ID: taskTasmota, - Type: "http", - Depends: TaskTCP80, + Type: tasks.Http, + Depends: TaskHttp, Config: map[string]interface{}{ "path": "//cm?cmnd=Module", "jq": ".Module", }, }) - taskList.Add(Task{ - ID: taskTPLink, - Type: "tcp", - Depends: TaskPing, - Config: map[string]interface{}{ - "port": 9999, // TP-Link Smart Home Protocol standard port - }, - }) - - // taskList.Add(Task{ - // ID: "volkszähler", - // Type: "http", - // Depends: TaskTCP80, + // taskList.Add(tasks.Task{ + // ID: taskTPLink, + // Type: tasks.Http, + // Depends: TaskHttp, // Config: map[string]interface{}{ - // "path": "/middleware.php/entity.json", - // "timeout": 500 * time.Millisecond, + // "ResponseHeader": map[string]string{ + // "Server": "TP-LINK Smart Plug", + // }, // }, // }) + + taskList.Add(tasks.Task{ + ID: "volkszähler", + Type: tasks.Http, + Depends: TaskHttp, + Config: map[string]interface{}{ + "path": "/middleware.php/entity.json", + "timeout": 500 * time.Millisecond, + }, + }) } diff --git a/detect/registry.go b/detect/registry.go deleted file mode 100644 index 7bfc65dfe..000000000 --- a/detect/registry.go +++ /dev/null @@ -1,30 +0,0 @@ -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/tasklist.go b/detect/tasklist.go index 44233ec8b..27a3fa0ad 100644 --- a/detect/tasklist.go +++ b/detect/tasklist.go @@ -4,23 +4,17 @@ import ( "fmt" "sync" + "github.com/andig/evcc/detect/tasks" "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 + tasks []tasks.Task + once sync.Once } -func (l *TaskList) Add(task Task) { +func (l *TaskList) Add(task tasks.Task) { + task.TaskHandler = l.handler(task) l.tasks = append(l.tasks, task) } @@ -42,7 +36,7 @@ func (l *TaskList) delete(i int) { } func (l *TaskList) sort() { - var res []Task + var res []tasks.Task for len(l.tasks) > 0 { last := len(l.tasks) @@ -72,54 +66,59 @@ func (l *TaskList) sort() { 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) handler(task tasks.Task) tasks.TaskHandler { + factory, err := tasks.Get(task.Type) + if err != nil { + panic("invalid task type " + task.Type) } + + // fmt.Println(task) + handler, err := factory(task.Config) + if err != nil { + panic("invalid config: " + err.Error()) + } + + return handler } -func (l *TaskList) Test(log *util.Logger, ip string) (res []Result) { - l.once.Do(func() { - l.sort() - l.createHandlers() - }) +func (l *TaskList) Test(log *util.Logger, id string, input tasks.ResultDetails) []tasks.Result { + l.once.Do(l.sort) - failed := make([]string, 0) + var all []tasks.Result + var inputs []tasks.ResultDetails -HANDLERS: - for id, handler := range l.handlers { - task := l.tasks[id] + if id == "" { + inputs = append(inputs, input) + } else { + log.DEBUG.Printf("ip: %s task: %s (%v)", input.IP, id, input) - 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, - }) + var task tasks.Task + for _, t := range l.tasks { + if t.ID == id { + task = t + break } - } else { - // log.INFO.Printf("ip: %s task: %s nok", ip, task.ID) - failed = append(failed, task.ID) + } + + inputs = task.Test(log, input) + for _, detail := range inputs { + all = append(all, tasks.Result{ + Task: task, + ResultDetails: detail, + }) } } - return res + // run dependent tasks + for _, task := range l.tasks { + if task.Depends == id { + // fmt.Println("task:", task) + for _, input := range inputs { + // fmt.Println("input:", input) + all = append(all, l.Test(log, task.ID, input)...) + } + } + } + + return all } diff --git a/detect/tasks/const.go b/detect/tasks/const.go new file mode 100644 index 000000000..bc5a37c96 --- /dev/null +++ b/detect/tasks/const.go @@ -0,0 +1,5 @@ +package tasks + +import "time" + +const timeout = 200 * time.Millisecond diff --git a/detect/http.go b/detect/tasks/http.go similarity index 68% rename from detect/http.go rename to detect/tasks/http.go index c2dcfe131..43e0c16d2 100644 --- a/detect/http.go +++ b/detect/tasks/http.go @@ -1,4 +1,4 @@ -package detect +package tasks import ( "fmt" @@ -12,8 +12,10 @@ import ( "github.com/itchyny/gojq" ) +const Http TaskType = "http" + func init() { - registry.Add("http", HttpHandlerFactory) + registry.Add(Http, HttpHandlerFactory) } type HttpResult struct { @@ -23,7 +25,6 @@ type HttpResult struct { func HttpHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { handler := HttpHandler{ Schema: "http", - Port: 80, Method: "GET", Codes: []int{200}, Header: map[string]string{ @@ -34,9 +35,11 @@ func HttpHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { 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) + switch handler.Schema { + case "http": + handler.Port = 80 + case "https": + handler.Port = 443 } if handler.Jq != "" { @@ -54,16 +57,25 @@ func HttpHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { type HttpHandler struct { query *gojq.Query Port int - optionalPort string Schema, Method, Path string Codes []int Header map[string]string + ResponseHeader 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, "/")) +func (h *HttpHandler) Test(log *util.Logger, in ResultDetails) []ResultDetails { + port := in.Port + if port == 0 { + port = h.Port + } + + if port == 0 { + panic("http: invalid port") + } + + uri := fmt.Sprintf("%s://%s:%d/%s", h.Schema, in.IP, port, strings.TrimLeft(h.Path, "/")) req, err := http.NewRequest(strings.ToUpper(h.Method), uri, nil) if err != nil { return nil @@ -94,13 +106,19 @@ func (h *HttpHandler) Test(log *util.Logger, ip string) []interface{} { } } - body, err := io.ReadAll(resp.Body) - if err != nil { - return nil + for k, v := range h.ResponseHeader { + if resp.Header.Get(k) != v { + return nil + } } var res HttpResult if h.query != nil { + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil + } + val, err := jq.Query(h.query, body) res.Jq = val @@ -110,7 +128,9 @@ func (h *HttpHandler) Test(log *util.Logger, ip string) []interface{} { } if err == nil { - return []interface{}{res} + out := in.Clone() + out.Port = port + return []ResultDetails{out} } return nil diff --git a/detect/keba.go b/detect/tasks/keba.go similarity index 75% rename from detect/keba.go rename to detect/tasks/keba.go index d0ed78c6e..6ebf2ff18 100644 --- a/detect/keba.go +++ b/detect/tasks/keba.go @@ -1,4 +1,4 @@ -package detect +package tasks import ( "sync" @@ -8,12 +8,14 @@ import ( "github.com/andig/evcc/util" ) -type KebaResult struct { - Addr, Serial string -} +const Keba TaskType = "keba" func init() { - registry.Add("keba", KEBAHandlerFactory) + registry.Add(Keba, KEBAHandlerFactory) +} + +type KebaResult struct { + Addr, Serial string } func KEBAHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { @@ -32,7 +34,7 @@ type KEBAHandler struct { Timeout time.Duration } -func (h *KEBAHandler) Test(log *util.Logger, ip string) []interface{} { +func (h *KEBAHandler) Test(log *util.Logger, in ResultDetails) []ResultDetails { h.mux.Lock() if h.listener == nil { @@ -48,16 +50,16 @@ func (h *KEBAHandler) Test(log *util.Logger, ip string) []interface{} { h.mux.Unlock() resC := make(chan keba.UDPMsg) - h.listener.Subscribe(ip, resC) + h.listener.Subscribe(in.IP, resC) - sender, err := keba.NewSender(log, ip) + sender, err := keba.NewSender(log, in.IP) if err != nil { log.ERROR.Println("keba:", err) return nil } timer := time.NewTimer(h.Timeout) -WAIT: + for { go func() { _ = sender.Send("report 1") @@ -70,17 +72,16 @@ WAIT: continue } - r := KebaResult{ + out := in.Clone() + out.KebaResult = &KebaResult{ Addr: t.Addr, Serial: t.Report.Serial, } - return []interface{}{r} + return []ResultDetails{out} case <-timer.C: - break WAIT + return nil } } - - return nil } diff --git a/detect/modbus.go b/detect/tasks/modbus.go similarity index 85% rename from detect/modbus.go rename to detect/tasks/modbus.go index ef3fcfb37..f859f8a34 100644 --- a/detect/modbus.go +++ b/detect/tasks/modbus.go @@ -1,4 +1,4 @@ -package detect +package tasks import ( "encoding/binary" @@ -14,21 +14,23 @@ import ( "github.com/volkszaehler/mbmd/meters/sunspec" ) +const Modbus TaskType = "modbus" + func init() { - registry.Add("modbus", ModbusHandlerFactory) + registry.Add(Modbus, ModbusHandlerFactory) } type ModbusResult struct { SlaveID uint8 - Model int - Point string - Value interface{} + Model int `json:",omitempty"` + Point string `json:",omitempty"` + Value interface{} `json:",omitempty"` } 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), + "uri": fmt.Sprintf("%s:%d", res.ResultDetails.IP, port), "model": "sunspec", "id": r.SlaveID, } @@ -38,7 +40,7 @@ func (r *ModbusResult) Configuration(handler TaskHandler, res Result) map[string func ModbusHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { handler := ModbusHandler{ - Port: 502, + // Port: 502, IDs: []uint8{1}, Models: []int{1}, Point: "Md", // Model @@ -165,8 +167,17 @@ func (h *ModbusHandler) testSunSpec(log *util.Logger, conn meters.Connection, de return false } -func (h *ModbusHandler) Test(log *util.Logger, ip string) (res []interface{}) { - addr := fmt.Sprintf("%s:%d", ip, h.Port) +func (h *ModbusHandler) Test(log *util.Logger, in ResultDetails) (res []ResultDetails) { + port := in.Port + if port == 0 { + port = h.Port + } + if port == 0 { + fmt.Println("modbus", in) + panic("modbus: invalid port") + } + + addr := fmt.Sprintf("%s:%d", in.IP, port) conn := meters.NewTCP(addr) dev := sunspec.NewDevice("sunspec") @@ -194,7 +205,9 @@ func (h *ModbusHandler) Test(log *util.Logger, ip string) (res []interface{}) { } if ok { - res = append(res, mr) + out := in.Clone() + out.ModbusResult = &mr + res = append(res, out) } } diff --git a/detect/mqtt.go b/detect/tasks/mqtt.go similarity index 76% rename from detect/mqtt.go rename to detect/tasks/mqtt.go index 21e35794c..ac45d697e 100644 --- a/detect/mqtt.go +++ b/detect/tasks/mqtt.go @@ -1,4 +1,4 @@ -package detect +package tasks import ( "errors" @@ -9,8 +9,10 @@ import ( mqtt "github.com/eclipse/paho.mqtt.golang" ) +const Mqtt TaskType = "mqtt" + func init() { - registry.Add("mqtt", MqttHandlerFactory) + registry.Add(Mqtt, MqttHandlerFactory) } func MqttHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { @@ -33,8 +35,8 @@ type MqttHandler struct { Timeout time.Duration } -func (h *MqttHandler) Test(log *util.Logger, ip string) []interface{} { - broker := fmt.Sprintf("%s:%d", ip, h.Port) +func (h *MqttHandler) Test(log *util.Logger, in ResultDetails) []ResultDetails { + broker := fmt.Sprintf("%s:%d", in.IP, h.Port) opt := mqtt.NewClientOptions() opt.AddBroker(broker) @@ -55,21 +57,18 @@ func (h *MqttHandler) Test(log *util.Logger, ip string) []interface{} { }) timer := time.NewTimer(timeout) - WAIT: + for { select { case <-recv: - break WAIT + out := in.Clone() + out.Topic = h.Topic + return []ResultDetails{out} case <-timer.C: - ok = false - break WAIT + return nil } } } - if ok { - return []interface{}{nil} - } - return nil } diff --git a/detect/ping.go b/detect/tasks/ping.go similarity index 80% rename from detect/ping.go rename to detect/tasks/ping.go index 032ad9313..a6c2a424a 100644 --- a/detect/ping.go +++ b/detect/tasks/ping.go @@ -1,4 +1,4 @@ -package detect +package tasks import ( "runtime" @@ -8,8 +8,10 @@ import ( "github.com/go-ping/ping" ) +const Ping TaskType = "ping" + func init() { - registry.Add("ping", PingHandlerFactory) + registry.Add(Ping, PingHandlerFactory) } func PingHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { @@ -28,8 +30,8 @@ type PingHandler struct { Timeout time.Duration } -func (h *PingHandler) Test(log *util.Logger, ip string) (res []interface{}) { - pinger, err := ping.NewPinger(ip) +func (h *PingHandler) Test(log *util.Logger, in ResultDetails) []ResultDetails { + pinger, err := ping.NewPinger(in.IP) if err != nil { panic(err) } @@ -60,5 +62,5 @@ func (h *PingHandler) Test(log *util.Logger, ip string) (res []interface{}) { return nil } - return []interface{}{nil} + return []ResultDetails{in} } diff --git a/detect/tasks/registry.go b/detect/tasks/registry.go new file mode 100644 index 000000000..01f24f830 --- /dev/null +++ b/detect/tasks/registry.go @@ -0,0 +1,32 @@ +package tasks + +import ( + "fmt" +) + +type TaskHandlerRegistry map[TaskType]func(map[string]interface{}) (TaskHandler, error) + +var registry TaskHandlerRegistry = make(map[TaskType]func(map[string]interface{}) (TaskHandler, error)) + +func (r TaskHandlerRegistry) Add(name TaskType, 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 +// } + +func Get(name TaskType) (func(map[string]interface{}) (TaskHandler, error), error) { + factory, exists := registry[name] + if !exists { + return nil, fmt.Errorf("charger type not registered: %s", name) + } + return factory, nil +} diff --git a/detect/sma.go b/detect/tasks/sma.go similarity index 72% rename from detect/sma.go rename to detect/tasks/sma.go index 53e0e5611..721bbd8a7 100644 --- a/detect/sma.go +++ b/detect/tasks/sma.go @@ -1,4 +1,4 @@ -package detect +package tasks import ( "crypto/tls" @@ -11,13 +11,15 @@ import ( "github.com/andig/evcc/util" ) -type SmaResult struct { - Addr, Serial string - Http bool -} +const Sma TaskType = "shm" func init() { - registry.Add("sma", SMAHandlerFactory) + registry.Add(Sma, SMAHandlerFactory) +} + +type ShmResult struct { + Serial string + Http bool } func SMAHandlerFactory(conf map[string]interface{}) (TaskHandler, error) { @@ -55,7 +57,7 @@ func (h *SMAHandler) httpAvailable(ip string) bool { return true } -func (h *SMAHandler) Test(log *util.Logger, ip string) (res []interface{}) { +func (h *SMAHandler) Test(log *util.Logger, in ResultDetails) (res []ResultDetails) { h.mux.Lock() if h.listener != nil { @@ -65,7 +67,7 @@ func (h *SMAHandler) Test(log *util.Logger, ip string) (res []interface{}) { var err error if h.listener, err = sma.New(log); err != nil { - log.ERROR.Println("sma:", err) + log.ERROR.Println("shm:", err) return nil } h.mux.Unlock() @@ -80,18 +82,20 @@ WAIT: case t := <-resC: // eliminate duplicates for _, r := range res { - if r.(SmaResult).Serial == t.Serial { + if r.ShmResult != nil && r.ShmResult.Serial == t.Serial { continue WAIT } } - r := SmaResult{ - Addr: t.Addr, - Serial: t.Serial, - Http: h.httpAvailable(t.Addr), + out := ResultDetails{ + IP: t.Addr, + ShmResult: &ShmResult{ + Serial: t.Serial, + Http: h.httpAvailable(t.Addr), + }, } - res = append(res, r) + res = append(res, out) case <-timer.C: break WAIT diff --git a/detect/tasks/tcp.go b/detect/tasks/tcp.go new file mode 100644 index 000000000..200904acc --- /dev/null +++ b/detect/tasks/tcp.go @@ -0,0 +1,54 @@ +package tasks + +import ( + "errors" + "fmt" + "net" + "time" + + "github.com/andig/evcc/util" +) + +const Tcp TaskType = "tcp" + +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 && len(handler.Ports) == 0 { + err = errors.New("missing port") + } + + handler.dialer = net.Dialer{Timeout: handler.Timeout} + return &handler, err +} + +type TcpHandler struct { + Ports []int + Timeout time.Duration + dialer net.Dialer +} + +func (h *TcpHandler) Test(log *util.Logger, in ResultDetails) (res []ResultDetails) { + for _, port := range h.Ports { + addr := fmt.Sprintf("%s:%d", in.IP, port) + conn, err := h.dialer.Dial("tcp", addr) + if err == nil { + defer conn.Close() + } + + if err == nil { + out := in.Clone() + out.Port = port + res = append(res, out) + } + } + + return res +} diff --git a/detect/tasks/types.go b/detect/tasks/types.go new file mode 100644 index 000000000..4df2ae22a --- /dev/null +++ b/detect/tasks/types.go @@ -0,0 +1,43 @@ +package tasks + +import ( + "github.com/andig/evcc/util" + "github.com/jinzhu/copier" +) + +type ResultDetails struct { + IP string + Port int `json:",omitempty"` + Topic string `json:",omitempty"` + ModbusResult *ModbusResult `json:",omitempty"` + KebaResult *KebaResult `json:",omitempty"` + ShmResult *ShmResult `json:",omitempty"` +} + +func (d *ResultDetails) Clone() ResultDetails { + var c ResultDetails + if err := copier.Copy(&c, *d); err != nil { + panic(err) + } + return c +} + +type Result struct { + Task + ResultDetails + Attributes map[string]interface{} // TODO remove, only used for post-processing +} + +type TaskType string + +type Task struct { + ID string + Type TaskType + Depends string + Config map[string]interface{} + TaskHandler +} + +type TaskHandler interface { + Test(log *util.Logger, in ResultDetails) []ResultDetails +} diff --git a/detect/tcp.go b/detect/tcp.go deleted file mode 100644 index e1d53913d..000000000 --- a/detect/tcp.go +++ /dev/null @@ -1,48 +0,0 @@ -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 index 73531cebd..b353bfa28 100644 --- a/detect/work.go +++ b/detect/work.go @@ -5,19 +5,13 @@ import ( "strings" "sync" + "github.com/andig/evcc/detect/tasks" "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 { +func workers(log *util.Logger, num int, tasks <-chan string, hits chan<- []tasks.Result) *sync.WaitGroup { var wg sync.WaitGroup for i := 0; i < num; i++ { wg.Add(1) @@ -30,21 +24,31 @@ func workers(log *util.Logger, num int, tasks <-chan string, hits chan<- []Resul return &wg } -func workunit(log *util.Logger, tasks <-chan string, hits chan<- []Result) { - for ip := range tasks { - res := taskList.Test(log, ip) +func workunit(log *util.Logger, ips <-chan string, hits chan<- []tasks.Result) { + for ip := range ips { + res := taskList.Test(log, "", tasks.ResultDetails{IP: ip}) hits <- res } } -func Work(log *util.Logger, num int, hosts []string) []Result { - tasks := make(chan string) - hits := make(chan []Result) +func Work(log *util.Logger, num int, hosts []string) []tasks.Result { + ip := make(chan string) + hits := make(chan []tasks.Result) done := make(chan struct{}) - wg := workers(log, num, tasks, hits) + // log.INFO.Println( + // "\n" + + // strings.Join( + // funk.Map(taskList.tasks, func(t tasks.Task) string { + // return fmt.Sprintf("task: %s\ttype: %s\tdepends: %s\n", t.ID, t.Type, t.Depends) + // }).([]string), + // "", + // ), + // ) - var res []Result + wg := workers(log, num, ip, hits) + + var res []tasks.Result go func() { for hits := range hits { res = append(res, hits...) @@ -53,10 +57,10 @@ func Work(log *util.Logger, num int, hosts []string) []Result { }() for _, host := range hosts { - tasks <- host + ip <- host } - close(tasks) + close(ip) wg.Wait() close(hits) @@ -65,11 +69,11 @@ func Work(log *util.Logger, num int, hosts []string) []Result { return postProcess(res) } -func postProcess(res []Result) []Result { +func postProcess(res []tasks.Result) []tasks.Result { for idx, hit := range res { - if sma, ok := hit.Details.(SmaResult); ok { - hit.Host = sma.Addr - } + // 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) @@ -82,10 +86,10 @@ func postProcess(res []Result) []Result { // sort by host sort.Slice(res, func(i, j int) bool { - if res[i].Host == res[j].Host { + if res[i].ResultDetails.IP == res[j].ResultDetails.IP { return res[i].Type < res[j].Type } - return res[i].Host < res[j].Host + return res[i].ResultDetails.IP < res[j].ResultDetails.IP }) return res diff --git a/go.mod b/go.mod index 06eebc7c6..a2b849cc0 100644 --- a/go.mod +++ b/go.mod @@ -19,7 +19,7 @@ require ( github.com/eclipse/paho.mqtt.golang v1.3.4 github.com/fatih/structs v1.1.0 github.com/felixge/httpsnoop v1.0.2 // indirect - github.com/go-ping/ping v0.0.0-20210506233800-ff8be3320020 + github.com/go-ping/ping v0.0.0-20210407214646-e4e642a95741 github.com/go-telegram-bot-api/telegram-bot-api v4.6.4+incompatible github.com/godbus/dbus/v5 v5.0.4 github.com/gokrazy/updater v0.0.0-20210130175436-d85b92498a28 @@ -39,6 +39,7 @@ require ( github.com/itchyny/gojq v0.12.3 github.com/itchyny/timefmt-go v0.1.3 // indirect github.com/jeremywohl/flatten v1.0.1 + github.com/jinzhu/copier v0.3.0 github.com/joeshaw/carwings v0.0.0-20210208214325-dacfdd3d7acc github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51 github.com/klauspost/compress v1.12.2 // indirect diff --git a/go.sum b/go.sum index f9816bbcb..d66b8878e 100644 --- a/go.sum +++ b/go.sum @@ -197,8 +197,8 @@ github.com/go-logfmt/logfmt v0.4.0/go.mod h1:3RMwSq7FuexP4Kalkev3ejPJsZTpXXBr9+V github.com/go-logfmt/logfmt v0.5.0/go.mod h1:wCYkCAKZfumFQihp8CzCvQ3paCTfi41vtzG1KdI/P7A= github.com/go-openapi/jsonpointer v0.19.5/go.mod h1:Pl9vOtqEWErmShwVjC8pYs9cog34VGT37dQOVbmoatg= github.com/go-openapi/swag v0.19.5/go.mod h1:POnQmlKehdgb5mhVOsnJFsivZCEZ/vjK9gh66Z9tfKk= -github.com/go-ping/ping v0.0.0-20210506233800-ff8be3320020 h1:mdi6AbCEoKCA1xKCmp7UtRB5fvGFlP92PvlhxgdvXEw= -github.com/go-ping/ping v0.0.0-20210506233800-ff8be3320020/go.mod h1:KmHOjTUmJh/l04ukqPoBWPEZr9jwN05h5NXQl5C+DyY= +github.com/go-ping/ping v0.0.0-20210407214646-e4e642a95741 h1:b0sLP++Tsle+s57tqg5sUk1/OQsC6yMCciVeqNzOcwU= +github.com/go-ping/ping v0.0.0-20210407214646-e4e642a95741/go.mod h1:35JbSyV/BYqHwwRA6Zr1uVDm1637YlNOU61wI797NPI= github.com/go-playground/assert/v2 v2.0.1/go.mod h1:VDjEfimB/XKnb+ZQfWdccd7VUvScMdVu0Titje2rxJ4= github.com/go-playground/locales v0.12.1/go.mod h1:IUMDtCfWo/w/mtMfIE/IG2K+Ey3ygWanZIBtBW0W2TM= github.com/go-playground/locales v0.13.0 h1:HyWk6mgj5qFqCT5fjGBuRArbVDfE4hi8+e8ceBS/t7Q= @@ -395,6 +395,8 @@ github.com/jarcoal/httpmock v1.0.4/go.mod h1:ATjnClrvW/3tijVmpL/va5Z3aAyGvqU3gCT 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/jinzhu/copier v0.3.0 h1:P5zN9OYSxmtzZmwgcVmt5Iu8egfP53BGMPAFgEksKPI= +github.com/jinzhu/copier v0.3.0/go.mod h1:24xnZezI2Yqac9J61UC6/dG/k76ttpq0DdJI3QmUvro= github.com/jmespath/go-jmespath v0.0.0-20180206201540-c2b33e8439af/go.mod h1:Nht3zPeWKUH0NzdCt2Blrr5ys8VGpn0CEB0cQHVjt7k= github.com/joeshaw/carwings v0.0.0-20191118152321-61b46581307a/go.mod h1:tB0OlpicmRVTL1Vksc5XRiYo+wkK2kl/GI7eGMIl6Rs= github.com/joeshaw/carwings v0.0.0-20210208214325-dacfdd3d7acc h1:lvs05O5riVMhA70SOK53fthEZ5lMDHHy51ZKDldr0aE= @@ -857,8 +859,6 @@ golang.org/x/sync v0.0.0-20200317015054-43a5402ce75a/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20200625203802-6e8e738ad208/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20210220032951-036812b2e83c h1:5KslGYwFpkhGh+Q16bwMP3cOontH8FOep7tGV86Y7SQ= -golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sys v0.0.0-20180823144017-11551d06cbcc/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= diff --git a/internal/charger/keba/listener.go b/internal/charger/keba/listener.go index 9d0928165..e48f87087 100644 --- a/internal/charger/keba/listener.go +++ b/internal/charger/keba/listener.go @@ -96,7 +96,10 @@ func (l *Listener) listen() { if body != OK { var report Report if err := json.Unmarshal([]byte(body), &report); err != nil { - l.log.WARN.Printf("recv: invalid message: %v", err) + // ignore error during detection when sending report request to localhost + if body != "report 1" { + l.log.WARN.Printf("recv: invalid message: %v", err) + } continue }