package cardata import ( "context" "fmt" "maps" "slices" "strings" "sync" "time" "github.com/evcc-io/evcc/api" "github.com/evcc-io/evcc/util" "github.com/spf13/cast" "golang.org/x/oauth2" ) const StreamingURL = "tls://customer.streaming-cardata.bmwgroup.com:9000" // Provider implements the vehicle api type Provider struct { mu sync.Mutex log *util.Logger api *API ts oauth2.TokenSource vin string container string rest map[string]TelematicData streaming map[string]StreamingData updated time.Time cache time.Duration } // NewProvider creates a vehicle api provider func NewProvider(ctx context.Context, log *util.Logger, api *API, ts oauth2.TokenSource, clientID, vin string, cache time.Duration) *Provider { v := &Provider{ log: log, api: api, ts: ts, vin: vin, cache: cache, rest: make(map[string]TelematicData), streaming: make(map[string]StreamingData), } mqtt := NewMqttConnector(context.Background(), log, clientID, ts) recvC := mqtt.Subscribe(vin) go func() { <-ctx.Done() mqtt.Unsubscribe(vin, recvC) }() go func() { for msg := range recvC { v.mu.Lock() maps.Copy(v.streaming, msg.Data) v.updated = time.Now() v.mu.Unlock() } }() return v } func (v *Provider) findOrCreateContainer() (string, error) { containers, err := v.api.GetContainers() if err != nil { return "", err } // obsolete containers keep streaming, resulting in duplicate messages defer v.deleteObsoleteContainers(containers) if i := slices.IndexFunc(containers, func(c Container) bool { return c.Name == containerName && c.Purpose == requiredVersion }); i >= 0 { return containers[i].ContainerId, nil } res, err := v.api.CreateContainer(CreateContainer{ Name: containerName, Purpose: requiredVersion, TechnicalDescriptors: requiredKeys, }) return res.ContainerId, err } func (v *Provider) deleteObsoleteContainers(containers []Container) { for _, c := range containers { if c.Name != containerName || c.Purpose == requiredVersion { continue } v.log.DEBUG.Printf("deleting obsolete container %s (%s)", c.ContainerId, c.Purpose) if err := v.api.DeleteContainer(c.ContainerId); err != nil { v.log.WARN.Printf("delete container %s: %v", c.ContainerId, err) } } } func (v *Provider) setupContainer() error { container, err := v.findOrCreateContainer() if err != nil { return fmt.Errorf("get container: %v", err) } v.container = container return nil } func (v *Provider) updateContainerData() error { res, err := v.api.GetTelematics(v.vin, v.container) if err != nil { return fmt.Errorf("get telematics: %v", err) } v.rest = res.TelematicData v.streaming = make(map[string]StreamingData) // reset streaming return nil } func (v *Provider) any(key string) (any, error) { v.mu.Lock() defer v.mu.Unlock() _, tokenErr := v.ts.Token() if tokenErr == nil && (v.updated.IsZero() || time.Since(v.updated) > v.cache) { if v.container == "" { if err := v.setupContainer(); err != nil { v.log.WARN.Println(err) } } if v.container != "" { if err := v.updateContainerData(); err != nil { v.log.WARN.Println(err) } } v.updated = time.Now() } if a, ok := v.streaming[key]; ok { return a.Value, nil } if el, ok := v.rest[key]; ok { return el.Value, nil } return nil, api.ErrNotAvailable } func isNilOrEmpty(val any) bool { return val == nil || val == "" } func (v *Provider) String(key string) (string, error) { res, err := v.any(key) if err != nil || isNilOrEmpty(res) { return "", api.ErrNotAvailable } return cast.ToStringE(res) } func (v *Provider) Int(key string) (int64, error) { res, err := v.any(key) if err != nil || isNilOrEmpty(res) { return 0, api.ErrNotAvailable } return cast.ToInt64E(res) } func (v *Provider) Float(key string) (float64, error) { res, err := v.any(key) if err != nil || isNilOrEmpty(res) { return 0, api.ErrNotAvailable } return cast.ToFloat64E(res) } var _ api.Battery = (*Provider)(nil) // Soc implements the api.Vehicle interface func (v *Provider) Soc() (float64, error) { if res, err := v.Float("vehicle.drivetrain.batteryManagement.header"); err == nil { return res, nil } return v.Float("vehicle.powertrain.electric.battery.stateOfCharge.displayed") } var _ api.ChargeState = (*Provider)(nil) // Status implements the api.ChargeState interface func (v *Provider) Status() (api.ChargeStatus, error) { // evaluate status first, since it's usually available through // mqtt, while hvStatus might only be available through rest // (https://github.com/evcc-io/evcc/pull/26235) cs, err := v.String("vehicle.drivetrain.electricEngine.charging.status") if err != nil { cs, _ = v.String("vehicle.drivetrain.electricEngine.charging.hvStatus") } if slices.Contains([]string{ "INITIALIZATION", // vehicle.drivetrain.electricEngine.charging.status "CHARGINGPAUSED", // vehicle.drivetrain.electricEngine.charging.status "CHARGINGENDED", // vehicle.drivetrain.electricEngine.charging.status "WAITING_FOR_CHARGING", // vehicle.drivetrain.electricEngine.charging.hvStatus "FINISHED_FULLY_CHARGED", // vehicle.drivetrain.electricEngine.charging.hvStatus "FINISHED_NOT_FULL", // vehicle.drivetrain.electricEngine.charging.hvStatus }, cs) { return api.StatusB, nil } if slices.Contains([]string{ "CHARGINGACTIVE", // vehicle.drivetrain.electricEngine.charging.status "CHARGING", // vehicle.drivetrain.electricEngine.charging.hvStatus }, cs) { return api.StatusC, nil } port, err := v.String("vehicle.body.chargingPort.status") if err != nil { port, err = v.String("vehicle.body.chargingPort.combinedStatus") if err != nil { return api.StatusNone, err } } status := api.StatusA // disconnected if port == "CONNECTED" { status = api.StatusB } return status, err } var _ api.VehicleFinishTimer = (*Provider)(nil) // FinishTime implements the api.VehicleFinishTimer interface func (v *Provider) FinishTime() (time.Time, error) { res, err := v.Int("vehicle.drivetrain.electricEngine.charging.timeRemaining") if err != nil { return time.Time{}, err } return time.Now().Add(time.Duration(res) * time.Minute), nil } var _ api.VehicleRange = (*Provider)(nil) // Range implements the api.VehicleRange interface func (v *Provider) Range() (int64, error) { if res, err := v.Int("vehicle.drivetrain.electricEngine.kombiRemainingElectricRange"); err == nil { return res, nil } return v.Int("vehicle.drivetrain.lastRemainingRange") } var _ api.VehicleOdometer = (*Provider)(nil) // Odometer implements the api.VehicleOdometer interface func (v *Provider) Odometer() (float64, error) { return v.Float("vehicle.vehicle.travelledDistance") } var _ api.SocLimiter = (*Provider)(nil) // GetLimitSoc implements the api.SocLimiter interface func (v *Provider) GetLimitSoc() (int64, error) { return v.Int("vehicle.powertrain.electric.battery.stateOfCharge.target") } var _ api.VehicleClimater = (*Provider)(nil) // Climater implements the api.VehicleClimater interface func (v *Provider) Climater() (bool, error) { activeStates := []string{"HEATING", "COOLING", "VENTILATION", "DEFROST"} res, err := v.String("vehicle.cabin.hvac.preconditioning.status.comfortState") if err == nil { return slices.Contains(activeStates, strings.TrimPrefix(strings.ToUpper(res), "COMFORT_")), nil } if res, err = v.String("vehicle.vehicle.preConditioning.activity"); err == nil { return slices.Contains(activeStates, strings.ToUpper(res)), nil } return false, err } var _ api.VehiclePosition = (*Provider)(nil) // Position implements the api.VehiclePosition interface func (v *Provider) Position() (float64, float64, error) { lat, err := v.Float("vehicle.cabin.infotainment.navigation.currentLocation.latitude") if err != nil { return 0, 0, err } lon, err := v.Float("vehicle.cabin.infotainment.navigation.currentLocation.longitude") if err != nil { return 0, 0, err } return lat, lon, nil }