booker_builtin.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313
  1. import os
  2. import time
  3. import json
  4. import threading
  5. import random
  6. from typing import List, Dict, Callable
  7. from vs_types import GroupConfig, VSPlgConfig, Task, VSQueryResult, AppointmentType
  8. from vs_plg_factory import VSPlgFactory
  9. from toolkit.thread_pool import ThreadPool
  10. from toolkit.vs_cloud_api import VSCloudApi
  11. from utils.safe_redis_cli import SafeRedisClient
  12. class BuiltinBookerGCO:
  13. """
  14. 非绑定模式 (公共内置账号池):
  15. - 只维护全局 target_instances 数量的实例。
  16. - 所有实例热机等待,发现信号后临时去云端 Pop 订单。
  17. """
  18. def __init__(self, cfg: GroupConfig, redis_conf: Dict, logger: Callable[[str], None] = None):
  19. self.m_cfg = cfg
  20. self.m_factory = VSPlgFactory()
  21. self.m_logger = logger
  22. self.m_tasks: List[Task] = []
  23. self.m_lock = threading.RLock()
  24. self.m_stop_event = threading.Event()
  25. self.redis_client = SafeRedisClient(redis_conf, self.m_logger)
  26. self.m_tracker_key = f"vs:worker:tasks_tracker:{self.m_cfg.identifier}"
  27. def _log(self, message):
  28. if self.m_logger:
  29. self.m_logger(f'[BUILTIN-BOOKER] [{self.m_cfg.identifier}] {message}')
  30. else:
  31. print(f'[BUILTIN-BOOKER] [{self.m_cfg.identifier}] {message}')
  32. def start(self):
  33. if not self.m_cfg.enable:
  34. return
  35. self._log("Starting Built-in Booker...")
  36. plugin_name = self.m_cfg.plugin_config.plugin_name
  37. class_name = "".join(part.title() for part in plugin_name.split('_'))
  38. plugin_path = os.path.join(self.m_cfg.plugin_config.lib_path, self.m_cfg.plugin_config.plugin_bin)
  39. self.m_factory.register_plugin(plugin_name, plugin_path, class_name)
  40. threading.Thread(target=self._booking_trigger_loop, daemon=True).start()
  41. threading.Thread(target=self._creator_loop, daemon=True).start()
  42. threading.Thread(target=self._maintain_loop, daemon=True).start()
  43. def stop(self):
  44. self._log("Stopping Booker...")
  45. self.m_stop_event.set()
  46. self._cleanup_all_tasks("booker stop")
  47. def update_config(self, new_cfg: GroupConfig):
  48. """
  49. 动态更新配置
  50. """
  51. with self.m_lock:
  52. if self.m_cfg.enable and not new_cfg.enable:
  53. self._log("Config dynamically updated: Group DISABLED. Will stop creating new tasks.")
  54. elif not self.m_cfg.enable and new_cfg.enable:
  55. self._log("Config dynamically updated: Group ENABLED.")
  56. else:
  57. self._log("Config dynamically updated: Parameters refreshed.")
  58. self.m_cfg = new_cfg
  59. def _cleanup_task(self, task: Task, reason: str = ""):
  60. try:
  61. if task and task.instance:
  62. task.instance.cleanup()
  63. self._log(f"🧹 Cleaned up built-in instance. Reason: {reason}")
  64. except Exception as e:
  65. self._log(f"Cleanup failed for built-in instance. Reason: {reason}. Error: {e}")
  66. def _remove_task(self, task: Task, reason: str = "", cleanup: bool = True):
  67. removed = False
  68. with self.m_lock:
  69. if task in self.m_tasks:
  70. self.m_tasks.remove(task)
  71. removed = True
  72. if cleanup and removed:
  73. self._cleanup_task(task, reason)
  74. return removed
  75. def _cleanup_all_tasks(self, reason: str = ""):
  76. with self.m_lock:
  77. tasks = list(self.m_tasks)
  78. self.m_tasks.clear()
  79. for task in tasks:
  80. self._cleanup_task(task, reason)
  81. def _get_redis_key(self, routing_key: str) -> str:
  82. return f"vs:signal:{routing_key}"
  83. def _maintain_loop(self):
  84. self._log("Maintain loop started.")
  85. while not self.m_stop_event.is_set():
  86. try:
  87. time.sleep(1.0)
  88. with self.m_lock:
  89. tasks_to_check = list(self.m_tasks)
  90. if not tasks_to_check:
  91. continue
  92. dead_tasks = []
  93. healthy_tasks = []
  94. now = time.time()
  95. for t in tasks_to_check:
  96. if now >= t.next_remote_ping:
  97. t.instance.keep_alive()
  98. if t.instance.health_check():
  99. healthy_tasks.append(t)
  100. t.next_remote_ping = now + random.gauss(self.m_cfg.booker.keep_alive, 5)
  101. else:
  102. dead_tasks.append(t)
  103. else:
  104. healthy_tasks.append(t)
  105. if dead_tasks:
  106. with self.m_lock:
  107. current_tasks = list(self.m_tasks)
  108. self.m_tasks = [t for t in self.m_tasks if t in healthy_tasks]
  109. for t in dead_tasks:
  110. if t in current_tasks:
  111. self._cleanup_task(t, "unhealthy or keep-alive failed")
  112. else:
  113. with self.m_lock:
  114. self.m_tasks = [t for t in self.m_tasks if t in healthy_tasks]
  115. except Exception as e:
  116. self._log(f'Maintain loop exception: {e}')
  117. def _booking_trigger_loop(self):
  118. self._log("Trigger loop started.")
  119. while not self.m_stop_event.is_set():
  120. try:
  121. time.sleep(0.1)
  122. now = time.time()
  123. for apt_type in self.m_cfg.appointment_types:
  124. redis_key = self._get_redis_key(apt_type.routing_key)
  125. raw_data = self.redis_client.get(redis_key)
  126. if not raw_data:
  127. continue
  128. try:
  129. data = json.loads(raw_data)
  130. query_result = VSQueryResult.model_validate(data['query_result'])
  131. query_result.apt_type = AppointmentType.model_validate(data['apt_type'])
  132. except Exception as parse_err:
  133. self._log(f"Data parsing error for {redis_key}: {parse_err}. Deleting corrupted signal.")
  134. self.redis_client.delete(redis_key)
  135. continue
  136. matching_tasks = []
  137. with self.m_lock:
  138. for task in self.m_tasks:
  139. if now < task.next_run or not task.book_allowed:
  140. continue
  141. if apt_type.routing_key not in task.acceptable_routing_keys:
  142. continue
  143. task.next_run = now + self.m_cfg.booker.booking_cooldown
  144. matching_tasks.append(task)
  145. if matching_tasks:
  146. threads = []
  147. for task in matching_tasks:
  148. self._log(f"🚀 Triggering BOOK for {apt_type.routing_key} | Order Ref: {task.task_ref}")
  149. t = threading.Thread(target=self._execute_book_job, args=(task, query_result))
  150. threads.append(t)
  151. t.start()
  152. for t in threads:
  153. t.join()
  154. except Exception as e:
  155. self._log(f"Booking trigger loop exception: {e}")
  156. def _execute_book_job(self, task: Task, query_result: VSQueryResult):
  157. queue_name = f"auto.{query_result.apt_type.routing_key}"
  158. task_id = None
  159. task_data = None
  160. booking_success = False
  161. try:
  162. task_data = VSCloudApi.Instance().get_vas_task_pop(queue_name)
  163. if not task_data:
  164. return
  165. task_id = task_data['id']
  166. order_id = task_data.get('order_id')
  167. self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + 30.0})
  168. user_input = task_data.get('user_inputs', {})
  169. book_res = task.instance.book(query_result, user_input)
  170. if book_res.success:
  171. booking_success = True
  172. self._log(f"✅ BOOK SUCCESS! Order: {order_id}")
  173. grab_info = {
  174. "account": book_res.account,
  175. "session_id": book_res.session_id,
  176. "urn": book_res.urn,
  177. "slot_date": book_res.book_date,
  178. "slot_time": book_res.book_time,
  179. "timestamp": int(time.time()),
  180. "payment_link": book_res.payment_link
  181. }
  182. def _update_cloud_success():
  183. try:
  184. VSCloudApi.Instance().update_vas_task(str(task_id), {"status": "grabbed", "grabbed_history": grab_info})
  185. push_content = (
  186. f"🎉 【预定成功通知】\n"
  187. f"━━━━━━━━━━━━━━━\n"
  188. f"订单编号: {order_id}\n"
  189. f"预约账号: {book_res.account}\n"
  190. f"预约日期: {book_res.book_date}\n"
  191. f"预约时间: {book_res.book_time}\n"
  192. f"预约编号: {book_res.urn}\n"
  193. f"支付链接: {book_res.payment_link if book_res.payment_link else '无需支付/暂无'}\n"
  194. f"━━━━━━━━━━━━━━━\n"
  195. )
  196. VSCloudApi.Instance().push_weixin_text(push_content)
  197. except Exception as e:
  198. self._log(f"Failed to update success state to cloud: {e}")
  199. ThreadPool.getInstance().enqueue(_update_cloud_success)
  200. self.redis_client.zrem(self.m_tracker_key, task_id)
  201. task.successful_bookings += 1
  202. max_b = self.m_cfg.booker.max_bookings_per_account
  203. if max_b > 0 and task.successful_bookings >= max_b:
  204. self._log(f"Account reached max bookings ({max_b}). Destroying instance.")
  205. self._remove_task(task, "max bookings reached")
  206. else:
  207. self._log(f"❌ BOOK FAILED for Order: {order_id}")
  208. except Exception as e:
  209. err_str = str(e)
  210. self._log(f"Exception during booking: {err_str}")
  211. rate_limited_indicators = [
  212. "42901" in err_str,
  213. "Rate limited" in err_str
  214. ]
  215. if any(rate_limited_indicators):
  216. self._remove_task(task, "booking rate limited")
  217. def _creator_loop(self):
  218. self._log("Creator loop started.")
  219. while not self.m_stop_event.is_set():
  220. try:
  221. time.sleep(1.0)
  222. with self.m_lock:
  223. current = len(self.m_tasks)
  224. target = self.m_cfg.booker.target_instances
  225. if current < target:
  226. self._spawn_worker()
  227. except Exception as e:
  228. self._log(f'Creator loop exception: {e}')
  229. def _spawn_worker(self):
  230. instance = None
  231. success = False
  232. plg_cfg = None
  233. try:
  234. plg_cfg = VSPlgConfig()
  235. plg_cfg.debug = self.m_cfg.debug
  236. plg_cfg.free_config = self.m_cfg.free_config
  237. plg_cfg.session_max_life = self.m_cfg.session_max_life
  238. if self.m_cfg.need_account:
  239. acc = VSCloudApi.Instance().get_next_account(self.m_cfg.booker.account_pool_id, self.m_cfg.booker.account_cd)
  240. plg_cfg.account = type(plg_cfg.account)(**acc)
  241. if self.m_cfg.need_proxy:
  242. proxy = VSCloudApi.Instance().get_next_proxy(self.m_cfg.proxy_pool, self.m_cfg.proxy_cd)
  243. plg_cfg.proxy = type(plg_cfg.proxy)(**proxy)
  244. instance = self.m_factory.create(self.m_cfg.identifier, self.m_cfg.plugin_config.plugin_name)
  245. instance.set_log(self.m_logger)
  246. instance.set_config(plg_cfg)
  247. instance.create_session()
  248. with self.m_lock:
  249. all_keys = [apt.routing_key for apt in self.m_cfg.appointment_types]
  250. self.m_tasks.append(
  251. Task(
  252. instance=instance,
  253. next_run=time.time(),
  254. task_ref=None,
  255. acceptable_routing_keys=all_keys,
  256. source_queue="built-in",
  257. book_allowed=True,
  258. next_remote_ping = time.time() + random.gauss(self.m_cfg.booker.keep_alive, 5)
  259. )
  260. )
  261. self._log(f"+++ Built-in Booker spawned: {plg_cfg.account.username}")
  262. except Exception as e:
  263. err_str = str(e)
  264. self._log(f"Spawn failed: {err_str}")
  265. rate_limited_indicators = [
  266. "42901" in err_str,
  267. "Rate limited" in err_str
  268. ]
  269. if any(rate_limited_indicators):
  270. if plg_cfg and plg_cfg.account.username != "Guest":
  271. VSCloudApi.lock_account(plg_cfg.account.id, self.m_cfg.login_backoff)
  272. finally:
  273. if not success:
  274. if instance:
  275. instance.cleanup()