refacto & improve remoteclient

This commit is contained in:
Florian Bezannier
2021-11-22 23:15:35 +01:00
parent ee36531e6b
commit 61f0d79d91
14 changed files with 438 additions and 333 deletions
+1
View File
@@ -19,4 +19,5 @@ pylint:
- W0603
- I0011
- W0511
- consider-using-f-string
+3 -3
View File
@@ -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)
+4 -3
View File
@@ -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
+12 -9
View File
@@ -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:
+13
View File
@@ -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:]
+277
View File
@@ -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
+23
View File
@@ -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()
@@ -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="
+39
View File
@@ -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
+19 -2
View File
@@ -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):
+27 -303
View File
@@ -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
+6 -5
View File
@@ -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
+5 -4
View File
@@ -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()
+7 -4
View File
@@ -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: