booker_order.py 18 KB

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