base.py 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154
  1. from abc import ABC, abstractmethod
  2. from typing import Optional, Dict, Union, final
  3. import logging
  4. logger = logging.getLogger("scada_productor_base")
  5. from pathlib import Path
  6. SCRIPT_DIR = Path(__file__).parent.absolute()
  7. import os
  8. import dotenv
  9. dotenv.load_dotenv()
  10. import requests
  11. class ScadaProductorBase(ABC):
  12. def __init__(self):
  13. self._api_base_url: Optional[str] = os.getenv("DEFAULT_API_BASE_URL")
  14. self._login_user: Optional[str] = os.getenv("DEFAULT_LOGIN_USER")
  15. self._login_password: Optional[str] = os.getenv("DEFAULT_LOGIN_PASSWORD")
  16. self._login_dep_id: Optional[str] = os.getenv("DEFAULT_LOGIN_DEP_ID")
  17. self._scada_secret: Optional[str] = os.getenv("DEFAULT_SCADA_SECRET")
  18. self._token: Optional[str] = None
  19. self._headers: Dict[str, str] = {"Content-Type": "application/json"}
  20. # ======================== Token 管理 ========================
  21. @final
  22. def _update_token(self, max_retries: int = 3) -> Union[str, bool]:
  23. """
  24. 向云平台请求登录,更新 JWT Token。
  25. Returns:
  26. str — 成功时返回新 token
  27. False — 重试耗尽后仍失败
  28. """
  29. url = f"{self._api_base_url}/api/v2/user/login"
  30. payload = {
  31. "UserName": self._login_user,
  32. "Password": self._login_password,
  33. "type": "account",
  34. "DepId": self._login_dep_id,
  35. }
  36. for attempt in range(1, max_retries + 1):
  37. try:
  38. resp = requests.post(url, json=payload, timeout=5)
  39. resp.raise_for_status()
  40. data = resp.json()
  41. if data.get("code") != 200:
  42. raise RuntimeError(f"登录失败: {data.get('msg', '未知错误')}")
  43. # 从响应中提取 token(兼容多种响应格式)
  44. resp_data = data.get("data")
  45. if isinstance(resp_data, dict):
  46. new_token = resp_data.get("token") or resp_data.get("Token")
  47. elif isinstance(resp_data, str):
  48. new_token = resp_data
  49. else:
  50. new_token = data.get("token")
  51. if not new_token:
  52. raise RuntimeError("登录响应中未找到 token")
  53. self._token = new_token
  54. self._headers["JWT-TOKEN"] = self._token
  55. logger.info(f"Token 更新成功 (第 {attempt} 次尝试)")
  56. return self._token
  57. except requests.Timeout:
  58. logger.warning(f"登录超时 (第 {attempt}/{max_retries} 次)", exc_info=True)
  59. except requests.RequestException as e:
  60. logger.warning(f"网络异常: {e} (第 {attempt}/{max_retries} 次)", exc_info=True)
  61. except (ValueError, KeyError) as e:
  62. logger.warning(f"响应解析异常: {e} (第 {attempt}/{max_retries} 次)", exc_info=True)
  63. except RuntimeError as e:
  64. logger.warning(f"{e} (第 {attempt}/{max_retries} 次)", exc_info=True)
  65. logger.error(f"Token 更新失败,已重试 {max_retries} 次")
  66. return False
  67. # ======================== 通用请求 ========================
  68. @final
  69. def _request(
  70. self,
  71. method: str,
  72. url: str,
  73. *,
  74. params: Optional[dict] = None,
  75. json: Optional[Union[dict, list]] = None,
  76. timeout: int = 30,
  77. max_retries: int = 3,
  78. ) -> Optional[dict]:
  79. """
  80. 带自动 Token 刷新的通用请求方法。
  81. 请求失败时(401 / 网络异常 / 响应异常)自动刷新 Token 并重试,
  82. 最多重试 max_retries 次。
  83. Parameters:
  84. method: "get" 或 "post"
  85. url: 完整请求 URL
  86. params: GET 查询参数
  87. json: POST 请求体
  88. timeout: 超时秒数
  89. max_retries: 最大重试次数
  90. Returns:
  91. dict — 成功时返回解析后的 JSON 响应体
  92. None — 重试耗尽仍失败
  93. """
  94. requester = getattr(requests, method.lower())
  95. kwargs = {"headers": self._headers, "timeout": timeout}
  96. if params is not None:
  97. kwargs["params"] = params
  98. if json is not None:
  99. kwargs["json"] = json
  100. for attempt in range(1, max_retries + 1):
  101. try:
  102. resp = requester(url, **kwargs)
  103. # Token 过期 → 刷新后重试(不消耗重试次数)
  104. if resp.status_code != 200:
  105. if resp.status_code == 601: logger.warning(f"请求返回 601,尝试刷新 Token")
  106. else : logger.warning(f"请求失败,尝试刷新 Token")
  107. if self._update_token():
  108. kwargs["headers"] = self._headers
  109. # 重试,不要浪费一次循环
  110. resp = requester(url, **kwargs)
  111. else:
  112. logger.error("Token 刷新失败,放弃重试")
  113. return None
  114. resp.raise_for_status()
  115. if resp.status_code == 200:
  116. data = resp.json()
  117. return data
  118. except requests.Timeout:
  119. logger.warning(f"请求超时 {url} (第 {attempt}/{max_retries} 次)", exc_info=True)
  120. except requests.RequestException as e:
  121. logger.warning(f"请求异常 {url}: {e} (第 {attempt}/{max_retries} 次)", exc_info=True)
  122. except (ValueError, KeyError) as e:
  123. logger.warning(f"响应解析异常 {url}: {e} (第 {attempt}/{max_retries} 次)", exc_info=True)
  124. except Exception as e:
  125. logger.warning(f"请求异常 {url}: {e}", exc_info=True)
  126. logger.error(f"请求失败,已重试 {max_retries} 次: {method.upper()} {url}")
  127. return None