import json import threading import uuid from datetime import datetime from json import JSONEncoder from hashlib import md5 from os import environ from time import sleep from oauth2_client.credentials_manager import ServiceInformation import paho.mqtt.client as mqtt from requests.exceptions import RequestException from urllib3.exceptions import InvalidHeader import psa_connectedcar as psac from libs.car import Cars, Car from libs.charging import Charging from libs.oauth import OpenIdCredentialManager, Oauth2PSACCApiConfig, OauthAPIClient from ecomix import Ecomix from libs.psa_constants import DELAYED_CHARGE, IMMEDIATE_CHARGE, PSA_CORRELATION_DATE_FORMAT, PSA_DATE_FORMAT, \ realm_info, MQTT_BRANDCODE, INPROGRESS, DEFAULT_PRECONDITIONING_PROGRAM from otp.otp import load_otp, save_otp, ConfigException, Otp from psa_connectedcar.rest import ApiException from mylogger import logger from libs.utils import rate_limit, parse_hour, RateLimitException from web.abrp import Abrp from web.db import Database AUTHORIZE_SERVICE = "https://api.mpsa.com/api/connectedcar/v2/oauth/authorize" REMOTE_URL = "https://api.groupe-psa.com/connectedcar/v4/virtualkey/remoteaccess/token?client_id=" SCOPE = ['openid profile'] MQTT_SERVER = "mwa.mpsa.com" MQTT_REQ_TOPIC = "psa/RemoteServices/from/cid/" MQTT_RESP_TOPIC = "psa/RemoteServices/to/cid/" MQTT_EVENT_TOPIC = "psa/RemoteServices/events/MPHRTServices/" MQTT_TOKEN_TTL = 890 CARS_FILE = "cars.json" DEFAULT_CONFIG_FILENAME = "config.json" def gen_correlation_id(date): date_str = date.strftime(PSA_CORRELATION_DATE_FORMAT)[:-3] uuid_str = str(uuid.uuid4()).replace("-", "") correlation_id = uuid_str + date_str return correlation_id class MyPSACC: def connect(self, user, password): self.manager.init_with_user_credentials_realm(user, password, self.realm) # pylint: disable=too-many-arguments def __init__(self, refresh_token, client_id, client_secret, remote_refresh_token, customer_id, realm, country_code, proxies=None, weather_api=None, abrp=None, co2_signal_api=None): self.realm = realm self.service_information = ServiceInformation(AUTHORIZE_SERVICE, realm_info[self.realm]['oauth_url'], client_id, client_secret, SCOPE, False) self.client_id = client_id self.manager = OpenIdCredentialManager(self.service_information) self.api_config = Oauth2PSACCApiConfig() self.api_config.set_refresh_callback(self.refresh_token) self.manager.refresh_token = refresh_token self.remote_refresh_token = remote_refresh_token self.remote_access_token = None self.vehicles_list = Cars.load_cars(CARS_FILE) self.customer_id = customer_id self._config_hash = None self.api_config.verify_ssl = False self.api_config.api_key['client_id'] = self.client_id self.api_config.api_key['x-introspect-realm'] = self.realm self.headers = { "x-introspect-realm": realm, "accept": "application/hal+json", "User-Agent": "okhttp/4.8.0", } self.remote_token_last_update = None self._record_enabled = False self.otp = None self.weather_api = weather_api self.country_code = country_code self.mqtt_client = None self.precond_programs = {} self.info_callback = [] self.info_refresh_rate = 120 if abrp is None: self.abrp = Abrp() else: self.abrp: Abrp = Abrp(**abrp) self.set_proxies(proxies) self.config_file = DEFAULT_CONFIG_FILENAME Ecomix.co2_signal_key = co2_signal_api self.refresh_thread = None def get_app_name(self): return realm_info[self.realm]['app_name'] @rate_limit(6, 1800) def refresh_token(self): # pylint: disable=protected-access self.manager._refresh_token() self.save_config() return True def api(self) -> psac.VehiclesApi: self.api_config.access_token = self.manager.access_token api_instance = psac.VehiclesApi(OauthAPIClient(self.api_config)) return api_instance def set_proxies(self, proxies): if proxies is None: self._proxies = dict(http='', https='') self.api_config.proxy = None else: self._proxies = proxies self.api_config.proxy = proxies['http'] self.abrp.proxies = proxies self.manager.proxies = self._proxies Otp.set_proxies(proxies) def get_vehicle_info(self, vin, cache=False): res = None car = self.vehicles_list.get_car_by_vin(vin) if cache and car.status is not None: res = car.status else: for _ in range(0, 2): try: res = self.api().get_vehicle_status(car.vehicle_id, extension=["odometer"]) if res is not None: car.status = res if self._record_enabled: self.record_info(car) return res except (ApiException, InvalidHeader) as ex: logger.error("get_vehicle_info: ApiException: %s", ex, exc_info_debug=True) car.status = res return res def __refresh_vehicle_info(self): if self.info_refresh_rate is not None: while True: try: logger.debug("refresh_vehicle_info") for car in self.vehicles_list: self.get_vehicle_info(car.vin) for callback in self.info_callback: callback() except: # pylint: disable=bare-except logger.exception("refresh_vehicle_info: ") sleep(self.info_refresh_rate) def start_refresh_thread(self): if self.refresh_thread is None: self.refresh_thread = threading.Thread(target=self.__refresh_vehicle_info) self.refresh_thread.setDaemon(True) self.refresh_thread.start() # monitor doesn't seem to work def new_monitor(self, vin, body): res = self.manager.post("https://api.groupe-psa.com/connectedcar/v4/user/vehicles/" + self.vehicles_list.get_car_by_vin(vin).id + "/status?client_id=" + self.client_id, headers=self.headers, data=body) data = res.json() return data def get_vehicles(self): try: res = self.api().get_vehicles_by_device() for vehicle in res.embedded.vehicles: self.vehicles_list.add(Car(vehicle.vin, vehicle.id, vehicle.brand, vehicle.label)) self.vehicles_list.save_cars() except (ApiException, InvalidHeader): logger.exception("get_vehicles:") return self.vehicles_list def load_otp(self, force_new=False): otp_session = load_otp() if otp_session is None or force_new: logger.error("Please redo otp config") return False self.otp = otp_session return True def get_sms_otp_code(self): res = self.manager.post( "https://api.groupe-psa.com/applications/cvs/v4/mobile/smsCode?client_id=" + self.client_id, headers={ "Connection": "Keep-Alive", "User-Agent": "okhttp/4.8.0", "x-introspect-realm": self.realm }) return res # 6 otp by day @rate_limit(6, 3600 * 24) def get_otp_code(self): try: otp_code = self.otp.get_otp_code() except ConfigException: logger.exception("get_otp_code:") self.load_otp(force_new=True) otp_code = self.otp.get_otp_code() save_otp(self.otp) return otp_code def get_remote_access_token(self, password): try: res = self.manager.post(REMOTE_URL + self.client_id, json={"grant_type": "password", "password": password}, headers=self.headers) data = res.json() self.remote_access_token = data["access_token"] self.remote_refresh_token = data["refresh_token"] return res except RequestException as e: logger.error("Can't refresh remote token %s", e) sleep(60) return None def _refresh_remote_token(self, force=False): bad_remote_token = self.remote_refresh_token is None res = None if not force and not bad_remote_token and self.remote_token_last_update: last_update: datetime = self.remote_token_last_update if (datetime.now() - last_update).total_seconds() < MQTT_TOKEN_TTL: return res try: self.refresh_token() if bad_remote_token: logger.error("remote_refresh_token isn't defined") else: res = self.manager.post(REMOTE_URL + self.client_id, json={"grant_type": "refresh_token", "refresh_token": self.remote_refresh_token}, headers=self.headers) data = res.json() logger.debug("refresh_remote_token: %s", data) if "access_token" in data: self.remote_access_token = data["access_token"] self.remote_refresh_token = data["refresh_token"] bad_remote_token = False else: logger.error("can't refresh_remote_token: %s\n Create a new one", data) bad_remote_token = True if bad_remote_token: otp_code = self.get_otp_code() res = self.get_remote_access_token(otp_code) self.remote_token_last_update = datetime.now() self.mqtt_client.username_pw_set("IMA_OAUTH_ACCESS_TOKEN", self.remote_access_token) self.save_config() return res except (RequestException, RateLimitException) as e: logger.exception("Can't refresh remote token %s", e) sleep(60) return None def __get_mqtt_customer_id(self): brand_code = self.customer_id[:2] return MQTT_BRANDCODE[brand_code] + self.customer_id[2:] # pylint: disable=unused-argument def __on_mqtt_connect(self, client, userdata, result_code, _): logger.info("Connected with result code %s", result_code) topics = [MQTT_RESP_TOPIC + self.__get_mqtt_customer_id() + "/#"] for car in self.vehicles_list: topics.append(MQTT_EVENT_TOPIC + car.vin) for topic in topics: client.subscribe(topic) logger.info("subscribe to %s", topic) # pylint: disable=unused-argument def _on_mqtt_disconnect(self, client, userdata, result_code): logger.warning("Disconnected with result code %d", result_code) if result_code == 1: self._refresh_remote_token(force=True) else: logger.warning(mqtt.error_string(result_code)) # pylint: disable=unused-argument def __on_mqtt_message(self, client, userdata, msg): try: logger.info("mqtt msg received: %s %s", msg.topic, msg.payload) data = json.loads(msg.payload) charge_info = None if msg.topic.startswith(MQTT_RESP_TOPIC): if "return_code" not in data: logger.debug("mqtt msg hasn't return code") elif data["return_code"] == "400": self._refresh_remote_token(force=True) logger.error("retry last request, token was expired") elif data["return_code"] != "0": logger.error('%s : %s', data["return_code"], data.get("reason", "?")) elif msg.topic.startswith(MQTT_EVENT_TOPIC): charge_info = data["charging_state"] self.precond_programs[data["vin"]] = data["precond_state"]["programs"] if charge_info is not None and charge_info['remaining_time'] != 0: try: car = self.vehicles_list.get_car_by_vin(vin=msg.topic.split("/")[-1]) if car and car.status.get_energy('Electric').charging.status != INPROGRESS: # fix a psa server bug where charge beginning without status api being properly updated logger.warning("charge begin but API isn't updated") sleep(60) self.wakeup(data["vin"]) except (IndexError, AttributeError, RateLimitException): logger.exception("on_mqtt_message:") except KeyError: logger.exception("on_mqtt_message:") def start_mqtt(self): self.mqtt_client = mqtt.Client(clean_session=True, protocol=mqtt.MQTTv311) if self.load_otp(): if environ.get("MQTT_LOG", "0") == "1": self.mqtt_client.enable_logger(logger=logger) if self._refresh_remote_token(): self.mqtt_client.tls_set_context() self.mqtt_client.on_connect = self.__on_mqtt_connect self.mqtt_client.on_message = self.__on_mqtt_message self.mqtt_client.on_disconnect = self._on_mqtt_disconnect self.mqtt_client.connect(MQTT_SERVER, 8885, 60) self.mqtt_client.loop_start() self.__keep_mqtt() return self.mqtt_client.is_connected() def __keep_mqtt(self): # avoid token expiration timeout = 3600 * 24 # 1 day if len(self.vehicles_list) > 0: try: self.wakeup(self.vehicles_list[0].vin) except RateLimitException: logger.exception("__keep_mqtt") t = threading.Timer(timeout, self.__keep_mqtt) t.setDaemon(True) t.start() def mqtt_request(self, vin, req_parameters): self.refresh_token() date = datetime.utcnow() date_str = date.strftime(PSA_DATE_FORMAT) data = {"access_token": self.remote_access_token, "customer_id": self.__get_mqtt_customer_id(), "correlation_id": gen_correlation_id(date), "req_date": date_str, "vin": vin, "req_parameters": req_parameters} return json.dumps(data) def __get_charge_hour(self, vin): data = self.get_vehicle_info(vin) hour_str = data.get_energy('Electric').charging.next_delayed_time try: return parse_hour(hour_str)[:2] except IndexError: logger.exception("Can't get charge hour: %s", hour_str) return None def get_charge_status(self, vin): data = self.get_vehicle_info(vin) status = data.get_energy('Electric').charging.status return status def __veh_charge_request(self, vin, hour, minute, charge_type): msg = self.mqtt_request(vin, {"program": {"hour": hour, "minute": minute}, "type": charge_type}) logger.info("veh_charge_request: %s", msg) self.mqtt_client.publish(MQTT_REQ_TOPIC + self.__get_mqtt_customer_id() + "/VehCharge", msg) return msg def change_charge_hour(self, vin, hour, miinute): self.__veh_charge_request(vin, hour, miinute, DELAYED_CHARGE) return True def charge_now(self, vin, now): if now: charge_type = IMMEDIATE_CHARGE else: charge_type = DELAYED_CHARGE hour, minute = self.__get_charge_hour(vin) res = self.__veh_charge_request(vin, hour, minute, charge_type) logger.info("charge_now: %s", res) return True def horn(self, vin, count): msg = self.mqtt_request(vin, {"nb_horn": count, "action": "activate"}) logger.info(msg) self.mqtt_client.publish(MQTT_REQ_TOPIC + self.__get_mqtt_customer_id() + "/Horn", msg) def lights(self, vin, duration: int): msg = self.mqtt_request(vin, {"action": "activate", "duration": duration}) logger.info(msg) self.mqtt_client.publish(MQTT_REQ_TOPIC + self.__get_mqtt_customer_id() + "/Lights", msg) @rate_limit(6, 60 * 20) def wakeup(self, vin): logger.info("ask wakeup to %s", vin) msg = self.mqtt_request(vin, {"action": "state"}) logger.info(msg) self.mqtt_client.publish(MQTT_REQ_TOPIC + self.__get_mqtt_customer_id() + "/VehCharge/state", msg) return True def lock_door(self, vin, lock: bool): if lock: value = "lock" else: value = "unlock" msg = self.mqtt_request(vin, {"action": value}) logger.info(msg) self.mqtt_client.publish(MQTT_REQ_TOPIC + self.__get_mqtt_customer_id() + "/Doors", msg) return True def preconditioning(self, vin, activate: bool): if activate: value = "activate" else: value = "deactivate" if vin in self.precond_programs: programs = self.precond_programs[vin] else: programs = DEFAULT_PRECONDITIONING_PROGRAM msg = self.mqtt_request(vin, {"asap": value, "programs": programs}) logger.info("preconditioning: %s", msg) self.mqtt_client.publish(MQTT_REQ_TOPIC + self.__get_mqtt_customer_id() + "/ThermalPrecond", msg) return True def save_config(self, name=None, force=False): if name is None: name = self.config_file config_str = json.dumps(self, cls=MyPeugeotEncoder, sort_keys=True, indent=4).encode("utf8") new_hash = md5(config_str).hexdigest() if force or self._config_hash != new_hash: with open(name, "wb") as f: f.write(config_str) self._config_hash = new_hash logger.info("save config change") @staticmethod def load_config(name="config.json"): with open(name, "r", encoding="utf-8") as f: config_str = f.read() config = dict(**json.loads(config_str)) if "country_code" not in config: config["country_code"] = input("What is your country code ? (ex: FR, GB, DE, ES...)\n") for new_el in ["abrp", "co2_signal_api"]: if new_el not in config: config[new_el] = None psacc = MyPSACC(**config) psacc.config_file = name return psacc def set_record(self, value: bool): self._record_enabled = value def record_info(self, car: Car): mileage = car.status.timed_odometer.mileage level = car.status.get_energy('Electric').level level_fuel = car.status.get_energy('Fuel').level charge_date = car.status.get_energy('Electric').updated_at moving = car.status.kinetic.moving longitude = car.status.last_position.geometry.coordinates[0] latitude = car.status.last_position.geometry.coordinates[1] altitude = car.status.last_position.geometry.coordinates[2] date = car.status.last_position.properties.updated_at if date is None: date = charge_date logger.debug("vin:%s longitude:%s latitude:%s date:%s mileage:%s level:%s charge_date:%s level_fuel:" "%s moving:%s", car.vin, longitude, latitude, date, mileage, level, charge_date, level_fuel, moving) Database.record_position(self.weather_api, car.vin, mileage, latitude, longitude, altitude, date, level, level_fuel, moving) self.abrp.call(car, Database.get_last_temp(car.vin)) try: charging_status = car.status.get_energy('Electric').charging.status charging_mode = car.status.get_energy('Electric').charging.charging_mode charging_rate = car.status.get_energy('Electric').charging.charging_rate autonomy = car.status.get_energy('Electric').autonomy Charging.record_charging(car, charging_status, charge_date, level, latitude, longitude, self.country_code, charging_mode, charging_rate, autonomy) logger.debug("charging_status:%s ", charging_status) except AttributeError: logger.error("charging status not available from api") def __iter__(self): for key, value in self.__dict__.items(): yield key, value class MyPeugeotEncoder(JSONEncoder): def default(self, mp: MyPSACC): # pylint: disable=arguments-renamed data = dict(mp) mpd = {"proxies": data["_proxies"], "refresh_token": mp.manager.refresh_token, "client_secret": mp.service_information.client_secret, "abrp": dict(mp.abrp)} for param in ["client_id", "realm", "remote_refresh_token", "customer_id", "weather_api", "country_code"]: mpd[param] = data[param] mpd["co2_signal_api"] = Ecomix.co2_signal_key return mpd