Merge pull request #274 from flobz/develop

Develop
This commit is contained in:
Florian BEZANNIER
2021-12-01 20:51:44 +01:00
committed by GitHub
29 changed files with 716 additions and 488 deletions
+1
View File
@@ -19,4 +19,5 @@ pylint:
- W0603
- I0011
- W0511
- consider-using-f-string
+9 -10
View File
@@ -7,15 +7,11 @@ from time import sleep
import pytz
from libs.psa.constants import DISCONNECTED, INPROGRESS, FINISHED
from libs.utils import RateLimitException
from my_psacc import MyPSACC
from mylogger import logger
DISCONNECTED = "Disconnected"
INPROGRESS = "InProgress"
FAILURE = "Failure"
STOPPED = "Stopped"
FINISHED = "Finished"
class ChargeControl:
MQTT_TIMEOUT = 60
@@ -42,7 +38,7 @@ class ChargeControl:
return self._stop_hour
def control_charge_with_ack(self, charge: bool):
self.psacc.charge_now(self.vin, charge)
self.psacc.remote_client.charge_now(self.vin, charge)
self.retry_count += 1
sleep(ChargeControl.MQTT_TIMEOUT)
vehicle_status = self.psacc.get_vehicle_info(self.vin)
@@ -51,7 +47,7 @@ class ChargeControl:
logger.warning("Car state isn't compatible with charging %s", status)
if (status == INPROGRESS) != charge:
logger.warning("retry to control the charge of %s", self.vin)
self.psacc.charge_now(self.vin, charge)
self.psacc.remote_client.charge_now(self.vin, charge)
self.retry_count += 1
return False
self.retry_count = 0
@@ -65,7 +61,10 @@ class ChargeControl:
else:
wakeup_timeout = self.wakeup_timeout
if (datetime.utcnow().replace(tzinfo=pytz.UTC) - last_update).total_seconds() > 60 * wakeup_timeout:
self.psacc.wakeup(self.vin)
try:
self.psacc.remote_client.wakeup(self.vin)
except RateLimitException:
logger.exception("force_update:")
def process(self):
now = datetime.now()
@@ -131,7 +130,7 @@ class ChargeControls(dict):
@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
+5 -3
View File
@@ -60,18 +60,20 @@ car_models = [
ElecModel("Mokka-e", 46, "opel:mokkae:20:48", r"VXKUKZKX.*"), # VXKUKZKXZM
ElecModel("Zaphira-e", 68, "peugeot:etraveler:21:75:opel", r"VXEVZZKX.*"), # VXEVZZKXZMZ
ElecModel("E-C4", 46, "citroen:ec4:21:50", r"VR7BCZKX.*"), # VR7BCZKXCM
CarModel("SUV 3008 Hybrid 225", 13.2, 43, reg=r"VF3M4DGZ.*"), #VF3M4DGZUMS
CarModel("SUV 3008 Hybrid 225", 13.2, 43, reg=r"VF3M4DGZ.*"), # VF3M4DGZUMS
CarModel("SUV 3008", 0, 53, reg=r"VF3MJEHZ.*"), # VF3MJEHZRJ
CarModel("308", 0, 56, reg=r"VF3L35GG.*"),
CarModel("208", 0, 44, reg=r"VR3UPHN[SE].*"), # VR3UPHNSSM VR3UPHNEKM
CarModel("2008", 0, 44, reg=r"VR3USHNS.*"), # VR3USHNSKM
CarModel("2008 II", 0, 45, reg=r"VR3USHNK.*"), # VR3USHNKKL
CarModel("SUV 5008 II", 0, 56, reg=r"VF3MRHNS.*"), # vf3mrhnsum
CarModel("SUV 5008 II 2018", 0, 56, reg=r"VF3MRHNY.*"), # VF3MRHNYHH
CarModel("N5008 GT-LINE 1.6L", 0, 60, reg=r"VF3MCBHZ.*"), # VF3MCBHZWJ
CarModel("C5 Aircross Hybrid", 13.2, 43, reg=r"VR7A4DGZ.*"), # VR7A4DGZSM
CarModel("DS7 Crossback E-Tense", 11.5, 43, reg="VR1J45GBUK.*"),
CarModel("DS7 Crossback E-Tense 300 4x4", 11.5, 43, reg="VR1J45GBUL.*"),
CarModel("DS7 Crossback E-Tense 300 4x4", 11.5, 43, reg="VR1J45GB.*"), # VR1J45GBUM VR1J45GBUL VR1J45GBUK
CarModel("508 SW Hybrid", 11.5, 45, reg=r"VR3F4DGZ.*"), # VR3F4DGZTL
CarModel("508 Hybrid", 11.5, 43, reg=r"VR3F3DGZ.*"), # VR3F3DGZTM
CarModel("508 SW 2.0 HDI 163CV", 0, 72, reg=r"VF38ERHH.*"), # VF38ERHH
CarModel("Grandland X Hybrid", 13.2, 43, reg=r"W0VZ4DGZ.*"), # W0VZ4DGZ2L
CarModel("Grandland X Hybrid 4x4", 13.2, 43, reg=r"W0VZ45GB.*") # W0VZ45GB3L
]
+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:]
+279
View File
@@ -0,0 +1,279 @@
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):
if self.load_otp():
self.mqtt_client = mqtt.Client(clean_session=True, protocol=mqtt.MQTTv311)
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()
logger.error("Can't configure MQTT Client")
return False
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()
message = mqtt_request.get_message_to_json(self.remoteCredentials.access_token)
logger.debug("%s %s", mqtt_request.topic, message)
self.mqtt_client.publish(mqtt_request.topic, message)
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
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 True
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 True
except (RequestException, RateLimitException) as e:
logger.exception("Can't refresh remote token %s", e)
sleep(60)
return False
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.get_car_by_vin(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.fromtimestamp(0)
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()
+38
View File
@@ -0,0 +1,38 @@
DELAYED_CHARGE = "delayed"
IMMEDIATE_CHARGE = "immediate"
PSA_CORRELATION_DATE_FORMAT = "%Y%m%d%H%M%S%f"
PSA_DATE_FORMAT = "%Y-%m-%dT%H:%M:%SZ"
realm_info = {
"clientsB2CPeugeot": {"oauth_url": "https://idpcvs.peugeot.com/am/oauth2/access_token", "app_name": "MyPeugeot"},
"clientsB2CCitroen": {"oauth_url": "https://idpcvs.citroen.com/am/oauth2/access_token", "app_name": "MyCitroen"},
"clientsB2CDS": {"oauth_url": "https://idpcvs.driveds.com/am/oauth2/access_token", "app_name": "MyDS"},
"clientsB2COpel": {"oauth_url": "https://idpcvs.opel.com/am/oauth2/access_token", "app_name": "MyOpel"},
"clientsB2CVauxhall": {"oauth_url": "https://idpcvs.vauxhall.co.uk/am/oauth2/access_token",
"app_name": "MyVauxhall"}
}
MQTT_BRANDCODE = {"AP": "AP",
"AC": "AC",
"DS": "AC",
"VX": "OV",
"OP": "OV"
}
DISCONNECTED = "Disconnected"
INPROGRESS = "InProgress"
FAILURE = "Failure"
STOPPED = "Stopped"
FINISHED = "Finished"
DEFAULT_PRECONDITIONING_PROGRAM = {
"program1": {"day": [0, 0, 0, 0, 0, 0, 0], "hour": 34, "minute": 7, "on": 0},
"program2": {"day": [0, 0, 0, 0, 0, 0, 0], "hour": 34, "minute": 7, "on": 0},
"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="
BRAND = {"com.psa.mym.myopel": {"realm": "clientsB2COpel", "brand_code": "OP", "app_name": "MyOpel"},
"com.psa.mym.mypeugeot": {"realm": "clientsB2CPeugeot", "brand_code": "AP", "app_name": "MyPeugeot"},
"com.psa.mym.mycitroen": {"realm": "clientsB2CCitroen", "brand_code": "AC", "app_name": "MyCitroen"},
"com.psa.mym.myds": {"realm": "clientsB2CDS", "brand_code": "DS", "app_name": "MyDS"},
"com.psa.mym.myvauxhall": {"realm": "clientsB2CVauxhall", "brand_code": "VX", "app_name": "MyVauxhall"}
}
+43
View File
@@ -0,0 +1,43 @@
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()
self.data = {}
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)
self.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 self.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
def __str__(self):
return "topic: " + self.topic + ": " + str(self.req_parameters)
+20 -3
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,11 +41,22 @@ 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):
super().__init__()
self.refresh_callback = None
self.refresh_callback = lambda: True
def set_refresh_callback(self, callback):
self.refresh_callback = callback
+54
View File
@@ -0,0 +1,54 @@
import json
import os
from androguard.core.bytecodes.apk import APK
from cryptography.hazmat.backends import default_backend
from cryptography.hazmat.primitives import serialization
from cryptography.hazmat.primitives.serialization import pkcs12
from libs.psa.constants import BRAND
class ApkParser:
def __init__(self, filename, country_code):
self.country_code = country_code
self.filename = filename
self.host_brandid_prod = None
self.site_code = None
self.culture = None
self.client_id = None
self.client_secret = None
@staticmethod
def __get_cultures_code(file, country_code):
cultures = json.loads(file)
return cultures[country_code]["languages"][0]
def retrieve_content_from_apk(self):
a = APK(self.filename)
package_name = a.get_package()
resources = a.get_android_resources() # .get_strings_resources()
self.client_id = resources.get_string(package_name, "PSA_API_CLIENT_ID_PROD")[1]
self.client_secret = resources.get_string(package_name, "PSA_API_CLIENT_SECRET_PROD")[1]
self.host_brandid_prod = resources.get_string(package_name, "HOST_BRANDID_PROD")[1]
self.culture = self.__get_cultures_code(a.get_file("res/raw/cultures.json"), self.country_code)
## Get Customer id
self.site_code = BRAND[package_name]["brand_code"] + "_" + self.country_code + "_ESP"
pfx_cert = a.get_file("assets/MWPMYMA1.pfx")
save_key_to_pem(pfx_cert, b"y5Y2my5B")
def save_key_to_pem(pfx_data, pfx_password):
private_key, certificate = pkcs12.load_key_and_certificates(pfx_data,
pfx_password, default_backend())[:2]
try:
os.mkdir("certs")
except FileExistsError:
pass
with open("certs/public.pem", "wb") as f:
f.write(certificate.public_bytes(encoding=serialization.Encoding.PEM))
with open("certs/private.pem", "wb") as f:
f.write(private_key.private_bytes(encoding=serialization.Encoding.PEM,
format=serialization.PrivateFormat.TraditionalOpenSSL,
encryption_algorithm=serialization.NoEncryption()))
@@ -1,84 +1,42 @@
#!/usr/bin/env python3
import json
import os
import traceback
from os import path
from androguard.core.bytecodes.apk import APK
import requests
from cryptography.hazmat.primitives import serialization
from cryptography.hazmat.primitives.serialization import pkcs12
from cryptography.hazmat.backends import default_backend
from charge_control import ChargeControl, ChargeControls
from libs.psa.constants import BRAND
from libs.psa.setup.apk_parser import ApkParser
from libs.psa.setup.github import urlretrieve_from_github
from my_psacc import MyPSACC
from mylogger import logger
BRAND = {"com.psa.mym.myopel": {"realm": "clientsB2COpel", "brand_code": "OP", "app_name": "MyOpel"},
"com.psa.mym.mypeugeot": {"realm": "clientsB2CPeugeot", "brand_code": "AP", "app_name": "MyPeugeot"},
"com.psa.mym.mycitroen": {"realm": "clientsB2CCitroen", "brand_code": "AC", "app_name": "MyCitroen"},
"com.psa.mym.myds": {"realm": "clientsB2CDS", "brand_code": "DS", "app_name": "MyDS"},
"com.psa.mym.myvauxhall": {"realm": "clientsB2CVauxhall", "brand_code": "VX", "app_name": "MyVauxhall"}
}
DOWNLOAD_URL = "https://github.com/flobz/psa_apk/raw/main/"
APP_VERSION = "1.33.0"
GITHUB_USER = "flobz"
GITHUB_REPO = "psa_apk"
def save_key_to_pem(pfx_data, pfx_password):
private_key, certificate = pkcs12.load_key_and_certificates(pfx_data,
pfx_password, default_backend())[:2]
try:
os.mkdir("certs")
except FileExistsError:
pass
with open("certs/public.pem", "wb") as f:
f.write(certificate.public_bytes(encoding=serialization.Encoding.PEM))
with open("certs/private.pem", "wb") as f:
f.write(private_key.private_bytes(encoding=serialization.Encoding.PEM,
format=serialization.PrivateFormat.TraditionalOpenSSL,
encryption_algorithm=serialization.NoEncryption()))
def urlretrieve(url, path):
with open(path, 'wb') as f:
r = requests.get(url, stream=True)
r.raise_for_status()
for chunk in r.iter_content(1024):
f.write(chunk)
def get_cultures_code(file, country_code):
cultures = json.loads(file)
return cultures[country_code]["languages"][0]
def get_content_from_apk(filename: str, country_code: str) -> ApkParser:
apk_parser = ApkParser(filename, country_code)
urlretrieve_from_github(GITHUB_USER, GITHUB_REPO, "", apk_parser.filename)
apk_parser.retrieve_content_from_apk()
return apk_parser
def firstLaunchConfig(package_name, client_email, client_password, country_code, # pylint: disable=too-many-locals
config_prefix=""):
filename = package_name.split(".")[-1] + ".apk"
if not path.exists(filename):
urlretrieve(DOWNLOAD_URL + filename, filename)
a = APK(filename)
package_name = a.get_package()
resources = a.get_android_resources() # .get_strings_resources()
client_id = resources.get_string(package_name, "PSA_API_CLIENT_ID_PROD")[1]
client_secret = resources.get_string(package_name, "PSA_API_CLIENT_SECRET_PROD")[1]
HOST_BRANDID_PROD = resources.get_string(package_name, "HOST_BRANDID_PROD")[1]
REMOTE_REFRESH_TOKEN = None
culture = get_cultures_code(a.get_file("res/raw/cultures.json"), country_code)
## Get Customer id
site_code = BRAND[package_name]["brand_code"] + "_" + country_code + "_ESP"
pfx_cert = a.get_file("assets/MWPMYMA1.pfx")
save_key_to_pem(pfx_cert, b"y5Y2my5B")
apk_parser = get_content_from_apk(filename, country_code)
try:
res = requests.post(HOST_BRANDID_PROD + "/GetAccessToken",
res = requests.post(apk_parser.host_brandid_prod + "/GetAccessToken",
headers={
"Connection": "Keep-Alive",
"Content-Type": "application/json",
"User-Agent": "okhttp/2.3.0"
},
params={"jsonRequest": json.dumps(
{"siteCode": site_code, "culture": "fr-FR", "action": "authenticate",
{"siteCode": apk_parser.site_code, "culture": "fr-FR", "action": "authenticate",
"fields": {"USR_EMAIL": {"value": client_email},
"USR_PASSWORD": {"value": client_password}}
}
@@ -87,7 +45,8 @@ def firstLaunchConfig(package_name, client_email, client_password, country_code,
token = res.json()["accessToken"]
except Exception as ex:
msg = traceback.format_exc() + f"\nHOST_BRANDID : {HOST_BRANDID_PROD} sitecode: {site_code}"
msg = traceback.format_exc() + f"\nHOST_BRANDID : {apk_parser.host_brandid_prod} " \
f"sitecode: {apk_parser.site_code}"
try:
msg += res.text
except: # pylint: disable=bare-except
@@ -98,11 +57,11 @@ def firstLaunchConfig(package_name, client_email, client_password, country_code,
res2 = requests.post(
f"https://mw-{BRAND[package_name]['brand_code'].lower()}-m2c.mym.awsmpsa.com/api/v1/user",
params={
"culture": culture,
"culture": apk_parser.culture,
"width": 1080,
"version": APP_VERSION
},
data=json.dumps({"site_code": site_code, "ticket": token}),
data=json.dumps({"site_code": apk_parser.site_code, "ticket": token}),
headers={
"Connection": "Keep-Alive",
"Content-Type": "application/json;charset=UTF-8",
@@ -125,7 +84,8 @@ def firstLaunchConfig(package_name, client_email, client_password, country_code,
logger.error(msg)
raise Exception(msg) from ex
# Psacc
psacc = MyPSACC(None, client_id, client_secret, REMOTE_REFRESH_TOKEN, customer_id, BRAND[package_name]["realm"],
psacc = MyPSACC(None, apk_parser.client_id, apk_parser.client_secret,
None, customer_id, BRAND[package_name]["realm"],
country_code)
psacc.connect(client_email, client_password)
psacc.save_config(name=config_prefix + "config.json")
+43
View File
@@ -0,0 +1,43 @@
from hashlib import sha1
import requests
from mylogger import logger
def get_github_sha_from_file(user, repo, directory, filename):
res = requests.get("https://api.github.com/repos/{}/{}/git/trees/main:{}".format(user, repo, directory)).json()
file_info = next((file for file in res["tree"] if file['path'] == filename))
return file_info["sha"]
def github_file_need_to_be_downloaded(user, repo, directory, filename):
try:
with open(filename, 'rb') as file_for_hash:
data = file_for_hash.read()
filesize = len(data)
prefix = "blob " + str(filesize) + "\0"
sha_of_downloaded_file = sha1(prefix.encode("utf-8") + data).hexdigest()
sha_of_git_file = get_github_sha_from_file(user, repo, directory, filename)
if sha_of_downloaded_file == sha_of_git_file:
logger.debug("locale file is the latest version")
return False
logger.debug("download last version of file")
except FileNotFoundError:
logger.debug("File not found, download file")
return True
def urlretrieve_from_github(user, repo, directory, filename, branch="main"):
if github_file_need_to_be_downloaded(user, repo, directory, filename):
with open(filename, 'wb') as f:
r = requests.get("https://github.com/{}/{}/raw/{}/{}{}".format(user, repo, branch, directory, filename),
headers={
"Accept": "application/vnd.github.VERSION.raw"
},
stream=True
)
r.raise_for_status()
for chunk in r.iter_content(1024):
f.write(chunk)
+13 -7
View File
@@ -26,19 +26,25 @@ def get_temp(latitude: str, longitude: str, api_key: str) -> float:
return None
class RateLimitException(Exception):
pass
def rate_limit(limit, every):
def limit_decorator(func):
semaphore = Semaphore(limit)
@wraps(func)
def wrapper(*args, **kwargs):
semaphore.acquire()
try:
return func(*args, **kwargs)
finally: # don't catch but ensure semaphore release
timer = Timer(every, semaphore.release)
timer.setDaemon(True) # allows the timer to be canceled on exit
timer.start()
if semaphore.acquire(blocking=False):
try:
return func(*args, **kwargs)
finally: # don't catch but ensure semaphore release
timer = Timer(every, semaphore.release)
timer.setDaemon(True) # allows the timer to be canceled on exit
timer.start()
else:
raise RateLimitException
return wrapper
+29 -323
View File
@@ -1,73 +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 otp.otp import load_otp, new_otp_session, 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
from web.abrp import Abrp
from web.db import Database
DELAYED_CHARGE = "delayed"
IMMEDIATE_CHARGE = "immediate"
PSA_CORRELATION_DATE_FORMAT = "%Y%m%d%H%M%S%f"
PSA_DATE_FORMAT = "%Y-%m-%dT%H:%M:%SZ"
realm_info = {
"clientsB2CPeugeot": {"oauth_url": "https://idpcvs.peugeot.com/am/oauth2/access_token", "app_name": "MyPeugeot"},
"clientsB2CCitroen": {"oauth_url": "https://idpcvs.citroen.com/am/oauth2/access_token", "app_name": "MyCitroen"},
"clientsB2CDS": {"oauth_url": "https://idpcvs.driveds.com/am/oauth2/access_token", "app_name": "MyDS"},
"clientsB2COpel": {"oauth_url": "https://idpcvs.opel.com/am/oauth2/access_token", "app_name": "MyOpel"},
"clientsB2CVauxhall": {"oauth_url": "https://idpcvs.vauxhall.co.uk/am/oauth2/access_token",
"app_name": "MyVauxhall"}
}
MQTT_BRANDCODE = {"AP": "AP",
"AC": "AC",
"DS": "AC",
"VX": "OV",
"OP": "OV"
}
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
# pylint: disable=too-many-instance-attributes,too-many-public-methods
class MyPSACC:
def connect(self, user, password):
self.manager.init_with_user_credentials_realm(user, password, self.realm)
@@ -84,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_now)
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
@@ -94,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:
@@ -116,20 +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)
def api(self) -> psac.VehiclesApi:
self.api_config.access_token = self.manager.access_token
api_instance = psac.VehiclesApi(OauthAPIClient(self.api_config))
@@ -137,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
@@ -184,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()
@@ -202,247 +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:
self.get_sms_otp_code()
otp_session = new_otp_session(otp_session)
self.otp = otp_session
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:
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
self.refresh_token()
try:
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 as e:
logger.error("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 and charge_info['rate'] == 0:
# 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 KeyError:
logger.exception("mqtt message:")
def start_mqtt(self):
self.load_otp()
self.mqtt_client = mqtt.Client(clean_session=True, protocol=mqtt.MQTTv311)
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:
self.wakeup(self.vehicles_list[0].vin)
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(3, 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 = {
"program1": {"day": [0, 0, 0, 0, 0, 0, 0], "hour": 34, "minute": 7, "on": 0},
"program2": {"day": [0, 0, 0, 0, 0, 0, 0], "hour": 34, "minute": 7, "on": 0},
"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}
}
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
@@ -456,7 +157,7 @@ class MyPSACC:
@staticmethod
def load_config(name="config.json"):
with open(name, "r") as f:
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:
@@ -506,14 +207,19 @@ class MyPSACC:
yield key, value
# pylint: disable=arguments-differ
class MyPeugeotEncoder(JSONEncoder):
def default(self, mp: MyPSACC):
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
def default(self, mp: MyPSACC): # pylint: disable=arguments-renamed
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 -7
View File
@@ -8,7 +8,7 @@ logging.addLevelName(DEBUG_LEVELV_NUM, "DEBUGV")
class CustomLogger(logging.Logger):
# pylint: disable=too-many-arguments,unused-argument,arguments-differ
# pylint: disable=too-many-arguments,unused-argument,arguments-renamed
def _log(self, level, msg, args, exc_info=None, extra=None, stack_info=False, exc_info_debug=False, **kwargs):
if exc_info_debug and self.isEnabledFor(logging.DEBUG):
exc_info = True
@@ -26,19 +26,18 @@ class CustomLogger(logging.Logger):
logging.setLoggerClass(CustomLogger)
# pylint: disable=invalid-name
logger = logging.getLogger("log")
file_handler = RotatingFileHandler(LOG_FILE, 'a', 1000000, 1, encoding='utf8')
formatter = logging.Formatter('%(asctime)s :: %(levelname)s :: %(message)s')
stream_handler = logging.StreamHandler()
def my_logger(file=LOG_FILE, handler_level=logging.INFO):
global logger
def my_logger(handler_level=logging.INFO):
logger.setLevel(handler_level)
formatter = logging.Formatter('%(asctime)s :: %(levelname)s :: %(message)s')
file_handler = RotatingFileHandler(file, 'a', 1000000, 1, encoding='utf8')
file_handler.setLevel(handler_level)
file_handler.setFormatter(formatter)
logger.addHandler(file_handler)
stream_handler = logging.StreamHandler()
stream_handler.setLevel(handler_level)
stream_handler.setFormatter(formatter)
logger.addHandler(stream_handler)
+10 -13
View File
@@ -16,7 +16,7 @@ from . import oaep
from .load import IWData
# pylint: disable=too-many-instance-attributes,invalid-name
# pylint: disable=invalid-name
CONFIG_NAME = "otp.bin"
@@ -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()
@@ -327,18 +327,15 @@ def load_otp(filename=CONFIG_NAME):
return None
def new_otp_session(old_otp_session: Otp = None, smscode=None, codepin=None):
def new_otp_session(smscode, codepin, old_otp_session: Otp = None, ):
if old_otp_session is None:
otp = Otp("bb8e981582b0f31353108fb020bead1c")
else:
otp = Otp("bb8e981582b0f31353108fb020bead1c", device_id=old_otp_session.device_id)
if smscode is None:
otp.smsCode = input("What is the code you just received by SMS ?")
otp.codepin = input("What is your app pin code ?")
else:
otp.smsCode = smscode
otp.codepin = codepin
otp.activation_start()
otp.activation_finalyze()
save_otp(otp)
return otp
otp.smsCode = smscode
otp.codepin = codepin
if otp.activation_start():
otp.activation_finalyze()
save_otp(otp)
return otp
return None
+1 -1
View File
@@ -1,5 +1,4 @@
#!/usr/bin/env python3
# pylint: disable=wrong-import-position
import os
import sys
from threading import Thread
@@ -12,6 +11,7 @@ if sys.version_info < (3, 6):
TestRequirements(DIR + "/requirements.txt").test_requirements()
# pylint: disable=wrong-import-position
import web.app
from libs.config import Config
from mylogger import logger
+38 -7
View File
@@ -7,6 +7,8 @@ from datetime import datetime, timedelta
from pytz import UTC
import libs.config
from libs.psa.setup.app_decoder import GITHUB_USER, GITHUB_REPO
from libs.psa.setup.github import github_file_need_to_be_downloaded
from psa_connectedcar import ApiClient
import psa_connectedcar as psacc
import reverse_geocode
@@ -22,15 +24,14 @@ from charge_control import ChargeControls
from test.utils import DATA_DIR, record_position, latitude, longitude, date0, date1, date2, date3, record_charging, \
vehicule_list, get_new_test_db, get_date, date4
from trip import Trips
from libs.utils import get_temp, parse_hour
from libs.utils import get_temp, parse_hour, rate_limit, RateLimitException
from web.db import Database
from web.figures import get_figures, get_battery_curve_fig, get_altitude_fig
from deepdiff import DeepDiff
def compare_dict(result, expected):
diff = DeepDiff(expected, result)
diff = DeepDiff(expected, result)
if diff != {}:
raise AssertionError(str(diff))
return True
@@ -95,7 +96,8 @@ class TestUnit(unittest.TestCase):
Ecomix.co2_signal_key = key
def_country = "FR"
Ecomix.get_data_from_co2_signal(latitude, longitude, def_country)
res = Ecomix.get_co2_from_signal_cache(datetime.utcnow().replace(tzinfo=UTC) - timedelta(minutes=5), datetime.now(), def_country)
res = Ecomix.get_co2_from_signal_cache(datetime.utcnow().replace(tzinfo=UTC) - timedelta(minutes=5),
datetime.now(), def_country)
assert isinstance(res, float)
def test_charge_control(self):
@@ -167,7 +169,7 @@ class TestUnit(unittest.TestCase):
end_level = 85
Charging.record_charging(car, "InProgress", date0, start_level, latitude, longitude, None, "slow", 20, 60)
Charging.record_charging(car, "InProgress", date1, 70, latitude, longitude, "FR", "slow", 20, 60)
Charging.record_charging(car, "InProgress", date1, 70, latitude, longitude, "FR", "slow",20, 60)
Charging.record_charging(car, "InProgress", date1, 70, latitude, longitude, "FR", "slow", 20, 60)
Charging.record_charging(car, "InProgress", date2, 80, latitude, longitude, "FR", "slow", 20, 60)
Charging.record_charging(car, "Stopped", date3, end_level, latitude, longitude, "FR", "slow", 20, 60)
chargings = Charging.get_chargings()
@@ -202,7 +204,7 @@ class TestUnit(unittest.TestCase):
'start_at': date0,
'consumption_by_temp': None,
'positions': {'lat': [latitude],
'long': [longitude]},
'long': [longitude]},
'duration': 40.0,
'speed_average': 28.5,
'distance': 19.0,
@@ -231,7 +233,7 @@ class TestUnit(unittest.TestCase):
'start_at': start,
'consumption_by_temp': None,
'positions': {'lat': [latitude],
'long': [longitude]},
'long': [longitude]},
'duration': 120.0,
'speed_average': 9.5,
'distance': 19.0,
@@ -253,6 +255,35 @@ class TestUnit(unittest.TestCase):
expected_res = [[2, 0, 0], [3, 14, 0], [0, 0, 2], [0, 30, 0]]
assert expected_res == [parse_hour(h) for h in ["PT2H", "PT3H14", "PT2S", "PT30M"]]
def test_rate_limit(self):
@rate_limit(2, 10)
def test_fct():
pass
test_fct()
test_fct()
try:
test_fct()
raise Exception("It should have raise RateLimitException")
except RateLimitException:
pass
def test_parse_apk(self):
from libs.psa.setup.app_decoder import get_content_from_apk
filename = "mypeugeot.apk"
try:
os.remove(filename)
except FileNotFoundError:
pass
assert get_content_from_apk(filename, "FR")
assert github_file_need_to_be_downloaded(GITHUB_USER, GITHUB_REPO, "", filename) is False
def test_file_need_to_be_updated(self):
filename = "mypeugeot.apk"
with open(filename, "w") as f:
f.write(" ")
assert github_file_need_to_be_downloaded(GITHUB_USER, GITHUB_REPO, "", filename) is True
if __name__ == '__main__':
my_logger(handler_level=os.environ.get("DEBUG_LEVEL", 20))
+3 -5
View File
@@ -11,7 +11,6 @@ from web.db import Database
class Points:
# pylint: disable= too-few-public-methods
def __init__(self, latitude, longitude):
self.latitude = latitude
self.longitude = longitude
@@ -21,7 +20,6 @@ class Points:
class Trip:
# pylint: disable= too-many-instance-attributes
def __init__(self):
self.start_at = None
self.end_at = None
@@ -58,7 +56,7 @@ class Trip:
try:
self.consumption_km = 100 * self.consumption / self.distance # kw/100 km
except TypeError:
raise ValueError("Distance not set")
raise ValueError("Distance not set") from TypeError
return self.consumption_km
def set_fuel_consumption(self, consumption) -> float:
@@ -135,8 +133,8 @@ class Trips(list):
logger.debugv("trip discarded")
return False
@staticmethod # noqa: MC0001
def get_trips(vehicles_list: Cars) -> Dict[str, "Trips"]:
@staticmethod
def get_trips(vehicles_list: Cars) -> Dict[str, "Trips"]: # noqa: MC0001
# pylint: disable=too-many-locals,too-many-statements,too-many-nested-blocks,too-many-branches
conn = Database.get_db()
vehicles = conn.execute("SELECT DISTINCT vin FROM position;").fetchall()
+3 -3
View File
@@ -1,6 +1,6 @@
import json
from datetime import datetime
from json import JSONDecodeError
from json.decoder import JSONDecodeError
import requests
@@ -14,7 +14,7 @@ class Abrp:
def __init__(self, token: str = "", abrp_enable_vin=None):
if abrp_enable_vin is None:
abrp_enable_vin = list()
abrp_enable_vin = []
self.token = token
self.abrp_enable_vin = set(abrp_enable_vin)
self.proxies = None
@@ -48,7 +48,7 @@ class Abrp:
verify=self.proxies is None)
logger.debug(response.text)
try:
return response.json()["status"] == "ok"
return json.loads(response.text)["status"] == "ok"
except (JSONDecodeError, KeyError):
logger.error("Bad response from ABRP API: %s", response.text)
return False
+3 -3
View File
@@ -12,13 +12,12 @@ try:
except ImportError:
from werkzeug import DispatcherMiddleware
from mylogger import logger
from mylogger import logger, file_handler
import importlib
# pylint: disable=invalid-name
app = None
dash_app = None
dispatcher = None
class MyProxyFix(ProxyFix):
@@ -47,9 +46,10 @@ def start_app(*args, **kwargs):
def config_flask(title, base_path, debug: bool, host, port, reloader=False, # pylint: disable=too-many-arguments
unminified=False, view="web.view.views"):
global app, dash_app, dispatcher
global app, dash_app
reload_view = app is not None
app = Flask(__name__)
app.logger.addHandler(file_handler)
try:
lang = locale.getlocale()[0].split("_")[0]
locale.setlocale(locale.LC_TIME, ".".join(locale.getlocale())) # make sure LC_TIME is set
+3 -4
View File
@@ -129,10 +129,9 @@ class Database:
db_file = Database.DEFAULT_DB_FILE
conn = CustomSqliteConnection(db_file, detect_types=sqlite3.PARSE_DECLTYPES | sqlite3.PARSE_COLNAMES)
conn.row_factory = sqlite3.Row
Database.__thread_lock.acquire()
if not Database.db_initialized:
Database.init_db(conn)
Database.__thread_lock.release()
with Database.__thread_lock:
if not Database.db_initialized:
Database.init_db(conn)
if update_callback:
conn.callbacks.append(Database.callback_fct)
return conn
+1 -2
View File
@@ -51,8 +51,7 @@ SUMMARY_CARDS = {"Average consumption": {"text": [card_value_div(AVG_CONSUM_KW,
# pylint: disable=too-many-locals
def get_figures(car: Car):
global consumption_fig, consumption_df, trips_map, consumption_fig_by_speed, table_fig, info, \
battery_table, consumption_fig_by_temp
global consumption_fig, trips_map, consumption_fig_by_speed, table_fig, battery_table, consumption_fig_by_temp
lats = [42, 41]
lons = [1, 2]
names = ["undefined", "undefined"]
-3
View File
@@ -64,9 +64,6 @@ class FigureFilter:
res = {table.src: table.date_columns for table in self.tables}
return res
def __get_table_src(self):
return [table.src for table in self.tables]
def __get_figures(self):
return {"graph": [graph.figure for graph in self.graphs],
"tables": [table.figure for table in self.tables],
+7 -6
View File
@@ -2,9 +2,9 @@ from dash import callback_context, html, dcc
from dash.exceptions import PreventUpdate
from flask import request
from app_decoder import firstLaunchConfig
from libs.psa.setup.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
@@ -94,7 +94,7 @@ config_otp_layout = dbc.Row(dbc.Col(className="col-md-12 col-lg-2 m-3", children
def log_layout():
with open(LOG_FILE, "r") as f:
with open(LOG_FILE, "r", encoding="utf-8") as f:
log_text = f.read()
return html.H3(className="m-2", children=["Log:", dbc.Container(
fluid=True,
@@ -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(smscode=sms_code, codepin=code_pin)
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()
+18 -6
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
@@ -12,6 +13,13 @@ CHARGE_SWITCH = "charge-switch"
PRECONDITIONING_SWITCH = "preconditioning-switch"
def convert_value_to_str(value):
try:
return str(int(value))
except TypeError:
return "-"
def get_control_tabs(config):
tabs = []
for car in config.myp.vehicles_list:
@@ -19,7 +27,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:
@@ -27,18 +35,22 @@ def get_control_tabs(config):
preconditionning_state = car.status.preconditionning.air_conditioning.status != "Disabled"
charging_state = car.status.get_energy('Electric').charging.status == "InProgress"
cards = {"Battery": {"text": [card_value_div("battery_value", "%",
value=str(int(car.status.get_energy('Electric').level)))],
value=convert_value_to_str(
car.status.get_energy('Electric').level))],
"src": "assets/images/battery-charge.svg"},
"Mileage": {"text": [card_value_div("mileage_value", "km",
value=str(int(car.status.timed_odometer.mileage)))],
value=convert_value_to_str(
car.status.timed_odometer.mileage))],
"src": "assets/images/mileage.svg"}
}
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:
+17 -10
View File
@@ -10,6 +10,7 @@ from flask import jsonify, request, Response as FlaskResponse
import web.utils
from libs.car import Cars, Car
from libs.utils import RateLimitException
from mylogger import logger
from trip import Trips
@@ -38,9 +39,10 @@ CONFIG = Config()
def add_header(el):
return dbc.Row([dbc.Col(dcc.Link(html.H1('My car info'), href="/", style={"text-decoration": "none"})),
return dbc.Row([dbc.Col(dcc.Link(html.H1('My car info'), href=dash_app.requests_pathname_external_prefix,
style={"text-decoration": "none"})),
dbc.Col(dcc.Link(html.Img(src="assets/images/settings.svg", width="30veh"),
href="/config",
href=dash_app.requests_pathname_external_prefix+"config",
className="float-end"))]), el
@@ -77,7 +79,7 @@ def create_callback(): # noqa: MC0001
[Input("battery-table", "data_timestamp")],
[State("battery-table", "data"),
State("battery-table", "data_previous")])
def capture_diffs_in_battery_table(timestamp, data, data_previous): # pylint: disable=unused-variable
def capture_diffs_in_battery_table(timestamp, data, data_previous):
if timestamp is None:
raise PreventUpdate
diff_data = diff_dashtable(data, data_previous, "start_at")
@@ -97,7 +99,7 @@ def create_callback(): # noqa: MC0001
Input("tab_battery_popup-close", "n_clicks")],
[State('battery-table', 'data'),
State("tab_battery_popup", "is_open")])
def get_battery_curve(active_cell, close, data, is_open): # pylint: disable=unused-argument, unused-variable
def get_battery_curve(active_cell, close, data, is_open): # pylint: disable=unused-argument
if is_open is None:
is_open = False
if active_cell is not None and active_cell["column_id"] in ["start_level", "end_level"] and not is_open:
@@ -109,7 +111,7 @@ def create_callback(): # noqa: MC0001
[Input("trips-table", "active_cell"),
Input("tab_trips_popup-close", "n_clicks")],
State("tab_trips_popup", "is_open"))
def get_altitude_graph(active_cell, close, is_open): # pylint: disable=unused-argument, unused-variable
def get_altitude_graph(active_cell, close, is_open): # pylint: disable=unused-argument
if is_open is None:
is_open = False
if active_cell is not None and active_cell["column_id"] in ["altitude_diff"] and not is_open:
@@ -147,7 +149,7 @@ STYLE_CACHE = None
def get_style():
global STYLE_CACHE
if not STYLE_CACHE:
with open(app.root_path + "/assets/style.json", "r") as f:
with open(app.root_path + "/assets/style.json", "r", encoding="utf-8") as f:
res = json.loads(f.read())
STYLE_CACHE = res
url_root = request.url_root
@@ -157,22 +159,27 @@ def get_style():
@app.route('/charge_now/<string:vin>/<int:charge>')
def charge_now(vin, charge):
return jsonify(CONFIG.myp.charge_now(vin, charge != 0))
return jsonify(CONFIG.myp.remote_client.charge_now(vin, charge != 0))
@app.route('/charge_hour')
def change_charge_hour():
return jsonify(CONFIG.myp.change_charge_hour(request.args['vin'], request.args['hour'], request.args['minute']))
return jsonify(CONFIG.myp.remote_client.change_charge_hour(request.args['vin'],
request.args['hour'],
request.args['minute']))
@app.route('/wakeup/<string:vin>')
def wakeup(vin):
return jsonify(CONFIG.myp.wakeup(vin))
try:
return jsonify(CONFIG.myp.remote_client.wakeup(vin))
except RateLimitException:
return jsonify({"error": "Wakeup rate limit exceeded"})
@app.route('/preconditioning/<string:vin>/<int:activate>')
def preconditioning(vin, activate):
return jsonify(CONFIG.myp.preconditioning(vin, activate))
return jsonify(CONFIG.myp.remote_client.preconditioning(vin, activate))
@app.route('/position/<string:vin>')