SoC: expose metrics for monitoring and fix duplicate vehicles kept in memory (#1024)
This commit is contained in:
parent
cf6d435b6e
commit
31ae1e2991
5 changed files with 129 additions and 19 deletions
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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{})
|
||||
|
||||
|
|
|
|||
35
soc/server/server/prom.go
Normal file
35
soc/server/server/prom.go
Normal file
|
|
@ -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))
|
||||
}
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue