chore: prune folder structure

This commit is contained in:
andig 2023-11-01 23:39:08 +01:00
parent 384cedae00
commit 54ab598602
15 changed files with 6 additions and 7 deletions

98
cmd/detect/analyze.go Normal file
View file

@ -0,0 +1,98 @@
package detect
import "github.com/evcc-io/evcc/cmd/detect/tasks"
type Criteria map[string]interface{}
func filter(list []tasks.Result, criteria []Criteria) (match []tasks.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 []tasks.Result
Found, Unique bool
}
type Summary struct {
Charger, Grid, PV, Charge, Battery, Meter TypeSummary
}
func summarize(res []tasks.Result) TypeSummary {
return TypeSummary{
Results: res,
Found: len(res) > 0,
Unique: len(res) == 1,
}
}
const (
tid = "task.id"
smaHttp = "details.http"
)
func Consolidate(res []tasks.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: taskPhoenixEMEth},
{tid: taskPhoenixEVEth},
{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),
}
}

291
cmd/detect/definitions.go Normal file
View file

@ -0,0 +1,291 @@
package detect
import (
"github.com/evcc-io/evcc/cmd/detect/tasks"
)
var (
taskList = &TaskList{}
sunspecIDs = []int{1, 2, 3, 71, 126, 200, 201, 202, 203, 204, 240} // modbus ids
chargeStatus = []int{0x41, 0x42, 0x43} // status values A..C
)
// public task ids
const (
TaskPing = "ping"
TaskHttp = "tcp_http"
TaskModbus = "tcp_modbus"
TaskSunspec = "sunspec"
)
// private task ids
const (
taskOpenwb = "openwb"
taskSMA = "sma"
taskKEBA = "keba"
taskE3DC = "e3dc_simple"
taskSonnen = "sonnen"
taskPowerwall = "powerwall"
taskWallbe = "wallbe"
taskPhoenixEMEth = "phx-em-eth"
taskPhoenixEVEth = "phx-ev-eth"
taskEVSEWifi = "evsewifi"
taskGoE = "go-e"
taskInverter = "inverter"
taskStrings = "strings"
taskBattery = "battery"
taskMeter = "meter"
taskFroniusWeb = "fronius-web"
taskTasmota = "tasmota"
taskShelly = "shelly"
// taskTPLink = "tplink"
)
func init() {
taskList.Add(tasks.Task{
ID: TaskPing,
Type: tasks.Ping,
})
taskList.Add(tasks.Task{
ID: taskSMA,
Type: tasks.Sma,
Depends: TaskPing,
})
taskList.Add(tasks.Task{
ID: taskKEBA,
Type: tasks.Keba,
Depends: TaskPing,
})
taskList.Add(tasks.Task{
ID: TaskModbus,
Type: tasks.Tcp,
Depends: TaskPing,
Config: map[string]interface{}{
"ports": []int{502, 1502},
},
})
taskList.Add(tasks.Task{
ID: TaskSunspec,
Type: tasks.Modbus,
Depends: TaskModbus,
Config: map[string]interface{}{
"ids": sunspecIDs,
"models": []int{1},
"point": "Mn",
},
})
taskList.Add(tasks.Task{
ID: taskInverter,
Type: tasks.Modbus,
Depends: TaskSunspec,
Config: map[string]interface{}{
"ids": sunspecIDs,
"models": []int{101, 103},
"point": "W",
"invalid": []int{0xFFFF},
},
})
taskList.Add(tasks.Task{
ID: taskStrings,
Type: tasks.Modbus,
Depends: TaskSunspec,
Config: map[string]interface{}{
"ids": sunspecIDs,
"models": []int{160},
"point": "N",
"invalid": []int{0xFFFF},
},
})
taskList.Add(tasks.Task{
ID: taskBattery,
Type: tasks.Modbus,
Depends: TaskSunspec,
Config: map[string]interface{}{
"ids": sunspecIDs,
"models": []int{124},
"point": "ChaSt",
"invalid": []int{0xFFFF},
},
})
taskList.Add(tasks.Task{
ID: taskMeter,
Type: tasks.Modbus,
Depends: TaskSunspec,
Config: map[string]interface{}{
"ids": sunspecIDs,
"models": []int{201, 203, 211, 213},
"point": "W",
},
})
taskList.Add(tasks.Task{
ID: taskE3DC,
Type: tasks.Modbus,
Depends: TaskModbus,
Config: map[string]interface{}{
"ids": []int{1, 2, 3, 4, 5, 6},
"address": 40000,
"type": "holding",
"decode": "uint16",
"values": []int{0xE3DC},
},
})
taskList.Add(tasks.Task{
ID: taskWallbe,
Type: tasks.Modbus,
Depends: TaskModbus,
Config: map[string]interface{}{
"ids": []int{255},
"address": 100,
"type": "input",
"decode": "uint16",
"values": chargeStatus,
},
})
taskList.Add(tasks.Task{
ID: taskPhoenixEMEth,
Type: tasks.Modbus,
Depends: TaskModbus,
Config: map[string]interface{}{
"ids": []int{180},
"address": 100,
"type": "input",
"decode": "uint16",
"values": chargeStatus,
},
})
taskList.Add(tasks.Task{
ID: taskPhoenixEVEth,
Type: tasks.Modbus,
Depends: TaskModbus,
Config: map[string]interface{}{
"ids": []int{255},
"address": 100,
"type": "input",
"decode": "uint16",
"values": chargeStatus,
},
})
taskList.Add(tasks.Task{
ID: taskOpenwb,
Type: tasks.Mqtt,
Depends: TaskPing,
Config: map[string]interface{}{
"topic": "openWB",
},
})
taskList.Add(tasks.Task{
ID: TaskHttp,
Type: tasks.Tcp,
Depends: TaskPing,
Config: map[string]interface{}{
"ports": []int{80, 443},
},
})
taskList.Add(tasks.Task{
ID: taskGoE,
Type: tasks.Http,
Depends: TaskHttp,
Config: map[string]interface{}{
"path": "/status",
"jq": ".car",
},
})
taskList.Add(tasks.Task{
ID: taskEVSEWifi,
Type: tasks.Http,
Depends: TaskHttp,
Config: map[string]interface{}{
"path": "/getParameters",
"jq": ".type",
},
})
taskList.Add(tasks.Task{
ID: taskSonnen,
Type: tasks.Http,
Depends: TaskPing,
Config: map[string]interface{}{
"port": 8080,
"path": "/api/v1/status",
"jq": ".GridFeedIn_W",
},
})
taskList.Add(tasks.Task{
ID: taskPowerwall,
Type: tasks.Http,
Depends: TaskHttp,
Config: map[string]interface{}{
"path": "/api/meters/aggregates",
"jq": ".load",
},
})
taskList.Add(tasks.Task{
ID: taskFroniusWeb,
Type: tasks.Http,
Depends: TaskHttp,
Config: map[string]interface{}{
"path": "/solar_api/GetAPIVersion.cgi",
"jq": ".BaseURL",
},
})
taskList.Add(tasks.Task{
ID: taskTasmota,
Type: tasks.Http,
Depends: TaskHttp,
Config: map[string]interface{}{
"path": "/cm?cmnd=Module",
"jq": ".Module",
},
})
// taskList.Add(tasks.Task{
// ID: taskTPLink,
// Type: tasks.Http,
// Depends: TaskHttp,
// Config: map[string]interface{}{
// "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",
"jq": ".version",
},
})
taskList.Add(tasks.Task{
ID: taskShelly,
Type: tasks.Http,
Depends: TaskHttp,
Config: map[string]interface{}{
"path": "/shelly",
"jq": ".type",
},
})
}

125
cmd/detect/tasklist.go Normal file
View file

@ -0,0 +1,125 @@
package detect
import (
"fmt"
"sync"
"github.com/evcc-io/evcc/cmd/detect/tasks"
"github.com/evcc-io/evcc/util"
)
type TaskList struct {
tasks []tasks.Task
once sync.Once
}
func (l *TaskList) Add(task tasks.Task) {
task.TaskHandler = l.handler(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 []tasks.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) 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, id string, input tasks.ResultDetails) []tasks.Result {
l.once.Do(l.sort)
var all []tasks.Result
var inputs []tasks.ResultDetails
if id == "" {
inputs = append(inputs, input)
} else {
var task tasks.Task
for _, t := range l.tasks {
if t.ID == id {
task = t
break
}
}
inputs = task.Test(log, input)
success := len(inputs) > 0
log.DEBUG.Printf("task: %s %v -> %v", id, input, success)
for _, detail := range inputs {
all = append(all, tasks.Result{
Task: task,
ResultDetails: detail,
})
}
}
// 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
}

View file

@ -0,0 +1,5 @@
package tasks
import "time"
const timeout = 200 * time.Millisecond

137
cmd/detect/tasks/http.go Normal file
View file

@ -0,0 +1,137 @@
package tasks
import (
"fmt"
"io"
"net/http"
"strings"
"time"
"github.com/evcc-io/evcc/util"
"github.com/evcc-io/evcc/util/jq"
"github.com/itchyny/gojq"
)
const Http TaskType = "http"
func init() {
registry.Add(Http, HttpHandlerFactory)
}
type HttpResult struct {
Jq interface{}
}
func HttpHandlerFactory(conf map[string]interface{}) (TaskHandler, error) {
handler := HttpHandler{
Schema: "http",
Method: "GET",
Codes: []int{200},
Header: map[string]string{
"Content-type": "application/json",
},
Timeout: 3 * timeout,
}
err := util.DecodeOther(conf, &handler)
switch handler.Schema {
case "http":
handler.Port = 80
case "https":
handler.Port = 443
}
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
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, 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
}
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
}
}
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
if val == nil || err != nil {
return nil
}
}
if err == nil {
out := in.Clone()
out.Port = port
return []ResultDetails{out}
}
return nil
}

87
cmd/detect/tasks/keba.go Normal file
View file

@ -0,0 +1,87 @@
package tasks
import (
"sync"
"time"
"github.com/evcc-io/evcc/charger/keba"
"github.com/evcc-io/evcc/util"
)
const Keba TaskType = "keba"
func init() {
registry.Add(Keba, KEBAHandlerFactory)
}
type KebaResult struct {
Addr, Serial string
}
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
}
func (h *KEBAHandler) Test(log *util.Logger, in ResultDetails) []ResultDetails {
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(in.IP, resC)
sender, err := keba.NewSender(log, in.IP)
if err != nil {
log.ERROR.Println("keba:", err)
return nil
}
timer := time.NewTimer(h.Timeout)
for {
go func() {
_ = sender.Send("report 1")
}()
select {
case t := <-resC:
log.INFO.Println(t)
if t.Report == nil {
continue
}
out := in.Clone()
out.KebaResult = &KebaResult{
Addr: t.Addr,
Serial: t.Report.Serial,
}
return []ResultDetails{out}
case <-timer.C:
return nil
}
}
}

219
cmd/detect/tasks/modbus.go Normal file
View file

@ -0,0 +1,219 @@
package tasks
import (
"encoding/binary"
"errors"
"fmt"
"net"
"strconv"
"time"
"github.com/evcc-io/evcc/util"
"github.com/evcc-io/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"
)
const Modbus TaskType = "modbus"
func init() {
registry.Add(Modbus, ModbusHandlerFactory)
}
type ModbusResult struct {
SlaveID uint8
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": net.JoinHostPort(res.ResultDetails.IP, strconv.Itoa(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 gridx.FuncCodeReadHoldingRegisters:
bytes, err = conn.ReadHoldingRegisters(h.op.OpCode, h.op.ReadLen)
case gridx.FuncCodeReadInputRegisters:
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.DEBUG.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())
case "count":
val = int(res.Count())
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, 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 := net.JoinHostPort(in.IP, strconv.Itoa(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 {
out := in.Clone()
out.ModbusResult = &mr
res = append(res, out)
}
}
return res
}

75
cmd/detect/tasks/mqtt.go Normal file
View file

@ -0,0 +1,75 @@
package tasks
import (
"errors"
"net"
"strconv"
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
"github.com/evcc-io/evcc/util"
)
const Mqtt TaskType = "mqtt"
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, in ResultDetails) []ResultDetails {
addr := net.JoinHostPort(in.IP, strconv.Itoa(h.Port))
opt := mqtt.NewClientOptions()
opt.AddBroker(addr)
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)
for {
select {
case <-recv:
out := in.Clone()
out.Topic = h.Topic
return []ResultDetails{out}
case <-timer.C:
return nil
}
}
}
return nil
}

67
cmd/detect/tasks/ping.go Normal file
View file

@ -0,0 +1,67 @@
package tasks
import (
"runtime"
"time"
"github.com/evcc-io/evcc/util"
ping "github.com/prometheus-community/pro-bing"
)
const Ping TaskType = "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, in ResultDetails) []ResultDetails {
pinger, err := ping.NewPinger(in.IP)
if err != nil {
panic(err)
}
if runtime.GOOS == "windows" {
pinger.Size = 548 // https://github.com/go-ping/ping/issues/168
pinger.SetPrivileged(true)
}
pinger.Count = h.Count
pinger.Timeout = h.Timeout
if err = pinger.Run(); err != nil {
log.FATAL.Println("ping:", err)
if runtime.GOOS != "windows" {
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 []ResultDetails{in}
}

View file

@ -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
}

92
cmd/detect/tasks/sma.go Normal file
View file

@ -0,0 +1,92 @@
package tasks
import (
"crypto/tls"
"fmt"
"net/http"
"strconv"
"sync"
"time"
"github.com/evcc-io/evcc/util"
"gitlab.com/bboehmke/sunny"
)
const Sma TaskType = "sma"
func init() {
registry.Add(Sma, SMAHandlerFactory)
}
type SmaResult struct {
Serial string
Http bool
}
func SMAHandlerFactory(conf map[string]interface{}) (TaskHandler, error) {
handler := SMAHandler{
Timeout: 5 * time.Second,
Password: "0000",
}
err := util.DecodeOther(conf, &handler)
return &handler, err
}
type SMAHandler struct {
mux sync.Mutex
handled bool
Timeout time.Duration
Password string
}
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, in ResultDetails) (res []ResultDetails) {
h.mux.Lock()
if h.handled {
h.mux.Unlock()
return nil
}
connection, err := sunny.NewConnection("")
if err != nil {
log.ERROR.Println("sma:", err)
return nil
}
devices := connection.SimpleDiscoverDevices(h.Password)
h.handled = true
h.mux.Unlock()
for _, device := range devices {
res = append(res, ResultDetails{
IP: device.Address().IP.String(),
SmaResult: &SmaResult{
Serial: strconv.FormatInt(int64(device.SerialNumber()), 10),
Http: h.httpAvailable(device.Address().IP.String()),
},
})
}
return res
}

54
cmd/detect/tasks/tcp.go Normal file
View file

@ -0,0 +1,54 @@
package tasks
import (
"errors"
"net"
"strconv"
"time"
"github.com/evcc-io/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 := net.JoinHostPort(in.IP, strconv.Itoa(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
}

43
cmd/detect/tasks/types.go Normal file
View file

@ -0,0 +1,43 @@
package tasks
import (
"github.com/evcc-io/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"`
SmaResult *SmaResult `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
}

96
cmd/detect/work.go Normal file
View file

@ -0,0 +1,96 @@
package detect
import (
"sort"
"strings"
"sync"
"github.com/evcc-io/evcc/cmd/detect/tasks"
"github.com/evcc-io/evcc/util"
"github.com/fatih/structs"
"github.com/jeremywohl/flatten"
)
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)
go func() {
workunit(log, tasks, hits)
wg.Done()
}()
}
return &wg
}
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) []tasks.Result {
ip := make(chan string)
hits := make(chan []tasks.Result)
done := make(chan struct{})
// log.INFO.Println(
// "\n" +
// strings.Join(
// lo.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),
// "",
// ),
// )
wg := workers(log, num, ip, hits)
var res []tasks.Result
go func() {
for hits := range hits {
res = append(res, hits...)
}
done <- struct{}{}
}()
for _, host := range hosts {
ip <- host
}
close(ip)
wg.Wait()
close(hits)
<-done
return postProcess(res)
}
func postProcess(res []tasks.Result) []tasks.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].ResultDetails.IP == res[j].ResultDetails.IP {
return res[i].Type < res[j].Type
}
return res[i].ResultDetails.IP < res[j].ResultDetails.IP
})
return res
}