booker_builtin.py 14 KB

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