jhs_rpc_spider.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330
  1. # -*- coding: utf-8 -*-
  2. # Author : Charley
  3. # Python : 3.10.8
  4. # Date : 2026/4/23 13:46
  5. import json
  6. import time
  7. import requests
  8. import inspect
  9. import schedule
  10. from loguru import logger
  11. from typing import Any, Dict
  12. from datetime import datetime
  13. from mysql_pool import MySQLConnectionPool
  14. from jhs_raw_codec_client import JhsRawCodecClient
  15. from tenacity import retry, stop_after_attempt, wait_fixed, retry_if_exception_type
  16. """
  17. 此项目基于 集换社3.36.2 版本 其他版本报错
  18. [2026-05-22 19:29:05.043] ERROR Error fetching page 9: [CODEC-ERROR]call failed: Error: unable to resolve instance for gc.b
  19. [2026-05-22 19:38:41.595] ERROR Error fetching page 4: [CODEC-ERROR]call failed: TypeError: not a function
  20. """
  21. class TokenExpiredError(Exception):
  22. """token 过期 / 未授权时抛出。
  23. 服务端会返回 HTTP 200,但 body 为 {"code":401,"error":"MARKET_UNAUTHORIZED",...},
  24. 因此 raise_for_status 无法拦截,需在解析 body 时主动识别。该异常故意不加入
  25. fetch_market_page 的 @retry 可重试类型——token 过期重试无意义,应立即中止整轮。
  26. """
  27. # TOKEN = "eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJlbnYiOiJwcm9kdWN0aW9uIiwic3ViIjoyODI3NDU4LCJpc3MiOiJodHRwOi8vYXBpLmppaHVhbnNoZS5jb20vYXBpL21hcmtldC9hdXRoL2xvZ2luLW9yLXNpZ251cCIsImlhdCI6MTc3NTYzNzQzNSwiZXhwIjoxNzgwODIxNDM1LCJuYmYiOjE3NzU2Mzc0MzUsImp0aSI6InhiT3NsdUJRTzVWeHRabHQifQ.uHz7M-U0ewPgi5Qzr5P4eJbSdIUO_i_hmVE-0jsaG2Y"
  28. # DEVICE_ID = "127.0.0.1:5557" # adb connect 127.0.0.1:5557
  29. DEVICE_ID = "25051FDD4S018P" # adb connect 127.0.0.1:5557
  30. CLI_TARGET_SEC = 2
  31. TIMEOUT_SEC = 15
  32. BASE_URL = "https://api.jihuanshe.com/api/market/auction-products"
  33. HEADERS = {
  34. "User-Agent": "Model/google,Pixel5 OS/30 Version/3.36.2",
  35. "Connection": "Keep-Alive",
  36. "Accept-Encoding": "gzip",
  37. "x-device-id": "6efe93931488e176",
  38. }
  39. logger.remove()
  40. logger.add("./logs/{time:YYYYMMDD}.log", encoding='utf-8', rotation="00:00",
  41. format="[{time:YYYY-MM-DD HH:mm:ss.SSS}] {level} {message}",
  42. level="DEBUG", retention="7 day")
  43. def after_log(retry_state):
  44. """
  45. retry 回调
  46. :param retry_state: RetryCallState 对象
  47. """
  48. # 检查 args 是否存在且不为空
  49. if retry_state.args and len(retry_state.args) > 0:
  50. log = retry_state.args[0] # 获取传入的 logger
  51. else:
  52. log = logger # 使用全局 logger
  53. if retry_state.outcome.failed:
  54. log.warning(
  55. f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} Times")
  56. else:
  57. log.info(f"Function '{retry_state.fn.__name__}', Attempt {retry_state.attempt_number} succeeded")
  58. @retry(stop=stop_after_attempt(5), wait=wait_fixed(1), after=after_log)
  59. def get_proxys(log):
  60. """
  61. 获取代理
  62. :return: 代理
  63. """
  64. tunnel = "x371.kdltps.com:15818"
  65. kdl_username = "t13753103189895"
  66. kdl_password = "o0yefv6z"
  67. try:
  68. proxies = {
  69. "http": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel},
  70. "https": "http://%(user)s:%(pwd)s@%(proxy)s/" % {"user": kdl_username, "pwd": kdl_password, "proxy": tunnel}
  71. }
  72. return proxies
  73. except Exception as e:
  74. log.error(f"Error getting proxy: {e}")
  75. raise e
  76. # 覆盖到所有可重试异常:
  77. # - json.JSONDecodeError: 响应体损坏
  78. # - TimeoutError: frida RPC 或 CLI 兜底超时(新增的超时保护路径)
  79. # - RuntimeError: [CODEC-ERROR] 等 JS 侧抛出的错误
  80. # - requests.RequestException: HTTP 层网络异常(连接超时 / 读超时 / 5xx 抛出等)
  81. @retry(stop=stop_after_attempt(3), wait=wait_fixed(2),
  82. retry=retry_if_exception_type(
  83. (json.JSONDecodeError, TimeoutError, RuntimeError, requests.RequestException)
  84. ),
  85. after=after_log)
  86. def fetch_market_page(
  87. log,
  88. page: int,
  89. token: str,
  90. client: JhsRawCodecClient,
  91. session: requests.Session,
  92. headers: Dict[str, str],
  93. timeout_sec: int = TIMEOUT_SEC,
  94. ) -> Dict[str, Any]:
  95. """
  96. 请求并解密单页数据。
  97. 复用方式:
  98. - `client` 和 `session` 由外层创建一次并长期复用
  99. - 调用本函数时只传不同 page 即可
  100. """
  101. log.info(f"Fetching page {page}......................")
  102. url_for_enc = f"{BASE_URL}?sorting=completed&page={page}&token={token}"
  103. enc = client.call({"op": "enc", "url": url_for_enc})
  104. raw_data = enc["raw_data"]
  105. resp = session.get(
  106. BASE_URL,
  107. headers=headers,
  108. params={"raw_data": raw_data, "token": token},
  109. timeout=timeout_sec,
  110. )
  111. resp.raise_for_status()
  112. body = resp.json()
  113. # 服务端 token 过期时返回 HTTP 200 但 body 无 raw_data({"code":401,"error":"MARKET_UNAUTHORIZED"}),
  114. # 需主动识别,否则会退化成没头没脑的 KeyError: 'raw_data' 并空跑满全部页
  115. if "raw_data" not in body:
  116. code = body.get("code")
  117. err = body.get("error")
  118. msg = body.get("msg")
  119. if code == 401 or err == "MARKET_UNAUTHORIZED":
  120. raise TokenExpiredError(f"token 已过期/未授权: code={code} error={err} msg={msg}")
  121. # 其他缺字段情况:抛 RuntimeError 触发 @retry(可能是偶发脏响应)
  122. raise RuntimeError(f"响应缺少 raw_data 字段: {json.dumps(body, ensure_ascii=False)[:200]}")
  123. response_raw_data = body["raw_data"]
  124. request_url_for_dec = f"{BASE_URL}?raw_data={raw_data}&token={token}"
  125. dec = client.call(
  126. {
  127. "op": "dec",
  128. "request_url": request_url_for_dec,
  129. "response_raw_data": response_raw_data,
  130. }
  131. )
  132. response_body = dec.get("response_body", "")
  133. parsed: Any = response_body
  134. if isinstance(response_body, str):
  135. try:
  136. parsed = json.loads(response_body)
  137. except Exception:
  138. log.error(f"Error parsing response body: {response_body}")
  139. pass
  140. return {
  141. "page": page,
  142. "enc": enc,
  143. "http_json": body,
  144. "dec": dec,
  145. "decoded": parsed,
  146. }
  147. def parse_data(resp_data, sql_pool):
  148. """
  149. 解析数据
  150. :param resp_data: 响应数据
  151. :param sql_pool: 数据库连接池
  152. """
  153. data_list = resp_data.get("raw_data", {}).get("data", [])
  154. info_list = []
  155. for data in data_list:
  156. seller_username = data.get("seller_username")
  157. product_id = data.get("auction_product_id")
  158. app_id = data.get("app_id")
  159. auction_product_name = data.get("auction_product_name")
  160. auction_product_images = data.get("auction_product_image")
  161. game_key = data.get("game_key")
  162. language_text = data.get("language_text")
  163. authenticator_name = data.get("authenticator_name")
  164. grading = data.get("grading")
  165. starting_price = data.get("starting_price")
  166. max_bid_price = data.get("max_bid_price")
  167. status = data.get("status")
  168. auction_product_start_timestamp = data.get('auction_product_start_timestamp')
  169. auction_product_start_time = datetime.fromtimestamp(auction_product_start_timestamp).strftime(
  170. '%Y-%m-%d %H:%M:%S') if auction_product_start_timestamp else None
  171. auction_product_end_timestamp = data.get('auction_product_end_timestamp')
  172. auction_product_end_time = datetime.fromtimestamp(auction_product_end_timestamp).strftime(
  173. '%Y-%m-%d %H:%M:%S') if auction_product_end_timestamp else None
  174. bid_count = data.get("bid_count")
  175. card_number = data.get("number")
  176. rarity = data.get("rarity")
  177. data_dict = {
  178. "seller_username": seller_username,
  179. "product_id": product_id,
  180. "app_id": app_id,
  181. "auction_product_name": auction_product_name,
  182. "auction_product_images": auction_product_images,
  183. "game_key": game_key,
  184. "language_text": language_text,
  185. "authenticator_name": authenticator_name,
  186. "grading": grading,
  187. "starting_price": starting_price,
  188. "max_bid_price": max_bid_price,
  189. "status": status,
  190. "auction_product_start_time": auction_product_start_time,
  191. "auction_product_end_time": auction_product_end_time,
  192. "bid_count": bid_count,
  193. "card_number": card_number,
  194. "rarity": rarity,
  195. }
  196. # print(data_dict)
  197. # print(type(data))
  198. info_list.append(data_dict)
  199. if info_list:
  200. sql_pool.insert_many(table="jhs_product_record", data_list=info_list, ignore=True)
  201. def get_market_list(log, token: str, sql_pool):
  202. """
  203. 分页抓取市场列表并入库。
  204. 支持连续失败自动重连 frida 会话——手机端 app 冷启动 / gadget 假死时,
  205. 避免整轮任务卡死或全页失败。
  206. Args:
  207. log: 日志对象。
  208. token (str): 用户登录 token。
  209. sql_pool: 数据库连接池实例。
  210. """
  211. page = 1
  212. max_page = 200
  213. # 连续失败达到阈值就重连 frida 会话(覆盖 app 冷启动类未加载 / gadget 掉线场景)
  214. consecutive_fails = 0
  215. reconnect_threshold = 3
  216. with JhsRawCodecClient(device_id=DEVICE_ID, cli_target_sec=CLI_TARGET_SEC) as codec_client:
  217. with requests.Session() as http_sess:
  218. while page < max_page:
  219. try:
  220. result = fetch_market_page(
  221. log=log,
  222. page=page,
  223. token=token,
  224. client=codec_client,
  225. session=http_sess,
  226. headers=HEADERS,
  227. )
  228. # print(page, result["decoded"])
  229. try:
  230. parse_data(result["decoded"], sql_pool)
  231. except Exception as e:
  232. log.error(f"Error parsing page {page}: {e}")
  233. consecutive_fails = 0
  234. except TokenExpiredError as e:
  235. # token 过期:重连 frida / 翻页都没用,立即中止整轮,等更新 token 后重跑
  236. log.error(f"token 已过期,请更新 jhs_token 表(id=1)后重跑,本轮中止: {e}")
  237. break
  238. except Exception as e:
  239. log.error(f"Error fetching page {page}: {e}")
  240. consecutive_fails += 1
  241. if consecutive_fails >= reconnect_threshold:
  242. log.warning(
  243. f"连续 {consecutive_fails} 页失败,尝试 reconnect frida 会话")
  244. try:
  245. codec_client.reconnect()
  246. consecutive_fails = 0
  247. except Exception as re:
  248. log.error(f"reconnect failed: {re}")
  249. page += 1
  250. @retry(stop=stop_after_attempt(100), wait=wait_fixed(3600), after=after_log)
  251. def jhs_rpc_main(log):
  252. """
  253. 主函数
  254. :param log: logger对象
  255. """
  256. log.info(
  257. f'开始运行 {inspect.currentframe().f_code.co_name} 爬虫任务....................................................')
  258. # 配置 MySQL 连接池
  259. sql_pool = MySQLConnectionPool(log=log)
  260. if not sql_pool:
  261. log.error("MySQL数据库连接失败")
  262. raise Exception("MySQL数据库连接失败")
  263. try:
  264. jhs_token = sql_pool.select_one('SELECT token FROM jhs_token WHERE id = 1')
  265. get_market_list(log, jhs_token[0], sql_pool)
  266. except Exception as e:
  267. log.error(f'{inspect.currentframe().f_code.co_name} error: {e}')
  268. finally:
  269. log.info(f'爬虫程序 {inspect.currentframe().f_code.co_name} 运行结束,等待下一轮的采集任务............')
  270. def schedule_task():
  271. """
  272. 设置定时任务
  273. """
  274. # jhs_rpc_main(log=logger)
  275. schedule.every().day.at("05:00").do(jhs_rpc_main, log=logger)
  276. while True:
  277. schedule.run_pending()
  278. time.sleep(1)
  279. if __name__ == "__main__":
  280. schedule_task()
  281. # jhs_rpc_main(log=logger)