From 61f0d79d91786bf9e1e9ac7907d9d8363423d843 Mon Sep 17 00:00:00 2001 From: Florian Bezannier Date: Mon, 22 Nov 2021 23:15:35 +0100 Subject: [PATCH] refacto & improve remoteclient --- .prospector.yaml | 1 + charge_control.py | 6 +- libs/car.py | 7 +- libs/config.py | 21 +- libs/psa/AccountInformation.py | 13 + libs/psa/RemoteClient.py | 277 ++++++++++++++++ libs/psa/RemoteCredentials.py | 23 ++ libs/{psa_constants.py => psa/constants.py} | 2 + libs/psa/mqtt_request.py | 39 +++ libs/{ => psa}/oauth.py | 21 +- my_psacc.py | 330 ++------------------ otp/otp.py | 11 +- web/view/config_views.py | 9 +- web/view/control.py | 11 +- 14 files changed, 438 insertions(+), 333 deletions(-) create mode 100644 libs/psa/AccountInformation.py create mode 100644 libs/psa/RemoteClient.py create mode 100644 libs/psa/RemoteCredentials.py rename libs/{psa_constants.py => psa/constants.py} (88%) create mode 100644 libs/psa/mqtt_request.py rename libs/{ => psa}/oauth.py (82%) diff --git a/.prospector.yaml b/.prospector.yaml index 59e3ee8..d68319b 100644 --- a/.prospector.yaml +++ b/.prospector.yaml @@ -19,4 +19,5 @@ pylint: - W0603 - I0011 - W0511 + - consider-using-f-string diff --git a/charge_control.py b/charge_control.py index 049902a..0aeae28 100644 --- a/charge_control.py +++ b/charge_control.py @@ -7,7 +7,7 @@ from time import sleep import pytz -from libs.psa_constants import DISCONNECTED, INPROGRESS, FINISHED +from libs.psa.constants import DISCONNECTED, INPROGRESS, FINISHED from libs.utils import RateLimitException from my_psacc import MyPSACC from mylogger import logger @@ -123,14 +123,14 @@ class ChargeControls(dict): config_str = json.dumps(chd, sort_keys=True, indent=4).encode('utf-8') new_hash = md5(config_str).hexdigest() if force or self._config_hash != new_hash: - with open(self.file_name, "wb") as f: + with open(self.file_name, "wb", encoding="utf-8") as f: f.write(config_str) self._config_hash = new_hash logger.info("save config change") @staticmethod def load_config(psacc: MyPSACC, name="charge_config.json"): - with open(name, "r") as file: + with open(name, "r", encoding="utf-8") as file: config_str = file.read() chd = json.loads(config_str) charge_control_list = ChargeControls(name) diff --git a/libs/car.py b/libs/car.py index 2b2949d..a2c7e4f 100644 --- a/libs/car.py +++ b/libs/car.py @@ -5,7 +5,8 @@ from mylogger import logger from libs.car_model import CarModel from libs.car_status import CarStatus -# pylint: disable=too-many-instance-attributes,too-many-arguments + +# pylint: disable=too-many-arguments class Car: def __init__(self, vin, vehicle_id, brand, label=None, battery_power=None, fuel_capacity=None, max_elec_consumption=None, max_fuel_consumption=None, abrp_name=None): @@ -118,7 +119,7 @@ class Cars(list): if name is None: name = self.config_filename config_str = json.dumps(self, default=lambda car: car.to_dict(), sort_keys=True, indent=4) - with open(name, "w") as file: + with open(name, "w", encoding="utf-8") as file: file.write(config_str) @staticmethod @@ -126,7 +127,7 @@ class Cars(list): if name is None: name = Cars().config_filename try: - with open(name, "r") as file: + with open(name, "r",encoding="utf-8") as file: json_str = file.read() cars = Cars.from_json(json.loads(json_str)) cars.config_filename = name diff --git a/libs/config.py b/libs/config.py index 30a8491..b56adba 100644 --- a/libs/config.py +++ b/libs/config.py @@ -10,7 +10,7 @@ from libs.charging import Charging from libs.elec_price import ElecPrice from my_psacc import MyPSACC from mylogger import logger, my_logger -from otp.otp import CONFIG_NAME as OTP_CONFIG_NAME +from otp.otp import CONFIG_NAME as OTP_CONFIG_NAME, ConfigException DEFAULT_NAME = "config.json" @@ -58,13 +58,16 @@ class Config(metaclass=Singleton): if self.args.remote_disable: logger.info("mqtt disabled") elif not self.args.web_conf or path.isfile(OTP_CONFIG_NAME): - if self.myp.mqtt_client is not None: - self.myp.mqtt_client.disconnect() - self.myp.start_mqtt() - if self.args.charge_control: - Config.chc = ChargeControls.load_config(self.myp, name=self.args.charge_control) - Config.chc.init() - self.myp.start_refresh_thread() + if self.myp.remote_client.mqtt_client is not None: + self.myp.remote_client.mqtt_client.disconnect() + try: + self.myp.remote_client.start() + if self.args.charge_control: + Config.chc = ChargeControls.load_config(self.myp, name=self.args.charge_control) + Config.chc.init() + self.myp.start_refresh_thread() + except ConfigException: + logger.error("start_remote_control failed redo otp config") def load_app(self) -> bool: my_logger(handler_level=int(self.args.debug)) @@ -86,7 +89,7 @@ class Config(metaclass=Singleton): else: self.is_good = False try: - self.is_good = self.myp.refresh_token() + self.is_good = self.myp.manager.refresh_token_now() if self.is_good: logger.info(str(self.myp.get_vehicles())) except OAuthError: diff --git a/libs/psa/AccountInformation.py b/libs/psa/AccountInformation.py new file mode 100644 index 0000000..9c8945a --- /dev/null +++ b/libs/psa/AccountInformation.py @@ -0,0 +1,13 @@ +from libs.psa.constants import MQTT_BRANDCODE + + +class AccountInformation: + def __init__(self, client_id, customer_id: str, realm: str, country_code: str): + self.client_id = client_id + self.customer_id = customer_id + self.realm = realm + self.country_code = country_code + + def get_mqtt_customer_id(self): + brand_code = self.customer_id[:2] + return MQTT_BRANDCODE[brand_code] + self.customer_id[2:] diff --git a/libs/psa/RemoteClient.py b/libs/psa/RemoteClient.py new file mode 100644 index 0000000..70b10a4 --- /dev/null +++ b/libs/psa/RemoteClient.py @@ -0,0 +1,277 @@ +import json +import threading +from datetime import datetime +from os import environ +from time import sleep + +import paho.mqtt.client as mqtt +from requests import RequestException + +from libs.car import Cars +from libs.psa.AccountInformation import AccountInformation +from libs.psa.RemoteCredentials import RemoteCredentials +from libs.psa.constants import INPROGRESS, DEFAULT_PRECONDITIONING_PROGRAM, IMMEDIATE_CHARGE, DELAYED_CHARGE, REMOTE_URL +from libs.psa.mqtt_request import MQTTRequest +from libs.psa.oauth import OpenIdCredentialManager +from libs.utils import RateLimitException, rate_limit, parse_hour +from mylogger import logger +from otp.otp import ConfigException, save_otp, load_otp + +MQTT_SERVER = "mwa.mpsa.com" +MQTT_RESP_TOPIC = "psa/RemoteServices/to/cid/" +MQTT_EVENT_TOPIC = "psa/RemoteServices/events/MPHRTServices/" +MQTT_TOKEN_TTL = 890 + + +class RemoteClient: + + def __init__(self, account_info: AccountInformation, vehicles_list: Cars, manager: OpenIdCredentialManager, + remoteCredentials: RemoteCredentials): + self.vehicles_list = vehicles_list + self.remoteCredentials: RemoteCredentials = remoteCredentials + self.manager = manager + self.precond_programs = {} + self.account_info = account_info + self.headers = { + "x-introspect-realm": self.account_info.realm, + "accept": "application/hal+json", + "User-Agent": "okhttp/4.8.0", + } + self.last_request = [] + self.mqtt_client = None + self.otp = None + + def __on_mqtt_connect(self, client, userdata, result_code, _): # pylint: disable=unused-argument + logger.info("Connected with result code %s", result_code) + topics = [MQTT_RESP_TOPIC + self.account_info.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) + + def _on_mqtt_disconnect(self, client, userdata, result_code): # pylint: disable=unused-argument + 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)) + + def __on_mqtt_message(self, client, userdata, msg): + try: + logger.info("mqtt msg received: %s %s", msg.topic, msg.payload) + logger.debug("client: %s userdata: %s", client, userdata) + 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) + if self.last_request: + logger.warning("last request is send again, token was expired") + last_request = self.last_request + self.last_request = None + self.publish(last_request, store=False) + else: + logger.error("Last request might have been send twice without success") + 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(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 veh_charge_request(self, vin, hour, minute, charge_type): + msg = self.mqtt_request(vin, {"program": {"hour": hour, "minute": minute}, "type": charge_type}, "/VehCharge") + logger.info("veh_charge_request: %s", msg) + self.publish(msg) + return msg + + def publish(self, mqtt_request: MQTTRequest, store=True): + self._refresh_remote_token() + self.mqtt_client.publish(mqtt_request.topic, + mqtt_request.get_message_to_json(self.remoteCredentials.access_token)) + if store: + self.last_request = [mqtt_request] + + def mqtt_request(self, vin, req_parameters, topic): + return MQTTRequest(topic, vin, req_parameters, self.account_info.get_mqtt_customer_id()) + + def _refresh_remote_token(self, force=False): + bad_remote_token = self.remoteCredentials.refresh_token is None + res = None + if not force and not bad_remote_token and self.remoteCredentials.last_update: + last_update: datetime = self.remoteCredentials.last_update + if (datetime.now() - last_update).total_seconds() < MQTT_TOKEN_TTL: + return res + try: + self.manager.refresh_token_now() + if bad_remote_token: + logger.error("remote_refresh_token isn't defined") + else: + res = self.manager.post(REMOTE_URL + self.account_info.client_id, + json={"grant_type": "refresh_token", + "refresh_token": self.remoteCredentials.refresh_token}, + headers=self.headers) + data = res.json() + logger.debug("refresh_remote_token: %s", data) + if "access_token" in data: + self.remoteCredentials.access_token = data["access_token"] + self.remoteCredentials.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.remoteCredentials.access_token) + return res + except (RequestException, RateLimitException) as e: + logger.exception("Can't refresh remote token %s", e) + sleep(60) + return None + + def get_sms_otp_code(self): + res = self.manager.post( + "https://api.groupe-psa.com/applications/cvs/v4/mobile/smsCode?client_id=" + self.account_info.client_id, + headers=self.headers) + 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.account_info.client_id, + json={"grant_type": "password", "password": password}, + headers=self.headers) + data = res.json() + self.remoteCredentials.access_token = data["access_token"] + self.remoteCredentials.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 horn(self, vin, count): + msg = self.mqtt_request(vin, {"nb_horn": count, "action": "activate"}, "/Horn") + logger.info(msg) + self.mqtt_client.publish(msg) + + def lights(self, vin, duration: int): + msg = self.mqtt_request(vin, {"action": "activate", "duration": duration}, "/Lights") + logger.info(msg) + self.publish(msg) + + @rate_limit(6, 60 * 20) + def wakeup(self, vin): + logger.info("ask wakeup to %s", vin) + msg = self.mqtt_request(vin, {"action": "state"}, "/VehCharge/state") + logger.info(msg) + self.publish(msg) + return True + + def lock_door(self, vin, lock: bool): + if lock: + value = "lock" + else: + value = "unlock" + + msg = self.mqtt_request(vin, {"action": value}, "/Doors") + logger.info(msg) + self.publish(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}, "/ThermalPrecond") + logger.info("Preconditioning: %s", msg) + self.publish(msg) + return True + + 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 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 __get_charge_hour(self, vin): + hour_str = self.vehicles_list[vin].status.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 diff --git a/libs/psa/RemoteCredentials.py b/libs/psa/RemoteCredentials.py new file mode 100644 index 0000000..5bbebe5 --- /dev/null +++ b/libs/psa/RemoteCredentials.py @@ -0,0 +1,23 @@ +from datetime import datetime + + +class RemoteCredentials: + def __init__(self, remote_refresh_token): + self._refresh_token = remote_refresh_token + self.access_token = None + self.update_callbacks = [] + self.last_update = datetime.now() + + def __update_callbacks(self): + self.last_update = datetime.now() + for update_callback in self.update_callbacks: + update_callback() + + @property + def refresh_token(self): + return self._refresh_token + + @refresh_token.setter + def refresh_token(self, remote_refresh_token): + self._refresh_token = remote_refresh_token + self.__update_callbacks() diff --git a/libs/psa_constants.py b/libs/psa/constants.py similarity index 88% rename from libs/psa_constants.py rename to libs/psa/constants.py index 3d0cbac..e660990 100644 --- a/libs/psa_constants.py +++ b/libs/psa/constants.py @@ -28,3 +28,5 @@ DEFAULT_PRECONDITIONING_PROGRAM = { "program3": {"day": [0, 0, 0, 0, 0, 0, 0], "hour": 34, "minute": 7, "on": 0}, "program4": {"day": [0, 0, 0, 0, 0, 0, 0], "hour": 34, "minute": 7, "on": 0} } +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=" \ No newline at end of file diff --git a/libs/psa/mqtt_request.py b/libs/psa/mqtt_request.py new file mode 100644 index 0000000..68fda28 --- /dev/null +++ b/libs/psa/mqtt_request.py @@ -0,0 +1,39 @@ +import json +from datetime import datetime, timedelta +from uuid import uuid4 + +from libs.psa.constants import PSA_DATE_FORMAT, PSA_CORRELATION_DATE_FORMAT + +MQTT_REQ_TOPIC = "psa/RemoteServices/from/cid/" + + +class MQTTRequest: + TIME_BEFORE_EXPIRATION = 30 + + def __init__(self, topic, vin, req_parameters, customer_id): + self.customer_id = customer_id + self.topic = MQTT_REQ_TOPIC + self.customer_id + topic + self.vin = vin + self.req_parameters = req_parameters + self.date = datetime.now() + + def get_message_to_json(self, remote_access_token): + return json.dumps(self.get_message(remote_access_token)) + + def get_message(self, remote_access_token): + date = datetime.utcnow() + date_str = date.strftime(PSA_DATE_FORMAT) + data = {"access_token": remote_access_token, "customer_id": self.customer_id, + "correlation_id": self.__gen_correlation_id(date), "req_date": date_str, "vin": self.vin, + "req_parameters": self.req_parameters} + return data + + def is_expired(self): + return self.date < datetime.now() - timedelta(seconds=self.TIME_BEFORE_EXPIRATION) + + @staticmethod + def __gen_correlation_id(date): + date_str = date.strftime(PSA_CORRELATION_DATE_FORMAT)[:-3] + uuid_str = str(uuid4()).replace("-", "") + correlation_id = uuid_str + date_str + return correlation_id \ No newline at end of file diff --git a/libs/oauth.py b/libs/psa/oauth.py similarity index 82% rename from libs/oauth.py rename to libs/psa/oauth.py index f2445a7..56b9685 100644 --- a/libs/oauth.py +++ b/libs/psa/oauth.py @@ -1,15 +1,21 @@ from http import HTTPStatus +from typing import Optional -from oauth2_client.credentials_manager import CredentialManager -from requests import Response +from oauth2_client.credentials_manager import CredentialManager, ServiceInformation +from requests import Response, RequestException import psa_connectedcar as psac +from libs.utils import rate_limit from mylogger import logger from psa_connectedcar import ApiClient from psa_connectedcar.rest import ApiException class OpenIdCredentialManager(CredentialManager): + def __init__(self, service_information: ServiceInformation, proxies: Optional[dict] = None): + super().__init__(service_information, proxies) + self.refresh_callbacks = [] + def _grant_password_request_realm(self, login: str, password: str, realm: str) -> dict: return dict(grant_type='password', username=login, @@ -35,6 +41,17 @@ class OpenIdCredentialManager(CredentialManager): def access_token(self): return self._access_token + @rate_limit(6, 1800) + def refresh_token_now(self): + try: + self._refresh_token() + for refresh_callback in self.refresh_callbacks: + refresh_callback() + return True + except RequestException as e: + logger.error("Can't refresh token %s", e) + return False + class Oauth2PSACCApiConfig(psac.Configuration): def __init__(self): diff --git a/my_psacc.py b/my_psacc.py index ec7dfd7..fdcb013 100644 --- a/my_psacc.py +++ b/my_psacc.py @@ -1,51 +1,32 @@ 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 libs.psa.AccountInformation import AccountInformation +from libs.psa.RemoteClient import RemoteClient +from libs.psa.RemoteCredentials import RemoteCredentials +from libs.psa.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 libs.psa.constants import realm_info, AUTHORIZE_SERVICE 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) @@ -62,9 +43,9 @@ class MyPSACC: 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.api_config.set_refresh_callback(self.manager.refresh_token) self.manager.refresh_token = refresh_token - self.remote_refresh_token = remote_refresh_token + self.account_info = AccountInformation(client_id, customer_id, realm, country_code) self.remote_access_token = None self.vehicles_list = Cars.load_cars(CARS_FILE) self.customer_id = customer_id @@ -72,18 +53,10 @@ class MyPSACC: 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: @@ -94,21 +67,16 @@ class MyPSACC: self.config_file = DEFAULT_CONFIG_FILENAME Ecomix.co2_signal_key = co2_signal_api self.refresh_thread = None + remote_credentials = RemoteCredentials(remote_refresh_token) + remote_credentials.update_callbacks.append(self.save_config) + self.remote_client = RemoteClient(self.account_info, + self.vehicles_list, + self.manager, + remote_credentials) def get_app_name(self): return realm_info[self.realm]['app_name'] - @rate_limit(6, 1800) - def refresh_token(self): - try: - # pylint: disable=protected-access - self.manager._refresh_token() - self.save_config() - return True - except RequestException as e: - logger.error("Can't refresh token %s", e) - return False - def api(self) -> psac.VehiclesApi: self.api_config.access_token = self.manager.access_token api_instance = psac.VehiclesApi(OauthAPIClient(self.api_config)) @@ -116,14 +84,12 @@ class MyPSACC: def set_proxies(self, proxies): if proxies is None: - self._proxies = dict(http='', https='') + 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) + self.manager.proxies = proxies def get_vehicle_info(self, vin, cache=False): res = None @@ -163,14 +129,6 @@ class MyPSACC: 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() @@ -181,251 +139,11 @@ class MyPSACC: 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): - 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 @@ -492,10 +210,16 @@ class MyPSACC: 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 + mpd = {"proxies": mp.manager.proxies, + "refresh_token": mp.manager.refresh_token, + "client_secret": mp.service_information.client_secret, + "abrp": dict(mp.abrp), + "remote_refresh_token": mp.remote_client.remoteCredentials.refresh_token, + "customer_id": mp.account_info.customer_id, + "client_id": mp.account_info.client_id, + "realm": mp.account_info.realm, + "country_code": mp.account_info.country_code, + "weather_api": mp.weather_api, + "co2_signal_api": Ecomix.co2_signal_key + } return mpd diff --git a/otp/otp.py b/otp/otp.py index 2c644c5..9db251d 100644 --- a/otp/otp.py +++ b/otp/otp.py @@ -189,7 +189,7 @@ class Otp: elif self.mode == Otp.OTP_MODE: self.challenge = xml["challenge"] return True - return False + raise ConfigException(xml) def activation_finalyze(self, random_bytes=None): R = self.get_r() @@ -334,7 +334,8 @@ def new_otp_session(smscode, codepin, old_otp_session: Otp = None, ): otp = Otp("bb8e981582b0f31353108fb020bead1c", device_id=old_otp_session.device_id) otp.smsCode = smscode otp.codepin = codepin - otp.activation_start() - otp.activation_finalyze() - save_otp(otp) - return otp + if otp.activation_start(): + otp.activation_finalyze() + save_otp(otp) + return otp + return None diff --git a/web/view/config_views.py b/web/view/config_views.py index a4e7bcc..087c14d 100644 --- a/web/view/config_views.py +++ b/web/view/config_views.py @@ -4,7 +4,7 @@ from flask import request from app_decoder import firstLaunchConfig from libs.config import Config -from mylogger import LOG_FILE +from mylogger import LOG_FILE, logger from otp.otp import new_otp_session from web.app import dash_app import dash_bootstrap_components as dbc @@ -144,7 +144,7 @@ def askCode(n_clicks): # pylint: disable=unused-argument ctx = callback_context if ctx.triggered: try: - config.myp.get_sms_otp_code() + config.myp.remote_client.get_sms_otp_code() return dbc.Alert("Sms sent", color="success") except Exception as e: res = str(e) @@ -161,13 +161,14 @@ def finishOtp(n_clicks, code_pin, sms_code): # pylint: disable=unused-argument ctx = callback_context if ctx.triggered: try: - otp_session = new_otp_session(sms_code, code_pin, config.myp.otp) - config.myp.otp = otp_session + otp_session = new_otp_session(sms_code, code_pin, config.myp.remote_client.otp) + config.myp.remote_client.otp = otp_session config.myp.save_config() config.start_remote_control() return dbc.Alert(["OTP config finish !!! ", html.A("Go to home", href=request.url_root)], color="success") except Exception as e: res = str(e) + logger.exception("finishOtp:") return dbc.Alert(res, color="danger") raise PreventUpdate() diff --git a/web/view/control.py b/web/view/control.py index b90d69d..57e285a 100644 --- a/web/view/control.py +++ b/web/view/control.py @@ -1,6 +1,7 @@ import dash_bootstrap_components as dbc from dash import html +from my_psacc import MyPSACC from mylogger import logger from web.tools.Button import Button from web.tools.Switch import Switch @@ -19,7 +20,7 @@ def get_control_tabs(config): label = car.vin else: label = car.label - myp = config.myp + myp: MyPSACC = config.myp el = [] buttons_row = [] if config.remote_control: @@ -35,10 +36,12 @@ def get_control_tabs(config): } el.append(dbc.Container(dbc.Row(children=create_card(cards)), fluid=True)) buttons_row.extend([Button(REFRESH_SWITCH, car.vin, - html.Img(src="assets/images/sync.svg", width="50px"), myp.wakeup).get_html(), - Switch(CHARGE_SWITCH, car.vin, "Charge", myp.charge_now, charging_state).get_html(), + html.Img(src="assets/images/sync.svg", width="50px"), + myp.remote_client.wakeup).get_html(), + Switch(CHARGE_SWITCH, car.vin, "Charge", myp.remote_client.charge_now, + charging_state).get_html(), Switch(PRECONDITIONING_SWITCH, car.vin, "Preconditioning", - myp.preconditioning, preconditionning_state).get_html()]) + myp.remote_client.preconditioning, preconditionning_state).get_html()]) except (AttributeError, TypeError): logger.exception("get_control_tabs:") if not config.offline: