vs_cloud_api.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370
  1. # toolkit/vs_cloud_api.py
  2. import requests
  3. import json
  4. import time
  5. from urllib.parse import quote
  6. from datetime import datetime
  7. from typing import Dict, Any, List, Optional
  8. import configure
  9. from vs_types import NotFoundError, PermissionDeniedError, RateLimiteddError, SessionExpiredOrInvalidError, BizLogicError
  10. from vs_log_macros import VSC_ERROR, VSC_INFO, VSC_WARN, VSC_DEBUG
  11. class VSCloudApi:
  12. """
  13. @brief VSCloudApi 的 Python 实现
  14. 用于对接云端服务 (打码、邮件、Session存储、任务调度)
  15. """
  16. _instance = None
  17. def __new__(cls, *args, **kwargs):
  18. if cls._instance is None:
  19. cls._instance = super(VSCloudApi, cls).__new__(cls)
  20. cls._instance.base_url = "https://api.text.skin"
  21. cls._instance.api_token = "Bearer tok_e946329a60ff45ba807f3f41b0e8b7fc"
  22. cls._instance.session = requests.Session()
  23. return cls._instance
  24. @staticmethod
  25. def Instance():
  26. return VSCloudApi()
  27. def _get_headers(self, content_type: str = "application/json") -> Dict[str, str]:
  28. return {
  29. "Authorization": self.api_token,
  30. "Content-Type": content_type,
  31. "Accept": "application/json, text/plain, */*"
  32. }
  33. def _perform_request(self, method, url, headers=None, data=None, json_data=None, params=None, timeout=5*60):
  34. """
  35. 统一 HTTP 请求封装
  36. """
  37. resp = self.session.request(method, url, headers=headers, data=data, json=json_data, params=params, timeout=timeout)
  38. VSC_DEBUG('vs_cloud', f'[perform request] {method} {url} {data} {json_data} {params} {resp.text}')
  39. if resp.status_code == 200:
  40. return resp
  41. elif resp.status_code == 401:
  42. raise SessionExpiredOrInvalidError()
  43. elif resp.status_code == 403:
  44. raise PermissionDeniedError()
  45. elif resp.status_code == 429:
  46. raise RateLimiteddError()
  47. else:
  48. raise BizLogicError(message=f"HTTP Error {resp.status_code}: {resp.text[:100]}")
  49. def get_dynamic_config(self, config_name: str) -> Dict:
  50. url = f'https://api.text.skin/api/dynamic-configurations/key/{quote(config_name)}'
  51. headers = self._get_headers()
  52. resp = self._perform_request('GET', url, headers=headers, timeout=10)
  53. result = resp.json()
  54. if result.get("code") == 0:
  55. data = result.get("data") or {}
  56. return data.get("config_value")
  57. else:
  58. raise BizLogicError(message=f"Get dynamic config biz error: {result.get('message')}")
  59. def get_vas_task_pop(self, routing_key: str) -> Dict:
  60. if configure.TEST_TASK:
  61. return configure.TEST_TASK
  62. url = f"{self.base_url}/api/vas/task/pop"
  63. params = {
  64. "queue_name": routing_key,
  65. }
  66. headers = self._get_headers()
  67. resp = self._perform_request('GET', url, params=params, headers=headers, timeout=10)
  68. result = resp.json()
  69. if result.get("code") == 0:
  70. return result.get("data", {})
  71. else:
  72. raise BizLogicError(message=f"Get vas task pop biz error: {result.get('message')}")
  73. def get_vas_task(self, task_id: str) -> Dict:
  74. url = f"{self.base_url}/api/vas/task/detail"
  75. params = {"task_id": task_id}
  76. headers = self._get_headers()
  77. resp = self._perform_request('GET', url, params=params, headers=headers, timeout=10)
  78. result = resp.json()
  79. if result.get("code") == 0:
  80. return result.get("data", {})
  81. else:
  82. raise BizLogicError(message=f"Get Task={task_id} error: {result.get('message')}")
  83. def update_vas_task(self, task_id: str, update_data: Dict[str, Any]) -> Dict:
  84. """
  85. 更新任务
  86. """
  87. url = f"{self.base_url}/api/vas/task/update"
  88. params = {"id": task_id}
  89. headers = self._get_headers()
  90. resp = self._perform_request('POST', url, params=params, json_data=update_data, headers=headers, timeout=10)
  91. result = resp.json()
  92. if result.get("code") == 0:
  93. return result.get("data", {})
  94. else:
  95. raise BizLogicError(message=f"Update vas task biz error: {result.get('message')}")
  96. def return_vas_task_to_queue(self, task_id: int):
  97. """
  98. 重入队列
  99. """
  100. url = f"{self.base_url}/api/vas/task/return_to_queue"
  101. params = {"task_id": task_id}
  102. headers = self._get_headers()
  103. resp = self._perform_request('POST', url, params=params, headers=headers, timeout=10)
  104. result = resp.json()
  105. if result.get("code") == 0:
  106. return result.get("data", {})
  107. else:
  108. raise BizLogicError(message=f"Return vas task to queue biz error: {result.get('message')}")
  109. def push_weixin_text(self, text:str):
  110. """
  111. 推送微信文本消息
  112. """
  113. url = f"{self.base_url}/api/wechat/send_no_token"
  114. payload = {"message": text}
  115. headers = self._get_headers()
  116. resp = self._perform_request('POST', url, json_data=payload, headers=headers, timeout=10)
  117. result = resp.json()
  118. if result.get("code") == 0:
  119. return result.get("data", {})
  120. else:
  121. raise BizLogicError(message=f"Return vas task to queue biz error: {result.get('message')}")
  122. def create_task(self, command: str, args: Dict) -> str:
  123. """
  124. 创建任务
  125. """
  126. url = f"{self.base_url}/api/tasks"
  127. headers = self._get_headers()
  128. payload = {
  129. "command": command,
  130. "args": args,
  131. "status": 0
  132. }
  133. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  134. result = resp.json()
  135. if result.get("code") == 0:
  136. data = result.get("data", {})
  137. task_id = data.get("id")
  138. if not task_id:
  139. raise BizLogicError(message=f"Task created but no ID returned. Resp: {data}")
  140. return str(task_id)
  141. else:
  142. raise BizLogicError(message=f"Create task failed ({command}): {result.get('message')}")
  143. def get_task_result(self, task_id: str, timeout: int = 120, interval: int = 3) -> Dict:
  144. """
  145. 轮询获取任务结果
  146. """
  147. url = f"{self.base_url}/api/tasks/{task_id}"
  148. headers = self._get_headers()
  149. start_time = time.time()
  150. while True:
  151. if time.time() - start_time > timeout:
  152. raise BizLogicError(message=f"Wait for task result timeout ({timeout}s). TaskID: {task_id}")
  153. try:
  154. resp = self._perform_request('GET', url, headers=headers, timeout=10)
  155. result = resp.json()
  156. if result.get("code") != 0:
  157. raise BizLogicError(message=f"API Error fetching task: {result.get('message')}")
  158. data = result.get("data", {})
  159. status = data.get("status")
  160. if status == 2:
  161. return data
  162. elif status == 3:
  163. error_msg = data.get("result", "Unknown error")
  164. raise BizLogicError(message=f"Task execution failed: {error_msg}")
  165. except Exception as e:
  166. if isinstance(e, BizLogicError):
  167. raise e
  168. VSC_WARN("vs_cloud", f"Polling exception: {str(e)}")
  169. time.sleep(interval)
  170. def create_http_session(
  171. self,
  172. session_id: str,
  173. cookies: Optional[str] = None,
  174. local_storage: Optional[str] = None,
  175. user_agent: Optional[str] = None,
  176. proxy: Optional[str] = None,
  177. page: Optional[str] = None,
  178. session_storage: Optional[str] = None
  179. ) -> Dict:
  180. """
  181. 创建 http session
  182. """
  183. url = f"{self.base_url}/api/http-session"
  184. headers = self._get_headers()
  185. payload = {
  186. "local_storage": local_storage,
  187. "session_storage": session_storage,
  188. "cookies": cookies,
  189. "user_agent": user_agent,
  190. "proxy": proxy,
  191. "page": page,
  192. "session_id": session_id
  193. }
  194. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  195. result = resp.json()
  196. if result.get("code") == 0:
  197. return result.get("data", {})
  198. else:
  199. raise BizLogicError(message=f"Create http session biz error: {result.get('message')}")
  200. def slot_refresh_start(
  201. self,
  202. routing_key: str,
  203. country: str = "",
  204. city: str = "",
  205. visa_type: str = "Tourist",
  206. snapshot_source: str = "worker",
  207. ):
  208. url = f"{self.base_url}/api/slot_refresh/start"
  209. payload = {
  210. "routing_key": routing_key,
  211. "country": country,
  212. "city": city,
  213. "visa_type": visa_type,
  214. "snapshot_source": snapshot_source,
  215. }
  216. headers = self._get_headers()
  217. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  218. result = resp.json()
  219. if result.get("code") == 0:
  220. return result.get("data", {})
  221. else:
  222. raise BizLogicError(message=f"Slot refresh start biz error: {result.get('message')}")
  223. def slot_refresh_success(self, routing_key: str, snapshot_source: str = "worker"):
  224. url = f'{self.base_url}/api/slot_refresh/success'
  225. payload = {
  226. "routing_key": routing_key,
  227. "snapshot_source": snapshot_source,
  228. }
  229. headers = self._get_headers()
  230. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  231. result = resp.json()
  232. if result.get("code") == 0:
  233. return result.get("data", {})
  234. else:
  235. raise BizLogicError(message=f"Slot refresh success biz error: {result.get('message')}")
  236. def slot_refresh_fail(self, routing_key: str, error: str, snapshot_source: str = "worker"):
  237. url = f'{self.base_url}/api/slot_refresh/fail'
  238. payload = {
  239. "routing_key": routing_key,
  240. "snapshot_source": snapshot_source,
  241. "error": error,
  242. }
  243. headers = self._get_headers()
  244. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  245. result = resp.json()
  246. if result.get("code") == 0:
  247. return result.get("data", {})
  248. else:
  249. raise BizLogicError(message=f"Slot refresh fail biz error: {result.get('message')}")
  250. def get_next_account(self, pool_name: str, account_cd: int = 60):
  251. if configure.TEST_ACCOUNT:
  252. return configure.TEST_ACCOUNT
  253. url = f'{self.base_url}/api/account/next'
  254. payload = {
  255. "pool_name": pool_name,
  256. "account_cd": account_cd
  257. }
  258. headers = self._get_headers()
  259. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  260. result = resp.json()
  261. if result.get("code") == 0:
  262. return result.get("data", {})
  263. else:
  264. raise BizLogicError(message=f"Get next account biz error: {result.get('message')}")
  265. def get_next_proxy(self, pools: List[str], proxy_cd: int = 60):
  266. if configure.TEST_PROXY:
  267. return configure.TEST_PROXY
  268. url = f'{self.base_url}/api/proxy/next-ip'
  269. payload = {
  270. "pools": pools,
  271. "proxy_cd": proxy_cd
  272. }
  273. headers = self._get_headers()
  274. resp = self._perform_request('POST', url, headers=headers, json_data=payload, timeout=10)
  275. result = resp.json()
  276. if result.get("code") == 0:
  277. return result.get("data", {})
  278. else:
  279. raise BizLogicError(message=f"Get next account biz error: {result.get('message')}")
  280. def slot_snapshot_report(self, query_payload: Dict[str, Any] = {}):
  281. url = f"{self.base_url}/api/slots/report"
  282. headers = self._get_headers()
  283. resp = self._perform_request("POST", url, headers=headers, json_data=query_payload, timeout=10)
  284. result = resp.json()
  285. if result.get("code") == 0:
  286. return result.get("data", {})
  287. else:
  288. raise BizLogicError(message=f"Slot refresh fail biz error: {result.get('message')}")
  289. def fetch_mail_content(
  290. self,
  291. email: str,
  292. sender: str,
  293. recipient: str,
  294. subject_keywords: str,
  295. body_keywords: str,
  296. sent_date: str,
  297. expiry: int
  298. ) -> Optional[str]:
  299. """
  300. 获取邮件内容
  301. """
  302. params = {
  303. "email": email,
  304. "sender": sender,
  305. "recipient": recipient,
  306. "subjectKeywords": subject_keywords,
  307. "bodyKeywords": body_keywords,
  308. "sentDate": sent_date,
  309. "expiry": str(expiry)
  310. }
  311. url = f"{self.base_url}/api/email-authorizations/fetch"
  312. headers = self._get_headers()
  313. resp = self._perform_request('POST', url, headers=headers, params=params, data="", timeout=5*60)
  314. result = resp.json()
  315. if result.get('code') == 0:
  316. data = result.get('data', {})
  317. return data.get('body', '')
  318. else:
  319. raise BizLogicError(message=f"Fetch mail content biz error: {result.get('message')}")
  320. def fetch_mail_content_from_top(
  321. self,
  322. email: str,
  323. sender: str,
  324. recipient: str,
  325. subject_keywords: str,
  326. body_keywords: str,
  327. top: int
  328. ) -> Optional[str]:
  329. """从顶部获取邮件内容"""
  330. params = {
  331. "email": email,
  332. "sender": sender,
  333. "recipient": recipient,
  334. "subjectKeywords": subject_keywords,
  335. "bodyKeywords": body_keywords,
  336. "top": str(top)
  337. }
  338. url = f"{self.base_url}/api/email-authorizations/fetch-top"
  339. headers = self._get_headers()
  340. resp = self._perform_request('POST', url, headers=headers, params=params, data="", timeout=5*60)
  341. result = resp.json()
  342. if result.get('code') == 0:
  343. data = result.get('data', {})
  344. return data.get('body', "")
  345. else:
  346. raise BizLogicError(message=f"Fetch mail content from top biz error: {result.get('message')}")