booker_order.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390
  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, AvailabilityStatus
  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 OrderBookerGCO:
  13. """
  14. 绑定模式 (订单自带账号):
  15. - 按城市队列维护热机配额。
  16. - 绝对的 1 对 1 关系:一个实例绑定一个云端订单。
  17. - 预订成功后,实例立即销毁。
  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_task_data_cache: Dict[str, dict] = {}
  28. self.m_tracker_key = f"vs:worker:tasks_tracker:{self.m_cfg.identifier}"
  29. self.heartbeat_ttl = 2*60.0
  30. def _log(self, message):
  31. if self.m_logger:
  32. self.m_logger(f'[ORDER-BOOKER] [{self.m_cfg.identifier}] {message}')
  33. else:
  34. print(f'[ORDER-BOOKER] [{self.m_cfg.identifier}] {message}')
  35. def start(self):
  36. if not self.m_cfg.enable:
  37. return
  38. self._log("Starting Order Booker...")
  39. plugin_name = self.m_cfg.plugin_config.plugin_name
  40. class_name = "".join(part.title() for part in plugin_name.split('_'))
  41. plugin_path = os.path.join(self.m_cfg.plugin_config.lib_path, self.m_cfg.plugin_config.plugin_bin)
  42. self.m_factory.register_plugin(plugin_name, plugin_path, class_name)
  43. threading.Thread(target=self._booking_trigger_loop, daemon=True).start()
  44. threading.Thread(target=self._creator_loop, daemon=True).start()
  45. threading.Thread(target=self._maintain_loop, daemon=True).start()
  46. threading.Thread(target=self._cache_refresh_loop, daemon=True).start()
  47. def stop(self):
  48. self._log("Stopping Booker...")
  49. self.m_stop_event.set()
  50. self._cleanup_all_tasks("booker stop")
  51. def update_config(self, new_cfg: GroupConfig):
  52. """
  53. 动态更新配置
  54. """
  55. with self.m_lock:
  56. if self.m_cfg.enable and not new_cfg.enable:
  57. self._log("Config dynamically updated: Group DISABLED. Will stop creating new tasks.")
  58. elif not self.m_cfg.enable and new_cfg.enable:
  59. self._log("Config dynamically updated: Group ENABLED.")
  60. else:
  61. self._log("Config dynamically updated: Parameters refreshed.")
  62. self.m_cfg = new_cfg
  63. def _cleanup_task(self, task: Task, reason: str = ""):
  64. try:
  65. if task and task.instance:
  66. task.instance.cleanup()
  67. self._log(f"🧹 Cleaned up instance for task={task.task_ref}. Reason: {reason}")
  68. except Exception as e:
  69. self._log(f"Cleanup failed for task={task.task_ref}. Reason: {reason}. Error: {e}")
  70. def _remove_task(self, task: Task, reason: str = "", cleanup: bool = True):
  71. removed = False
  72. with self.m_lock:
  73. if task in self.m_tasks:
  74. self.m_tasks.remove(task)
  75. removed = True
  76. task_id = task.task_ref
  77. self.m_task_data_cache.pop(task_id, None)
  78. if cleanup and removed:
  79. self._cleanup_task(task, reason)
  80. return removed
  81. def _cleanup_all_tasks(self, reason: str = ""):
  82. with self.m_lock:
  83. tasks = list(self.m_tasks)
  84. self.m_tasks.clear()
  85. self.m_task_data_cache.clear()
  86. for task in tasks:
  87. self._cleanup_task(task, reason)
  88. def _get_redis_key(self, routing_key: str) -> str:
  89. return f"vs:signal:{routing_key}"
  90. def _maintain_loop(self):
  91. self._log("Maintain loop started.")
  92. while not self.m_stop_event.is_set():
  93. try:
  94. time.sleep(1)
  95. with self.m_lock:
  96. tasks_to_check = list(self.m_tasks)
  97. if not tasks_to_check:
  98. continue
  99. healthy_tasks = []
  100. dead_tasks = []
  101. now = time.time()
  102. for t in tasks_to_check:
  103. if now >= t.next_remote_ping:
  104. t.instance.keep_alive()
  105. if t.instance.health_check():
  106. healthy_tasks.append(t)
  107. t.next_remote_ping = now + random.gauss(self.m_cfg.booker.keep_alive, 5)
  108. else:
  109. dead_tasks.append(t)
  110. else:
  111. healthy_tasks.append(t)
  112. if healthy_tasks:
  113. new_deadline = time.time() + self.heartbeat_ttl
  114. mapping = {str(t.task_ref): new_deadline for t in healthy_tasks}
  115. self.redis_client.bulk_zadd(self.m_tracker_key, mapping)
  116. if dead_tasks:
  117. mapping = {str(t.task_ref): 0 for t in dead_tasks}
  118. self.redis_client.bulk_zadd(self.m_tracker_key, mapping)
  119. if dead_tasks:
  120. with self.m_lock:
  121. current_tasks = list(self.m_tasks)
  122. self.m_tasks = [t for t in self.m_tasks if t in healthy_tasks]
  123. for t in dead_tasks:
  124. if t in current_tasks:
  125. self._cleanup_task(t, "unhealthy or keep-alive failed")
  126. else:
  127. with self.m_lock:
  128. self.m_tasks = [t for t in self.m_tasks if t in healthy_tasks]
  129. except Exception as e:
  130. self._log(f'Maintain loop exception:{e}')
  131. def _cache_refresh_loop(self):
  132. self._log("Cache refresh loop started.")
  133. refresh_interval = 15 * 60
  134. while not self.m_stop_event.is_set():
  135. try:
  136. time.sleep(1)
  137. with self.m_lock:
  138. tasks_to_check = {
  139. tid: data.get('_last_refresh', 0)
  140. for tid, data in self.m_task_data_cache.items()
  141. }
  142. if not tasks_to_check:
  143. continue
  144. now = time.time()
  145. for tid, last_refresh in tasks_to_check.items():
  146. if now - last_refresh >= refresh_interval:
  147. fresh_data = VSCloudApi.Instance().get_vas_task(tid)
  148. if fresh_data:
  149. fresh_data['_last_refresh'] = time.time()
  150. with self.m_lock:
  151. if tid in self.m_task_data_cache:
  152. self.m_task_data_cache[tid] = fresh_data
  153. time.sleep(0.5)
  154. except Exception as e:
  155. self._log(f'Cache refresh loop exception: {e}')
  156. def _is_date_of_interest(self, task, query_result: VSQueryResult) -> bool:
  157. """
  158. 判断 query_result 中的可用日期,是否在 task 的意向日期范围内。
  159. """
  160. if query_result.availability_status != AvailabilityStatus.Available:
  161. return True
  162. task_id = task.task_ref
  163. task_data = self.m_task_data_cache.get(str(task_id), {})
  164. user_input = task_data.get('user_inputs', {})
  165. expected_end_date = (
  166. user_input.get('expected_end_date')
  167. or '2100-01-01'
  168. )
  169. available_date = query_result.earliest_date
  170. dt = available_date.strftime("%Y-%m-%d")
  171. return dt <= expected_end_date
  172. def _booking_trigger_loop(self):
  173. self._log("Trigger loop started.")
  174. while not self.m_stop_event.is_set():
  175. try:
  176. time.sleep(0.1)
  177. now = time.time()
  178. for apt_type in self.m_cfg.appointment_types:
  179. redis_key = self._get_redis_key(apt_type.routing_key)
  180. raw_data = self.redis_client.get(redis_key)
  181. if not raw_data:
  182. continue
  183. try:
  184. data = json.loads(raw_data)
  185. query_result = VSQueryResult.model_validate(data['query_result'])
  186. query_result.apt_type = AppointmentType.model_validate(data['apt_type'])
  187. except Exception as parse_err:
  188. self._log(f"Data parsing error for {redis_key}: {parse_err}. Deleting corrupted signal.")
  189. self.redis_client.delete(redis_key)
  190. continue
  191. matching_tasks = []
  192. with self.m_lock:
  193. for task in self.m_tasks:
  194. if now < task.next_run or not task.book_allowed:
  195. continue
  196. if apt_type.routing_key not in task.acceptable_routing_keys:
  197. continue
  198. if not self._is_date_of_interest(task, query_result):
  199. continue
  200. task.next_run = now + self.m_cfg.booker.booking_cooldown
  201. matching_tasks.append(task)
  202. if matching_tasks:
  203. threads = []
  204. for task in matching_tasks:
  205. self._log(f"🚀 Triggering BOOK for {apt_type.routing_key} | Order Ref: {task.task_ref}")
  206. t = threading.Thread(target=self._execute_book_job, args=(task, query_result))
  207. threads.append(t)
  208. t.start()
  209. for t in threads:
  210. t.join()
  211. except Exception as e:
  212. self._log(f"Booking trigger loop exception: {e}")
  213. def _execute_book_job(self, task: Task, query_result: VSQueryResult):
  214. task_id = task.task_ref
  215. task_data = None
  216. try:
  217. with self.m_lock:
  218. task_data = self.m_task_data_cache.get(str(task_id))
  219. if not task_data or task_data.get('status') in ['grabbed', 'pause', 'completed', 'cancelled']:
  220. self._log(f"Bound Task={task_id} is no longer valid or already processed. Removing instance.")
  221. self._remove_task(task, "bound task no longer valid")
  222. self.redis_client.zrem(self.m_tracker_key, task_id)
  223. order_id = task_data.get('order_id')
  224. user_input = task_data.get('user_inputs', {})
  225. book_res = task.instance.book(query_result, user_input)
  226. if book_res.success:
  227. self._log(f"✅ BOOK SUCCESS! Order: {order_id}. Destroying instance.")
  228. grab_info = {
  229. "account": book_res.account,
  230. "session_id": book_res.session_id,
  231. "urn": book_res.urn,
  232. "slot_date": book_res.book_date,
  233. "slot_time": book_res.book_time,
  234. "timestamp": int(time.time()),
  235. "payment_link": book_res.payment_link
  236. }
  237. def _update_cloud_success():
  238. try:
  239. VSCloudApi.Instance().update_vas_task(str(task_id), {"status": "grabbed", "grabbed_history": grab_info})
  240. push_content = (
  241. f"🎉 【预定成功通知】\n"
  242. f"━━━━━━━━━━━━━━━\n"
  243. f"订单编号: {order_id}\n"
  244. f"预约账号: {book_res.account}\n"
  245. f"预约日期: {book_res.book_date}\n"
  246. f"预约时间: {book_res.book_time}\n"
  247. f"预约编号: {book_res.urn}\n"
  248. f"支付链接: {book_res.payment_link if book_res.payment_link else '无需支付/暂无'}\n"
  249. f"━━━━━━━━━━━━━━━\n"
  250. )
  251. VSCloudApi.Instance().push_weixin_text(push_content)
  252. except Exception as e:
  253. self._log(f"Failed to update success state to cloud: {e}")
  254. ThreadPool.getInstance().enqueue(_update_cloud_success)
  255. self.redis_client.zrem(self.m_tracker_key, task_id)
  256. self._remove_task(task, "booking success")
  257. else:
  258. self._log(f"❌ BOOK FAILED for Order: {order_id}. Will retry on next signal.")
  259. except Exception as e:
  260. err_str = str(e)
  261. self._log(f"Exception during booking: {err_str}")
  262. rate_limited_indicators = [
  263. "42901" in err_str,
  264. "Rate limited" in err_str
  265. ]
  266. if any(rate_limited_indicators):
  267. self._remove_task(task, "booking rate limited")
  268. def _creator_loop(self):
  269. self._log("Creator loop started.")
  270. while not self.m_stop_event.is_set():
  271. try:
  272. time.sleep(1)
  273. for apt in self.m_cfg.appointment_types:
  274. r_key = apt.routing_key
  275. with self.m_lock:
  276. active = sum(1 for t in self.m_tasks if t.source_queue == r_key)
  277. target = self.m_cfg.booker.target_instances
  278. if active < target:
  279. self._spawn_worker(r_key)
  280. except Exception as e:
  281. self._log(f'Creator loop exception:{e}')
  282. def _spawn_worker(self, target_routing_key: str):
  283. instance = None
  284. success = False
  285. task_id = None
  286. try:
  287. queue_name = f"auto.{target_routing_key}"
  288. task_data = VSCloudApi.Instance().get_vas_task_pop(queue_name)
  289. if not task_data:
  290. return
  291. task_id = task_data['id']
  292. with self.m_lock:
  293. self.m_task_data_cache[str(task_id)] = task_data
  294. self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + 8*60.0})
  295. user_inputs = task_data.get('user_inputs', {})
  296. plg_cfg = VSPlgConfig()
  297. plg_cfg.debug = self.m_cfg.debug
  298. plg_cfg.free_config = self.m_cfg.free_config
  299. plg_cfg.session_max_life = self.m_cfg.session_max_life
  300. plg_cfg.account.username = user_inputs.get("username", "")
  301. plg_cfg.account.password = user_inputs.get("password", "")
  302. if not plg_cfg.account.username:
  303. return
  304. acceptable_keys = [target_routing_key]
  305. if self.m_cfg.need_proxy:
  306. proxy = VSCloudApi.Instance().get_next_proxy(self.m_cfg.proxy_pool, self.m_cfg.proxy_cd)
  307. plg_cfg.proxy = type(plg_cfg.proxy)(**proxy)
  308. instance = self.m_factory.create(self.m_cfg.identifier, self.m_cfg.plugin_config.plugin_name)
  309. instance.set_log(self.m_logger)
  310. instance.set_config(plg_cfg)
  311. instance.create_session()
  312. with self.m_lock:
  313. self.m_tasks.append(
  314. Task(
  315. instance=instance,
  316. next_run=time.time(),
  317. task_ref=task_id,
  318. acceptable_routing_keys=acceptable_keys,
  319. source_queue=target_routing_key,
  320. book_allowed=True,
  321. next_remote_ping=time.time() + random.gauss(self.m_cfg.booker.keep_alive, 5)
  322. )
  323. )
  324. success = True
  325. self._log(f"+++ Order Booker spawned: {plg_cfg.account.username} (Target: {acceptable_keys})")
  326. except Exception as e:
  327. err_str = str(e)
  328. self._log(f"Order Booker spawn failed: {err_str}")
  329. rate_limited_indicators = [
  330. "42901" in err_str,
  331. "Rate limited" in err_str
  332. ]
  333. if any(rate_limited_indicators):
  334. if task_id is not None:
  335. self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + self.m_cfg.login_backoff})
  336. finally:
  337. if not success:
  338. if task_id:
  339. with self.m_lock:
  340. self.m_task_data_cache.pop(str(task_id), None)
  341. if instance:
  342. instance.cleanup()