import os import time import json import threading import random from datetime import datetime from typing import List, Dict, Callable, Optional from vs_types import GroupConfig, VSPlgConfig, Task, VSQueryResult, AppointmentType, AvailabilityStatus from vs_plg_factory import VSPlgFactory from toolkit.thread_pool import ThreadPool from toolkit.vs_cloud_api import VSCloudApi from utils.safe_redis_cli import SafeRedisClient class BaseBookerGCO: """ Booker 基类,封装公共的基础设施与通用的生命周期管理、逻辑循环等。 """ TAG = "BOOKER" def __init__(self, cfg: GroupConfig, redis_conf: Dict, logger: Callable[[str], None] = None): self.m_cfg = cfg self.m_factory = VSPlgFactory() self.m_logger = logger self.m_tasks: List[Task] = [] self.m_lock = threading.RLock() self.m_stop_event = threading.Event() self.redis_client = SafeRedisClient(redis_conf, self.m_logger) self.m_tracker_key = f"vs:worker:tasks_tracker:{self.m_cfg.identifier}" def _log(self, message: str): prefix = f'[{self.TAG}] [{self.m_cfg.identifier}]' if self.m_logger: self.m_logger(f'{prefix} {message}') else: print(f'{prefix} {message}') def start(self): if not self.m_cfg.enable: return self._log(f"Starting {self.TAG}...") plugin_name = self.m_cfg.plugin_config.plugin_name class_name = "".join(part.title() for part in plugin_name.split('_')) plugin_path = os.path.join(self.m_cfg.plugin_config.lib_path, self.m_cfg.plugin_config.plugin_bin) self.m_factory.register_plugin(plugin_name, plugin_path, class_name) threading.Thread(target=self._booking_trigger_loop, daemon=True).start() threading.Thread(target=self._creator_loop, daemon=True).start() threading.Thread(target=self._maintain_loop, daemon=True).start() self._start_additional_threads() def _start_additional_threads(self): """子类扩展线程的 Hook""" pass def stop(self): self._log("Stopping Booker...") self.m_stop_event.set() self._cleanup_all_tasks("booker stop") def update_config(self, new_cfg: GroupConfig): """动态更新配置""" with self.m_lock: if self.m_cfg.enable and not new_cfg.enable: self._log("Config dynamically updated: Group DISABLED. Will stop creating new tasks.") elif not self.m_cfg.enable and new_cfg.enable: self._log("Config dynamically updated: Group ENABLED.") else: self._log("Config dynamically updated: Parameters refreshed.") self.m_cfg = new_cfg def _cleanup_task(self, task: Task, reason: str = ""): try: if task and task.instance: task.instance.cleanup() ref_str = f" for task={task.task_ref}" if task.task_ref else "" self._log(f"🧹 Cleaned up instance{ref_str}. Reason: {reason}") except Exception as e: self._log(f"Cleanup failed for instance. Reason: {reason}. Error: {e}") def _remove_task(self, task: Task, reason: str = "", cleanup: bool = True): removed = False with self.m_lock: if task in self.m_tasks: self.m_tasks.remove(task) removed = True self._on_task_removed(task) if cleanup and removed: self._cleanup_task(task, reason) return removed def _on_task_removed(self, task: Task): """子类在 task 被移除时的自定义 Hook(如清理缓存)""" pass def _cleanup_all_tasks(self, reason: str = ""): with self.m_lock: tasks = list(self.m_tasks) self.m_tasks.clear() self._on_all_tasks_cleaned() for task in tasks: self._cleanup_task(task, reason) def _on_all_tasks_cleaned(self): """子类在所有 task 被清空时的自定义 Hook""" pass def _get_redis_key(self, routing_key: str) -> str: return f"vs:signal:{routing_key}" def _is_within_active_hours(self) -> bool: """判断当前是否在允许创建实例的时间段内""" start_str = self.m_cfg.active_time_start end_str = self.m_cfg.active_time_end current_bj_time = datetime.utcnow().time() start_time = datetime.strptime(start_str, "%H:%M").time() end_time = datetime.strptime(end_str, "%H:%M").time() return start_time <= current_bj_time <= end_time def _maintain_loop(self): self._log("Maintain loop started.") while not self.m_stop_event.is_set(): try: time.sleep(1.0) with self.m_lock: tasks_to_check = list(self.m_tasks) if not tasks_to_check: continue dead_tasks = [] healthy_tasks = [] now = time.time() for t in tasks_to_check: if now >= t.next_remote_ping: t.instance.keep_alive() if t.instance.health_check(): healthy_tasks.append(t) t.next_remote_ping = now + random.gauss(self.m_cfg.booker.keep_alive, 5) else: dead_tasks.append(t) else: healthy_tasks.append(t) self._on_maintain_ping(healthy_tasks, dead_tasks) if dead_tasks: with self.m_lock: current_tasks = list(self.m_tasks) self.m_tasks = [t for t in self.m_tasks if t in healthy_tasks] for t in dead_tasks: if t in current_tasks: self._cleanup_task(t, "unhealthy or keep-alive failed") else: with self.m_lock: self.m_tasks = [t for t in self.m_tasks if t in healthy_tasks] except Exception as e: self._log(f'Maintain loop exception: {e}') def _on_maintain_ping(self, healthy_tasks: List[Task], dead_tasks: List[Task]): """子类维护循环中对于健康/异常 Task 的额外处理""" pass def _is_date_of_interest(self, task: Task, query_result: VSQueryResult) -> bool: """判断 query_result 中的可用日期是否在 task 意向范围内(默认 True,Order 模式重写)""" return True def _booking_trigger_loop(self): self._log("Trigger loop started.") while not self.m_stop_event.is_set(): try: time.sleep(0.1) now = time.time() for apt_type in self.m_cfg.appointment_types: redis_key = self._get_redis_key(apt_type.routing_key) raw_data = self.redis_client.get(redis_key) if not raw_data: continue try: data = json.loads(raw_data) query_result = VSQueryResult.model_validate(data['query_result']) query_result.apt_type = AppointmentType.model_validate(data['apt_type']) except Exception as parse_err: self._log(f"Data parsing error for {redis_key}: {parse_err}. Deleting corrupted signal.") self.redis_client.delete(redis_key) continue matching_tasks = [] with self.m_lock: for task in self.m_tasks: if now < task.next_run or not task.book_allowed: continue if apt_type.routing_key not in task.acceptable_routing_keys: continue if not self._is_date_of_interest(task, query_result): continue task.next_run = now + self.m_cfg.booker.booking_cooldown matching_tasks.append(task) if matching_tasks: threads = [] for task in matching_tasks: self._log(f"🚀 Triggering BOOK for {apt_type.routing_key} | Order Ref: {task.task_ref}") t = threading.Thread(target=self._execute_book_job, args=(task, query_result)) threads.append(t) t.start() for t in threads: t.join() except Exception as e: self._log(f"Booking trigger loop exception: {e}") def _execute_book_job(self, task: Task, query_result: VSQueryResult): raise NotImplementedError def _creator_loop(self): raise NotImplementedError def _push_success_notification(self, task_id, order_id, book_res): """成功落单后同步云端与推送微信通知的公共方法""" grab_info = { "account": book_res.account, "session_id": book_res.session_id, "urn": book_res.urn, "slot_date": book_res.book_date, "slot_time": book_res.book_time, "timestamp": int(time.time()), "payment_link": book_res.payment_link } def _update_cloud_success(): try: VSCloudApi.Instance().update_vas_task(str(task_id), {"status": "grabbed", "grabbed_history": grab_info}) push_content = ( f"🎉 【预定成功通知】\n" f"━━━━━━━━━━━━━━━\n" f"订单编号: {order_id}\n" f"预约账号: {book_res.account}\n" f"预约日期: {book_res.book_date}\n" f"预约时间: {book_res.book_time}\n" f"预约编号: {book_res.urn}\n" f"支付链接: {book_res.payment_link if book_res.payment_link else '无需支付/暂无'}\n" f"━━━━━━━━━━━━━━━\n" ) VSCloudApi.Instance().push_weixin_text(push_content) except Exception as e: self._log(f"Failed to update success state to cloud: {e}") ThreadPool.getInstance().enqueue(_update_cloud_success) self.redis_client.zrem(self.m_tracker_key, task_id) class BuiltinBookerGCO(BaseBookerGCO): """ 非绑定模式 (公共内置账号池): - 只维护全局 target_instances 数量的实例。 - 所有实例热机等待,发现信号后临时去云端 Pop 订单。 """ TAG = "BUILTIN-BOOKER" def _creator_loop(self): self._log("Creator loop started.") while not self.m_stop_event.wait(1.0): try: if not self._is_within_active_hours(): continue with self.m_lock: current = len(self.m_tasks) target = self.m_cfg.booker.target_instances if current < target: self._spawn_worker() except Exception as e: self._log(f'Creator loop exception: {e}') def _spawn_worker(self): instance = None success = False plg_cfg = None try: plg_cfg = VSPlgConfig() plg_cfg.debug = self.m_cfg.debug plg_cfg.free_config = self.m_cfg.free_config plg_cfg.session_max_life = self.m_cfg.session_max_life if self.m_cfg.need_account: acc = VSCloudApi.Instance().get_next_account(self.m_cfg.booker.account_pool_id, self.m_cfg.booker.account_cd) plg_cfg.account = type(plg_cfg.account)(**acc) if self.m_cfg.need_proxy: proxy = VSCloudApi.Instance().get_next_proxy(self.m_cfg.proxy_pool, self.m_cfg.proxy_cd) plg_cfg.proxy = type(plg_cfg.proxy)(**proxy) instance = self.m_factory.create(self.m_cfg.identifier, self.m_cfg.plugin_config.plugin_name) instance.set_log(self.m_logger) instance.set_config(plg_cfg) instance.create_session() success = True with self.m_lock: all_keys = [apt.routing_key for apt in self.m_cfg.appointment_types] self.m_tasks.append( Task( instance=instance, next_run=time.time(), task_ref=None, acceptable_routing_keys=all_keys, source_queue="built-in", book_allowed=True, next_remote_ping=time.time() + random.gauss(self.m_cfg.booker.keep_alive, 5) ) ) success = True self._log(f"+++ Built-in Booker spawned: {plg_cfg.account.username}") except Exception as e: err_str = str(e) self._log(f"Spawn failed: {err_str}") rate_limited_indicators = [ "42901" in err_str, "Rate limited" in err_str ] if any(rate_limited_indicators): if plg_cfg and plg_cfg.account.username != "Guest": VSCloudApi.lock_account(plg_cfg.account.id, self.m_cfg.login_backoff) finally: if not success: if instance: instance.cleanup() def _execute_book_job(self, task: Task, query_result: VSQueryResult): queue_name = f"auto.{query_result.apt_type.routing_key}" task_id = None task_data = None try: task_data = VSCloudApi.Instance().get_vas_task_pop(queue_name) if not task_data: return task_id = task_data['id'] order_id = task_data.get('order_id') self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + 30.0}) user_input = task_data.get('user_inputs', {}) book_res = task.instance.book(query_result, user_input) if book_res.success: self._log(f"✅ BOOK SUCCESS! Order: {order_id}") self._push_success_notification(task_id, order_id, book_res) task.successful_bookings += 1 max_b = self.m_cfg.booker.max_bookings_per_account if max_b > 0 and task.successful_bookings >= max_b: self._log(f"Account reached max bookings ({max_b}). Destroying instance.") self._remove_task(task, "max bookings reached") else: self._log(f"❌ BOOK FAILED for Order: {order_id}") except Exception as e: err_str = str(e) self._log(f"Exception during booking: {err_str}") rate_limited_indicators = [ "42901" in err_str, "Rate limited" in err_str ] if any(rate_limited_indicators): self._remove_task(task, "booking rate limited") class OrderBookerGCO(BaseBookerGCO): """ 绑定模式 (订单自带账号): - 按城市队列维护热机配额。 - 绝对的 1 对 1 关系:一个实例绑定一个云端订单。 - 预订成功后,实例立即销毁。 """ TAG = "ORDER-BOOKER" def __init__(self, cfg: GroupConfig, redis_conf: Dict, logger: Callable[[str], None] = None): super().__init__(cfg, redis_conf, logger) self.m_task_data_cache: Dict[str, dict] = {} self.heartbeat_ttl = 2 * 60.0 def _start_additional_threads(self): threading.Thread(target=self._cache_refresh_loop, daemon=True).start() def _on_task_removed(self, task: Task): task_id = task.task_ref if task_id: self.m_task_data_cache.pop(str(task_id), None) def _on_all_tasks_cleaned(self): self.m_task_data_cache.clear() def _on_maintain_ping(self, healthy_tasks: List[Task], dead_tasks: List[Task]): if healthy_tasks: new_deadline = time.time() + self.heartbeat_ttl mapping = {str(t.task_ref): new_deadline for t in healthy_tasks} self.redis_client.bulk_zadd(self.m_tracker_key, mapping) if dead_tasks: mapping = {str(t.task_ref): 0 for t in dead_tasks} self.redis_client.bulk_zadd(self.m_tracker_key, mapping) def _cache_refresh_loop(self): self._log("Cache refresh loop started.") refresh_interval = 15 * 60 while not self.m_stop_event.is_set(): try: time.sleep(1) with self.m_lock: tasks_to_check = { tid: data.get('_last_refresh', 0) for tid, data in self.m_task_data_cache.items() } if not tasks_to_check: continue now = time.time() for tid, last_refresh in tasks_to_check.items(): if now - last_refresh >= refresh_interval: fresh_data = VSCloudApi.Instance().get_vas_task(tid) if fresh_data: fresh_data['_last_refresh'] = time.time() with self.m_lock: if tid in self.m_task_data_cache: self.m_task_data_cache[tid] = fresh_data time.sleep(0.5) except Exception as e: self._log(f'Cache refresh loop exception: {e}') def _is_date_of_interest(self, task: Task, query_result: VSQueryResult) -> bool: if query_result.availability_status != AvailabilityStatus.Available: return True task_id = task.task_ref task_data = self.m_task_data_cache.get(str(task_id), {}) user_input = task_data.get('user_inputs', {}) expected_end_date = ( user_input.get('expected_end_date') or '2100-01-01' ) available_date = query_result.earliest_date dt = available_date.strftime("%Y-%m-%d") return dt <= expected_end_date def _execute_book_job(self, task: Task, query_result: VSQueryResult): task_id = task.task_ref task_data = None try: with self.m_lock: task_data = self.m_task_data_cache.get(str(task_id)) if not task_data or task_data.get('status') in ['grabbed', 'pause', 'completed', 'cancelled']: self._log(f"Bound Task={task_id} is no longer valid or already processed. Removing instance.") self._remove_task(task, "bound task no longer valid") self.redis_client.zrem(self.m_tracker_key, task_id) return order_id = task_data.get('order_id') user_input = task_data.get('user_inputs', {}) book_res = task.instance.book(query_result, user_input) if book_res.success: self._log(f"✅ BOOK SUCCESS! Order: {order_id}. Destroying instance.") self._push_success_notification(task_id, order_id, book_res) self._remove_task(task, "booking success") else: self._log(f"❌ BOOK FAILED for Order: {order_id}. Will retry on next signal.") except Exception as e: err_str = str(e) self._log(f"Exception during booking: {err_str}") rate_limited_indicators = [ "42901" in err_str, "Rate limited" in err_str ] if any(rate_limited_indicators): self._remove_task(task, "booking rate limited") def _creator_loop(self): self._log("Creator loop started.") while not self.m_stop_event.wait(1.0): try: if not self._is_within_active_hours(): continue for apt in self.m_cfg.appointment_types: r_key = apt.routing_key with self.m_lock: active = sum(1 for t in self.m_tasks if t.source_queue == r_key) target = self.m_cfg.booker.target_instances if active < target: self._spawn_worker(r_key) except Exception as e: self._log(f'Creator loop exception:{e}') def _spawn_worker(self, target_routing_key: str): instance = None success = False task_id = None try: queue_name = f"auto.{target_routing_key}" task_data = VSCloudApi.Instance().get_vas_task_pop(queue_name) if not task_data: return task_id = task_data['id'] with self.m_lock: self.m_task_data_cache[str(task_id)] = task_data self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + 8*60.0}) user_inputs = task_data.get('user_inputs', {}) plg_cfg = VSPlgConfig() plg_cfg.debug = self.m_cfg.debug plg_cfg.free_config = self.m_cfg.free_config plg_cfg.session_max_life = self.m_cfg.session_max_life plg_cfg.account.username = user_inputs.get("username", "") plg_cfg.account.password = user_inputs.get("password", "") if not plg_cfg.account.username: return acceptable_keys = [target_routing_key] if self.m_cfg.need_proxy: proxy = VSCloudApi.Instance().get_next_proxy(self.m_cfg.proxy_pool, self.m_cfg.proxy_cd) plg_cfg.proxy = type(plg_cfg.proxy)(**proxy) instance = self.m_factory.create(self.m_cfg.identifier, self.m_cfg.plugin_config.plugin_name) instance.set_log(self.m_logger) instance.set_config(plg_cfg) instance.create_session() with self.m_lock: self.m_tasks.append( Task( instance=instance, next_run=time.time(), task_ref=task_id, acceptable_routing_keys=acceptable_keys, source_queue=target_routing_key, book_allowed=True, next_remote_ping=time.time() + random.gauss(self.m_cfg.booker.keep_alive, 5) ) ) success = True self._log(f"+++ Order Booker spawned: {plg_cfg.account.username} (Target: {acceptable_keys})") except Exception as e: err_str = str(e) self._log(f"Order Booker spawn failed: {err_str}") rate_limited_indicators = [ "42901" in err_str, "Rate limited" in err_str ] if any(rate_limited_indicators): if task_id is not None: self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + self.m_cfg.login_backoff}) finally: if not success: if task_id: with self.m_lock: self.m_task_data_cache.pop(str(task_id), None) if instance: instance.cleanup()