diff --git a/internal/vehicle/cloud.go b/internal/vehicle/cloud.go index d3c6ef860..84aee6612 100644 --- a/internal/vehicle/cloud.go +++ b/internal/vehicle/cloud.go @@ -3,6 +3,7 @@ package vehicle import ( "context" "errors" + "strings" "time" "github.com/andig/evcc/api" @@ -101,7 +102,12 @@ func (v *Cloud) chargeState() (float64, error) { defer cancel() res, err := v.client.SoC(ctx, req) - if errors.As(err, &cloud.ErrVehicleNotAvailable) && v.prepareVehicle() == nil { + + if err != nil && strings.Contains(err.Error(), api.ErrMustRetry.Error()) { + return 0, api.ErrMustRetry + } + + if err != nil && strings.Contains(err.Error(), cloud.ErrVehicleNotAvailable.Error()) && v.prepareVehicle() == nil { req.VehicleId = v.vehicleID res, err = v.client.SoC(ctx, req) } diff --git a/soc/client/main.go b/soc/client/main.go index b3f0ae509..11eb14580 100644 --- a/soc/client/main.go +++ b/soc/client/main.go @@ -1,11 +1,13 @@ package main import ( + "errors" "fmt" "log" "math" "os" "strings" + "time" "github.com/andig/evcc/api" "github.com/andig/evcc/internal/vehicle" @@ -17,7 +19,7 @@ func usage() { soc Usage: - soc vehicle [--log level] [--param value [...]] + soc brand [--log level] [--param value [...]] `) } @@ -70,11 +72,25 @@ func main() { } case "soc": - soc, err := v.SoC() - if err != nil { - log.Fatal(err) + start := time.Now() + for { + if time.Since(start) > time.Minute { + log.Fatal(api.ErrTimeout) + } + + soc, err := v.SoC() + if err != nil { + if errors.As(err, &api.ErrMustRetry) { + time.Sleep(5 * time.Second) + continue + } + + log.Fatal(err) + } + + fmt.Println(int(math.Round(soc))) + break } - fmt.Println(int(math.Round(soc))) default: log.Fatal("invalid action:", action) diff --git a/soc/server/server/init.go b/soc/server/server/init.go index 4a6f7ef43..f64e848aa 100644 --- a/soc/server/server/init.go +++ b/soc/server/server/init.go @@ -7,7 +7,6 @@ import ( "log" "net" - "github.com/andig/evcc/api" "github.com/andig/evcc/soc/cert/ca" "github.com/andig/evcc/soc/cert/server" "github.com/andig/evcc/soc/proto/pb" @@ -26,6 +25,8 @@ func init() { if tlsConfig, err = loadTLSCredentials(); err != nil { log.Fatalf("cannot load TLS credentials: %v", err) } + + registerMetrics() } func loadTLSCredentials() (*tls.Config, error) { @@ -65,7 +66,7 @@ func Run() { grpcServer := grpc.NewServer(serverOptions...) pb.RegisterVehicleServer(grpcServer, &VehicleServer{ - vehicles: make(map[string]map[int64]api.Vehicle), + registry: make(map[string][]*VehicleContainer), }) pb.RegisterAuthServer(grpcServer, &AuthServer{}) diff --git a/soc/server/server/prom.go b/soc/server/server/prom.go new file mode 100644 index 000000000..67fdd3e4c --- /dev/null +++ b/soc/server/server/prom.go @@ -0,0 +1,35 @@ +package server + +import ( + "log" + + "github.com/prometheus/client_golang/prometheus" +) + +var ( + promActiveVehicles *prometheus.GaugeVec +) + +func registerMetrics() { + promActiveVehicles = prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Namespace: "soc", + Name: "active_vehicles", + }, []string{"token", "brand"}) + + if err := prometheus.Register(promActiveVehicles); err != nil { + log.Fatal(err) + } +} + +func updateActiveVehiclesMetric(token, typ string, delta int) { + g, err := promActiveVehicles.GetMetricWith(prometheus.Labels{ + "token": token, + "brand": typ, + }) + if err != nil { + log.Println("get metrics:", err) + return + } + + g.Add(float64(delta)) +} diff --git a/soc/server/server/vehicle.go b/soc/server/server/vehicle.go index bef122bb6..88a4e3f77 100644 --- a/soc/server/server/vehicle.go +++ b/soc/server/server/vehicle.go @@ -1,8 +1,12 @@ package server import ( + "bytes" "context" + "crypto/sha256" + "fmt" "log" + "sort" "sync/atomic" "github.com/andig/evcc/api" @@ -13,8 +17,14 @@ import ( var vehicleID int64 +type VehicleContainer struct { + id int64 + hash []byte + vehicle api.Vehicle +} + type VehicleServer struct { - vehicles map[string]map[int64]api.Vehicle + registry map[string][]*VehicleContainer pb.UnimplementedVehicleServer } @@ -25,18 +35,19 @@ type vehicler interface { func (s *VehicleServer) vehicle(r vehicler) (api.Vehicle, error) { token := r.GetToken() - vehicles, ok := s.vehicles[token] + vehicles, ok := s.registry[token] if !ok { return nil, cloud.ErrVehicleNotAvailable } id := r.GetVehicleId() - v, ok := vehicles[id] - if !ok { - return nil, cloud.ErrVehicleNotAvailable + for _, c := range vehicles { + if c.id == id { + return c.vehicle, nil + } } - return v, nil + return nil, cloud.ErrVehicleNotAvailable } func stringMapToInterface(in map[string]string) map[string]interface{} { @@ -49,6 +60,51 @@ func stringMapToInterface(in map[string]string) map[string]interface{} { return res } +func (s *VehicleServer) addVehicleToRegistry(token, typ string, config map[string]string, v api.Vehicle) int64 { + id := atomic.AddInt64(&vehicleID, 1) + + // sort config keys + var keys []string + for k := range config { + keys = append(keys, k) + } + sort.Strings(keys) + + // hash config + h := sha256.New() + _, _ = h.Write([]byte(typ)) + for _, k := range keys { + _, _ = h.Write([]byte(k)) + _, _ = h.Write([]byte(config[k])) + } + hash := h.Sum(nil) + + // find vehicle by hash and update it + for _, c := range s.registry[token] { + if bytes.Equal(c.hash, hash) { + c.vehicle = v + c.id = id + return id + } + } + + // register new vehicle + c := VehicleContainer{ + id: id, + hash: hash, + vehicle: v, + } + s.registry[token] = append(s.registry[token], &c) + + h.Reset() + _, _ = h.Write([]byte(token)) + thash := fmt.Sprintf("%x", h.Sum(nil)) + + updateActiveVehiclesMetric(thash, typ, 1) + + return id +} + func (s *VehicleServer) New(ctx context.Context, r *pb.NewRequest) (*pb.NewReply, error) { authorized, token, claims, err := isAuthorized(r) if err != nil { @@ -69,11 +125,7 @@ func (s *VehicleServer) New(ctx context.Context, r *pb.NewRequest) (*pb.NewReply return nil, err } - id := atomic.AddInt64(&vehicleID, 1) - if s.vehicles[token] == nil { - s.vehicles[token] = make(map[int64]api.Vehicle) - } - s.vehicles[token][id] = v + id := s.addVehicleToRegistry(token, typ, config, v) res := pb.NewReply{ VehicleId: id,