Hujiarui před 2 měsíci
rodič
revize
33f0454f2a
6 změnil soubory, kde provedl 323 přidání a 247 odebrání
  1. 73 66
      booker_builtin.py
  2. 129 125
      booker_order.py
  3. 32 49
      sentinel.py
  4. 7 7
      test/test_publish_slot.py
  5. 75 0
      utils/safe_redis_cli.py
  6. 7 0
      vs_plg.py

+ 73 - 66
booker_builtin.py

@@ -3,15 +3,15 @@ import time
 import json
 import threading
 import random
-import redis
 from typing import List, Dict, Callable
 
-import configure
 from vs_types import GroupConfig, VSPlgConfig, Task, VSQueryResult, AppointmentType
 from vs_plg_factory import VSPlgFactory 
 from toolkit.thread_pool import ThreadPool 
 from toolkit.vs_cloud_api import VSCloudApi
 from toolkit.backoff import ExponentialBackoff
+from utils.safe_redis_cli import SafeRedisClient
+
 
 class BuiltinBookerGCO:
     """
@@ -26,7 +26,7 @@ class BuiltinBookerGCO:
         self.m_tasks: List[Task] = []
         self.m_lock = threading.RLock()
         self.m_stop_event = threading.Event()
-        self.redis_client = redis.Redis(**redis_conf)
+        self.redis_client = SafeRedisClient(redis_conf, self.m_logger)
         self.m_pending_builtin = 0
         
         self.m_tracker_key = f"vs:worker:tasks_tracker:{self.m_cfg.identifier}"
@@ -38,6 +38,8 @@ class BuiltinBookerGCO:
     def _log(self, message):
         if self.m_logger:
             self.m_logger(f'[BUILTIN-BOOKER] [{self.m_cfg.identifier}] {message}')
+        else:
+            print(f'[BUILTIN-BOOKER] [{self.m_cfg.identifier}] {message}')
 
     def start(self):
         if not self.m_cfg.enable:
@@ -56,12 +58,25 @@ class BuiltinBookerGCO:
         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:
-            instance = getattr(task, 'instance', None)
-            if instance and hasattr(instance, 'cleanup'):
-                instance.cleanup()
+            if task and task.instance:
+                task.instance.cleanup()
                 self._log(f"🧹 Cleaned up built-in instance. Reason: {reason}")
         except Exception as e:
             self._log(f"Cleanup failed for built-in instance. Reason: {reason}. Error: {e}")
@@ -89,20 +104,19 @@ class BuiltinBookerGCO:
     def _maintain_loop(self):
         self._log("Maintain loop started.")
         while not self.m_stop_event.is_set():
-            time.sleep(1.0)
-            now = time.time()
-            
-            with self.m_lock:
-                tasks_to_check = list(self.m_tasks)
-            
-            if not tasks_to_check:
-                continue
-            
-            healthy_tasks = []
-            dead_tasks = []
-            for t in tasks_to_check:
-                if now >= t.next_remote_ping:
-                    try:
+            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)
@@ -111,28 +125,27 @@ class BuiltinBookerGCO:
                         else:
                             dead_tasks.append(t)
                             self._log(f"♻️ Instance unhealthy. Will be removed.")
-                    except Exception as e:
-                        dead_tasks.append(t)
-                        self._log(f"Instance keep-alive failed: {e}")
+                    else:
+                        healthy_tasks.append(t)
+                
+                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:
-                    healthy_tasks.append(t)
-            
-            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]
+                    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 _booking_trigger_loop(self):
         self._log("Trigger loop started.")
         while not self.m_stop_event.is_set():
             try:
-                time.sleep(1.0)
+                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)
@@ -171,8 +184,7 @@ class BuiltinBookerGCO:
                             t.join() 
                     
             except Exception as e:
-                self._log(f"Trigger loop error: {e}")
-                time.sleep(2)
+                self._log(f"Booking trigger loop exception: {e}")
 
     def _execute_book_job(self, task: Task, query_result: VSQueryResult):
         queue_name = f"auto.{query_result.apt_type.routing_key}"
@@ -245,10 +257,12 @@ class BuiltinBookerGCO:
                     t_fails = task_meta.get('booking_failures', 0) + 1
                     task_meta['booking_failures'] = t_fails
                     
-                    try:
-                        VSCloudApi.Instance().update_vas_task(task_id, {"meta": task_meta})
-                    except Exception as cloud_err:
-                        self._log(f"Failed to update task meta: {cloud_err}")
+                    def _update_cloud_meta():
+                        try:
+                            VSCloudApi.Instance().update_vas_task(str(task_id), {"meta": task_meta})
+                        except Exception as cloud_err:
+                            self._log(f"Failed to update task meta: {cloud_err}")
+                    ThreadPool.getInstance().enqueue(_update_cloud_meta) 
                         
                     t_cd = self.task_backoff.calculate(t_fails)
                     self._log(f"⏳ Task={task_id} (Booking Attempt {t_fails}) suspended for {t_cd:.1f}s.")
@@ -264,18 +278,21 @@ class BuiltinBookerGCO:
         spawn_interval = 10.0
         group_cd_key = f"vs:group:cooldown:{self.m_cfg.identifier}"
         while not self.m_stop_event.is_set():
-            time.sleep(2.0)
-            if self.redis_client.exists(group_cd_key):
-                continue
-            with self.m_lock:
-                current = len(self.m_tasks)
-                pending = self.m_pending_builtin
-                target = self.m_cfg.booker.target_instances
-            if (current + pending) < target:
-                now = time.time()
-                if now - self.m_last_spawn_time >= spawn_interval:
-                    self.m_last_spawn_time = now 
-                    self._spawn_worker()
+            try:
+                time.sleep(1.0)
+                if self.redis_client.exists(group_cd_key):
+                    continue
+                with self.m_lock:
+                    current = len(self.m_tasks)
+                    pending = self.m_pending_builtin
+                    target = self.m_cfg.booker.target_instances
+                if (current + pending) < target:
+                    now = time.time()
+                    if now - self.m_last_spawn_time >= spawn_interval:
+                        self.m_last_spawn_time = now 
+                        self._spawn_worker()
+            except Exception as e:
+                self._log(f'Creator loop exception: {e}')
 
     def _spawn_worker(self):
         with self.m_lock:
@@ -290,18 +307,11 @@ class BuiltinBookerGCO:
 
                 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.id = acc['id']
-                    plg_cfg.account.username = acc['username']
-                    plg_cfg.account.password = acc['password']
+                    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.id = proxy['id']
-                    plg_cfg.proxy.ip = proxy['ip']
-                    plg_cfg.proxy.port = proxy['port']
-                    plg_cfg.proxy.proto = proxy['proto']
-                    plg_cfg.proxy.username = proxy['username']
-                    plg_cfg.proxy.password = proxy['password']
+                    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)
@@ -347,11 +357,8 @@ class BuiltinBookerGCO:
                     group_fail_key = f"vs:group:failures:{self.m_cfg.identifier}"
                     group_cd_key = f"vs:group:cooldown:{self.m_cfg.identifier}"
                     
-                    # 更新全局(机器组)失败次数
                     g_fails = self.redis_client.incr(group_fail_key)
-                    # 计算退避时间
                     g_cd = self.group_backoff.calculate(g_fails)
-                    # 设置 Redis 全局冷却保护阀
                     self.redis_client.set(group_cd_key, "1", ex=int(g_cd))
                     self._log(f"📉 [Rate Limited] Group '{self.m_cfg.identifier}' failed {g_fails} times. Global Backoff: {g_cd:.1f}s.")
 

+ 129 - 125
booker_order.py

@@ -3,14 +3,16 @@ import time
 import json
 import threading
 import random
-import redis
-from typing import List, Dict, Callable, Any, Optional
+from datetime import datetime
+from typing import List, Dict, Callable, Any
 
 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 toolkit.backoff import ExponentialBackoff
+from utils.safe_redis_cli import SafeRedisClient
+
 
 class OrderBookerGCO:
     """
@@ -27,7 +29,8 @@ class OrderBookerGCO:
         self.m_lock = threading.RLock()
         self.m_stop_event = threading.Event()
 
-        self.redis_client = redis.Redis(**redis_conf)
+        self.redis_client = SafeRedisClient(redis_conf, self.m_logger)
+        
         self.m_pending_order_by_queue: Dict[str, int] = {}        
         self.m_last_spawn_times: Dict[str, float] = {}
         self.m_task_data_cache: Dict[str, dict] = {}
@@ -41,6 +44,8 @@ class OrderBookerGCO:
     def _log(self, message):
         if self.m_logger:
             self.m_logger(f'[ORDER-BOOKER] [{self.m_cfg.identifier}] {message}')
+        else:
+            print(f'[ORDER-BOOKER] [{self.m_cfg.identifier}] {message}')
 
     def start(self):
         if not self.m_cfg.enable:
@@ -60,15 +65,28 @@ class OrderBookerGCO:
         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:
-            instance = getattr(task, 'instance', None)
-            if instance and hasattr(instance, 'cleanup'):
-                instance.cleanup()
-                self._log(f"🧹 Cleaned up instance for task={getattr(task, 'task_ref', None)}. Reason: {reason}")
+            if task and task.instance:
+                task.instance.cleanup()
+                self._log(f"🧹 Cleaned up instance for task={task.task_ref}. Reason: {reason}")
         except Exception as e:
-            self._log(f"Cleanup failed for task={getattr(task, 'task_ref', None)}. Reason: {reason}. Error: {e}")
+            self._log(f"Cleanup failed for task={task.task_ref}. Reason: {reason}. Error: {e}")
 
     def _remove_task(self, task: Task, reason: str = "", cleanup: bool = True):
         removed = False
@@ -76,7 +94,7 @@ class OrderBookerGCO:
             if task in self.m_tasks:
                 self.m_tasks.remove(task)
                 removed = True
-            task_id = str(getattr(task, 'task_ref', ''))
+            task_id = task.task_ref
             self.m_task_data_cache.pop(task_id, None)
                 
         if cleanup and removed:
@@ -96,26 +114,20 @@ class OrderBookerGCO:
             
     def _maintain_loop(self):
         self._log("Maintain loop started.")
-        heartbeat_interval = 30 
         while not self.m_stop_event.is_set():
-            for _ in range(heartbeat_interval):
-                if self.m_stop_event.is_set():
-                    return
-                time.sleep(1.0)
-            
-            with self.m_lock:
-                tasks_to_check = list(self.m_tasks)
+            try:
+                time.sleep(1)
+                with self.m_lock:
+                    tasks_to_check = list(self.m_tasks)
+                    
+                if not tasks_to_check:
+                    continue
                 
-            if not tasks_to_check:
-                continue
-            
-            healthy_tasks = []
-            dead_tasks = []
-            now = time.time()
-            
-            for t in tasks_to_check:
-                if now >= t.next_remote_ping:
-                    try:
+                healthy_tasks = []
+                dead_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)
@@ -125,78 +137,67 @@ class OrderBookerGCO:
                         else:
                             dead_tasks.append(t)
                             self._log(f"♻️ Instance for task={t.task_ref} unhealthy.")
-                    except Exception as e:
-                        dead_tasks.append(t)
-                        self._log(f"♻️ Instance for task={t.task_ref} keep-alive failed: {e}.")
-                else:
-                    healthy_tasks.append(t)
-            
-            if healthy_tasks:
-                try:
-                    pipeline = self.redis_client.pipeline()
+                    else:
+                        healthy_tasks.append(t)
+                
+                if healthy_tasks:
                     new_deadline = time.time() + self.heartbeat_ttl
-                    for t in healthy_tasks:
-                        if t.task_ref is not None:
-                            pipeline.zadd(self.m_tracker_key, {str(t.task_ref): new_deadline})
-                    pipeline.execute()
-                    self._log(f"💓 Heartbeat sent. Renewed {len(healthy_tasks)} tasks.")
-                except Exception as e:
-                    self._log(f"Redis Heartbeat update failed: {e}")
+                    mapping = {str(t.task_ref): new_deadline for t in healthy_tasks if t.task_ref is not None}
+                    self.redis_client.bulk_zadd(self.m_tracker_key, mapping)
+                    # self._log(f"💓 Heartbeat sent. Renewed {len(healthy_tasks)} tasks.")
 
-            if dead_tasks:
-                try:
-                    pipeline = self.redis_client.pipeline()
-                    for t in dead_tasks:
-                        if t.task_ref is not None:
-                            pipeline.zadd(self.m_tracker_key, {str(t.task_ref): 0})
-                    pipeline.execute()
+                if dead_tasks:
+                    mapping = {str(t.task_ref): 0 for t in dead_tasks if t.task_ref is not None}
+                    self.redis_client.bulk_zadd(self.m_tracker_key, mapping)
                     self._log(f"🗑️ Handed over {len(dead_tasks)} dead tasks to Sweeper.")
-                except Exception as e:
-                    pass
-            
-            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]
+                
+                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 _cache_refresh_loop(self):
         self._log("Cache refresh loop started.")
-        refresh_interval = 15*60
+        refresh_interval = 15 * 60
         
         while not self.m_stop_event.is_set():
-            for _ in range(refresh_interval):
-                if self.m_stop_event.is_set():
-                    return
-                time.sleep(1.0)
-            with self.m_lock:
-                task_ids = list(self.m_task_data_cache.keys())
-            if not task_ids:
-                continue
-            for tid in task_ids:
-                if self.m_stop_event.is_set():
-                    break
-                try:
-                    fresh_data = VSCloudApi.Instance().get_vas_task(tid)
-                    if fresh_data:
-                        with self.m_lock:
-                            if tid in self.m_task_data_cache:
-                                self.m_task_data_cache[tid] = fresh_data
-                except Exception:
-                    pass
-                time.sleep(0.5)
+            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, query_result) -> bool:
+    def _is_date_of_interest(self, task, query_result: VSQueryResult) -> bool:
         """
-        判断 query_result 中的可用日期,
-        是否在 task 的意向日期范围内。
+        判断 query_result 中的可用日期,是否在 task 的意向日期范围内。
         """
-
         if query_result.availability_status != AvailabilityStatus.Available:
             return True
         task_id = task.task_ref
@@ -211,7 +212,8 @@ class OrderBookerGCO:
             or '2100-01-01'
         )
         available_date = query_result.earliest_date
-        return available_date <= expected_end_date
+        dt = available_date.strftime("%Y-%m-%d")
+        return dt <= expected_end_date
 
     def _booking_trigger_loop(self):
         self._log("Trigger loop started.")
@@ -256,8 +258,7 @@ class OrderBookerGCO:
                         for t in threads:
                             t.join() 
             except Exception as e:
-                self._log(f"Trigger loop error: {e}")
-                time.sleep(2)
+                self._log(f"Booking trigger loop exception: {e}")
 
     def _execute_book_job(self, task: Task, query_result: VSQueryResult):
         task_id = task.task_ref
@@ -270,7 +271,6 @@ class OrderBookerGCO:
                 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', {})
@@ -333,33 +333,37 @@ class OrderBookerGCO:
                         except Exception as cloud_err:
                             self._log(f"Failed to update task meta: {cloud_err}")
                     ThreadPool.getInstance().enqueue(_update_cloud_meta)   
+                    
                     t_cd = self.task_backoff.calculate(t_fails)
                     self._log(f"⏳ Task={task_id} (Booking Attempt {t_fails}) suspended for {t_cd:.1f}s.")
                     self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + t_cd})
-
+                    
     def _creator_loop(self):
         self._log("Creator loop started.")
         spawn_interval = 10.0
         while not self.m_stop_event.is_set():
-            time.sleep(2.0)
-            now = time.time()
-            for apt in self.m_cfg.appointment_types:
-                r_key = apt.routing_key
-                
-                queue_cd_key = f"vs:queue:cooldown:{r_key}"
-                if self.redis_client.exists(queue_cd_key):
-                    continue
-                
-                with self.m_lock:
-                    active = sum(1 for t in self.m_tasks if getattr(t, 'source_queue', '') == r_key)
-                    pending = self.m_pending_order_by_queue.get(r_key, 0)
-                    target = self.m_cfg.booker.target_instances
+            try:
+                time.sleep(1)
+                now = time.time()
+                for apt in self.m_cfg.appointment_types:
+                    r_key = apt.routing_key
+                    queue_cd_key = f"vs:queue:cooldown:{r_key}"
                 
-                if (active + pending) < target:
-                    last_spawn = self.m_last_spawn_times.get(r_key, 0.0)
-                    if now - last_spawn >= spawn_interval:
-                        self.m_last_spawn_times[r_key] = now 
-                        self._spawn_worker(r_key)
+                    if self.redis_client.exists(queue_cd_key):
+                        continue
+
+                    with self.m_lock:
+                        active = sum(1 for t in self.m_tasks if t.source_queue == r_key)
+                        pending = self.m_pending_order_by_queue.get(r_key, 0)
+                        target = self.m_cfg.booker.target_instances
+                    
+                    if (active + pending) < target:
+                        last_spawn = self.m_last_spawn_times.get(r_key, 0.0)
+                        if now - last_spawn >= spawn_interval:
+                            self.m_last_spawn_times[r_key] = now 
+                            self._spawn_worker(r_key)
+            except Exception as e:
+                self._log(f'Creator loop exception:{e}')
 
     def _spawn_worker(self, target_routing_key: str):
         with self.m_lock: 
@@ -396,12 +400,7 @@ class OrderBookerGCO:
                 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.id = proxy['id']
-                    plg_cfg.proxy.ip = proxy['ip']
-                    plg_cfg.proxy.port = proxy['port']
-                    plg_cfg.proxy.proto = proxy['proto']
-                    plg_cfg.proxy.username = proxy['username']
-                    plg_cfg.proxy.password = proxy['password']
+                    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)
@@ -421,9 +420,10 @@ class OrderBookerGCO:
                             next_remote_ping=time.time() + random.randint(55, 65)   
                         )
                     )
-                    queue_fail_key = f"vs:queue:failures:{target_routing_key}"
-                    self.redis_client.delete(queue_fail_key)                    
+                    
                 success = True
+                queue_fail_key = f"vs:queue:failures:{target_routing_key}"
+                self.redis_client.delete(queue_fail_key)                    
                 self._log(f"+++ Order Booker spawned: {plg_cfg.account.username} (Target: {acceptable_keys})")
             except Exception as e:
                 err_str = str(e)
@@ -445,23 +445,28 @@ class OrderBookerGCO:
                     is_rate_limited = True
                     queue_fail_key = f"vs:queue:failures:{target_routing_key}"
                     queue_cd_key = f"vs:queue:cooldown:{target_routing_key}"
+                    
                     q_fails = self.redis_client.incr(queue_fail_key)
                     q_cd = self.queue_backoff.calculate(q_fails)
                     self.redis_client.set(queue_cd_key, "1", ex=int(q_cd))
                     self._log(f"📉 [Rate Limited] Queue '{target_routing_key}' failed {q_fails} times. Global Backoff: {q_cd:.1f}s.")
+
                     if task_id is not None:
                         task_meta = task_data.get('meta') or {}
                         t_fails = task_meta.get('spawn_failures', 0) + 1
                         task_meta['spawn_failures'] = t_fails
                         
-                        try:
-                            VSCloudApi.Instance().update_vas_task(str(task_id), {"meta": task_meta})
-                        except Exception as cloud_err:
-                            self._log(f"Failed to update task meta: {cloud_err}")
+                        def _update_cloud_meta():
+                            try:
+                                VSCloudApi.Instance().update_vas_task(str(task_id), {"meta": task_meta})
+                            except Exception as cloud_err:
+                                self._log(f"Failed to update task meta: {cloud_err}")
+                        ThreadPool.getInstance().enqueue(_update_cloud_meta) 
                         
                         t_cd = self.account_backoff.calculate(t_fails)
                         self._log(f"⏳ Task={task_id} (Attempt {t_fails}) suspended for {t_cd:.1f}s.")
-                        self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + t_cd})       
+                        self.redis_client.zadd(self.m_tracker_key, {str(task_id): time.time() + t_cd})  
+         
             finally:
                 with self.m_lock: 
                     self.m_pending_order_by_queue[target_routing_key] = max(0, self.m_pending_order_by_queue[target_routing_key] - 1)
@@ -469,7 +474,6 @@ class OrderBookerGCO:
                 if not success and task_id is not None and not is_rate_limited:
                     self.redis_client.zadd(self.m_tracker_key, {str(task_id): 0})
                     self._log(f"♻️ Task={task_id} failed normal spawn. Instantly handed over to Sweeper.")
-                    
                     with self.m_lock:
                         self.m_task_data_cache.pop(str(task_id), None)
                         

+ 32 - 49
sentinel.py

@@ -3,7 +3,6 @@ import time
 import json
 import random
 import threading
-import redis
 from typing import List, Dict, Callable
 
 from vs_types import GroupConfig, VSPlgConfig, Task, QueryWaitMode
@@ -11,6 +10,7 @@ from vs_plg_factory import VSPlgFactory
 from toolkit.thread_pool import ThreadPool 
 from toolkit.vs_cloud_api import VSCloudApi
 from toolkit.backoff import ExponentialBackoff
+from utils.safe_redis_cli import SafeRedisClient
 
 class SentinelGCO:
     def __init__(self, cfg: GroupConfig, redis_conf: Dict, logger: Callable[[str], None] = None):
@@ -21,10 +21,9 @@ class SentinelGCO:
         self.m_lock = threading.RLock()
         self.m_stop_event = threading.Event()
         
-        self.redis_client = redis.Redis(**redis_conf)
+        self.redis_client = SafeRedisClient(redis_conf, self.m_logger)
         self.m_pending_builtin = 0
         
-        # 1. 全局建连退避:起步 1 分钟,封顶 1 小时 (保护登录接口)
         self.group_backoff = ExponentialBackoff(base_delay=60.0, max_delay=3600.0, factor=2.0)
         self.m_last_spawn_time = 0.0
         self.m_spawn_interval = 120
@@ -33,6 +32,8 @@ class SentinelGCO:
     def _log(self, message):
         if self.m_logger:
             self.m_logger(f'[SENTINEL] [{self.m_cfg.identifier}] {message}')
+        else:
+            print(f'[SENTINEL] [{self.m_cfg.identifier}] {message}')
             
     def _get_average_interval(self) -> float:
         """计算当前组平均的查询间隔(秒)"""
@@ -47,11 +48,9 @@ class SentinelGCO:
     
     def update_config(self, new_cfg: GroupConfig):
         """
-        动态更新配置。侵入性小,仅替换配置对象。
-        现有的 creator_loop 和 monitor_loop 下一次循环读取时即生效。
+        动态更新配置
         """
         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:
@@ -84,7 +83,7 @@ class SentinelGCO:
 
     def _cleanup_task(self, task: Task, reason: str):
         try:
-            if task and task.instance and hasattr(task.instance, "cleanup"):
+            if task and task.instance:
                 self._log(f"Cleaning up sentinel instance. reason={reason}")
                 task.instance.cleanup()
         except Exception as e:
@@ -108,7 +107,7 @@ class SentinelGCO:
         self.m_last_group_query_time = 0.0 
         while not self.m_stop_event.is_set():
             try:
-                time.sleep(0.5)
+                time.sleep(1)
                 now = time.time()
                 
                 with self.m_lock:
@@ -117,20 +116,15 @@ class SentinelGCO:
                 active_tasks = []
                 dead_tasks = []
                 for t in tasks_to_check:
-                    
-                    if not t.is_querying:
+                    if t.is_querying:
                         active_tasks.append(t)
                         continue
                     
-                    try:
-                        if t.instance.health_check():
-                            active_tasks.append(t)
-                        else:
-                            dead_tasks.append(t)
-                    except Exception as e:
+                    if t.instance.health_check():
+                        active_tasks.append(t)
+                    else:
                         dead_tasks.append(t)
-                        self._log(f"Health check failed: {e}")
-
+          
                 if dead_tasks:
                     with self.m_lock:
                         current_tasks = list(self.m_tasks)
@@ -207,30 +201,30 @@ class SentinelGCO:
 
             except Exception as e:
                 self._log(f"Monitor loop error: {e}")
-                time.sleep(2)
 
     def _creator_loop(self):
         self._log("Creator loop started.")
         group_cd_key = f"vs:group:cooldown:{self.m_cfg.identifier}"
         
         while not self.m_stop_event.is_set():
-            time.sleep(2)
-            with self.m_lock:
-                
+            try:
+                time.sleep(1)
                 if self.redis_client.exists(group_cd_key):
                     continue
+                with self.m_lock:
+                    current = len(self.m_tasks)
+                    pending = self.m_pending_builtin
+                    target = self.m_cfg.sentinel.target_instances
                 
-                current = len(self.m_tasks)
-                pending = self.m_pending_builtin
-                target = self.m_cfg.sentinel.target_instances
-            
-            if (current + pending) < target:
-                now = time.time()
-                if now - self.m_last_spawn_time >= self.m_spawn_interval:
-                    with self.m_lock:
-                        self.m_last_spawn_time = now
-                    self._log(f"Staggered spawn triggered. Next spawn in {self.m_spawn_interval:.1f}s")
-                    self._spawn_sentinel_worker()
+                if (current + pending) < target:
+                    now = time.time()
+                    if now - self.m_last_spawn_time >= self.m_spawn_interval:
+                        with self.m_lock:
+                            self.m_last_spawn_time = now
+                        self._log(f"Staggered spawn triggered. Next spawn in {self.m_spawn_interval:.1f}s")
+                        self._spawn_sentinel_worker()
+            except Exception as e:
+                self._log(f'Creator loop exception:{e}')
 
     def _spawn_sentinel_worker(self):
         with self.m_lock:
@@ -250,18 +244,11 @@ class SentinelGCO:
                     plg_cfg.account.username = "Guest"
                 else:
                     acc = VSCloudApi.Instance().get_next_account(self.m_cfg.sentinel.account_pool_id, self.m_cfg.sentinel.account_cd)
-                    plg_cfg.account.id = acc['id']
-                    plg_cfg.account.username = acc['username']
-                    plg_cfg.account.password = acc['password']
+                    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.id = proxy['id']
-                    plg_cfg.proxy.ip = proxy['ip']
-                    plg_cfg.proxy.port = proxy['port']
-                    plg_cfg.proxy.proto = proxy['proto']
-                    plg_cfg.proxy.username = proxy['username']
-                    plg_cfg.proxy.password = proxy['password']
+                    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)
@@ -272,8 +259,8 @@ class SentinelGCO:
                     self.m_tasks.append(
                         Task(instance=instance,qw_cfg=self.m_cfg.query_wait,next_run=time.time(), book_allowed=False))
                 
-                    group_fail_key = f"vs:group:failures:{self.m_cfg.identifier}"
-                    self.redis_client.delete(group_fail_key)
+                group_fail_key = f"vs:group:failures:{self.m_cfg.identifier}"
+                self.redis_client.delete(group_fail_key)
                 
                 success = True
                 self._log(f"+++ Sentinel spawned: {plg_cfg.account.username}")
@@ -305,11 +292,7 @@ class SentinelGCO:
                     
             finally:
                 if not success and instance is not None:
-                    try:
-                        if hasattr(instance, "cleanup"):
-                            instance.cleanup()
-                    except Exception as e:
-                        self._log(f"Cleanup failed after spawn failure: {e}")
+                    instance.cleanup()
                 with self.m_lock:
                     self.m_pending_builtin = max(0, self.m_pending_builtin - 1)
 

+ 7 - 7
test/test_publish_slot.py

@@ -11,22 +11,22 @@ r = redis.Redis(
 )
 
 # 使用的键名(原 Channel 名称)
-key = "vs:signal:slot.bjs.fr.tourist"
+key = "vs:signal:slot.lon.fr.tourist"
 
 # 消息体
 message = {
-    "group_id": "tls.cn.bjs.fr",
+    "group_id": "tls.gb.fr",
     "apt_type": {
         "weight": 10,
-        "routing_key": "slot.bjs.fr.tourist",
-        "city": "Beijing",
+        "routing_key": "slot.lon.fr.tourist",
+        "city": "London",
         "visa_type": "Tourist",
         "country": "France"
     },
     "query_result": {
-        "routing_key": "slot.bjs.fr.tourist",
+        "routing_key": "slot.lon.fr.tourist",
         "country": "France",
-        "city": "Beijing",
+        "city": "London",
         "visa_type": "Tourist",
         "availability_status": "Available",
         "earliest_date": "2026-06-06",
@@ -36,7 +36,7 @@ message = {
                 "times": [
                     {
                         "time": "16:30",
-                        "label": "pta"
+                        "label": ""
                     }
                 ]
             }

+ 75 - 0
utils/safe_redis_cli.py

@@ -0,0 +1,75 @@
+import time
+import redis
+from typing import List, Dict, Callable, Any
+
+
+class SafeRedisClient:
+    def __init__(self, redis_conf: Dict, logger: Callable[[str], None] = None):
+        safe_conf = dict(redis_conf)
+        safe_conf.setdefault('socket_timeout', 5.0)          # 读写超时防死等
+        safe_conf.setdefault('socket_connect_timeout', 5.0)  # 建连超时防死等
+        safe_conf.setdefault('socket_keepalive', True)       # TCP 操作系统级保活
+        safe_conf.setdefault('health_check_interval', 30)    # 自动探测断连并重连
+        
+        self._client = redis.Redis(**safe_conf)
+        self._logger = logger
+
+    def _log(self, msg: str):
+        if self._logger:
+            self._logger(f"[Safe-Redis] {msg}")
+        else:
+            print(f"[Safe-Redis] {msg}")
+
+    def _safe_execute(self, func, default_return, *args, **kwargs):
+        """通用安全执行器:遇错自动重试一次,彻底失败返回安全默认值"""
+        try:
+            return func(*args, **kwargs)
+        except Exception as e:
+            func_name = func.__name__ if hasattr(func, '__name__') else 'operation'
+            self._log(f"Attempt 1 failed ({func_name}): {e}. Retrying in 0.5s...")
+            time.sleep(0.5)
+            try:
+                return func(*args, **kwargs)
+            except Exception as e2:
+                self._log(f"Attempt 2 failed ({func_name}): {e2}. Returning safe fallback: {default_return}.")
+                return default_return
+    
+    def exists(self, name) -> bool:
+        return bool(self._safe_execute(self._client.exists, False, name))
+
+    def get(self, name) -> Any:
+        return self._safe_execute(self._client.get, None, name)
+
+    def set(self, name, value, ex=None) -> bool:
+        return self._safe_execute(self._client.set, False, name, value, ex=ex)
+
+    def setex(self, name, time_s, value) -> bool:
+        return self._safe_execute(self._client.setex, False, name, time_s, value)
+
+    def delete(self, *names) -> int:
+        if not names: return 0
+        return self._safe_execute(self._client.delete, 0, *names)
+
+    def incr(self, name, amount=1) -> int:
+        # 如果增加失败,默认返回 1,防止触发外部数学运算崩溃(同时保证退避策略能生效)
+        return self._safe_execute(self._client.incr, 1, name, amount=amount)
+
+    def zadd(self, name, mapping) -> int:
+        if not mapping: return 0
+        return self._safe_execute(self._client.zadd, 0, name, mapping)
+
+    def zrem(self, name, *values) -> int:
+        if not values: return 0
+        return self._safe_execute(self._client.zrem, 0, name, *values)
+
+    def bulk_zadd(self, name, mapping) -> bool:
+        """封装 Pipeline,用于一次性安全地写入大量 zset 元素"""
+        if not mapping:
+            return True
+        def _op():
+            pipe = self._client.pipeline()
+            for k, v in mapping.items():
+                pipe.zadd(name, {k: v})
+            pipe.execute()
+            return True
+        return self._safe_execute(_op, False)

+ 7 - 0
vs_plg.py

@@ -72,4 +72,11 @@ class IVSPlg(ABC):
         """
         @brief 设置日志输出工具, 该函数不允许抛异常
         """
+        pass
+    
+    @abstractmethod
+    def cleanup(self) -> None:
+        """
+        @brief 清理实例使用的资源
+        """
         pass