"""渲染引擎 Feature Flag 解析器。 封装渲染引擎选择逻辑,支持: - 环境变量作为默认值(RENDER_ENGINE=legacy/unified) - Redis Feature Flag 运行时覆盖(白名单 + 百分比 + 全局开关) - 定时刷新,支持热更新不重启 worker 使用方式: resolver = RenderEngineResolver(redis_url="redis://...", default_engine="legacy") engine = resolver.get_engine(user_id="user123") # engine: "legacy" 或 "unified" """ from __future__ import annotations import logging import threading from typing import Optional from packages.adapters.redis.feature_flag_store import ( FeatureFlagConfig, FeatureFlagStore, InMemoryFeatureFlagStore, RedisFeatureFlagStore, ) logger = logging.getLogger(__name__) # Feature Flag 名称常量 FLAG_RENDER_ENGINE = "render_engine" # 引擎常量 ENGINE_LEGACY = "legacy" ENGINE_UNIFIED = "unified" VALID_ENGINES = {ENGINE_LEGACY, ENGINE_UNIFIED} class RenderEngineResolver: """渲染引擎选择器。 判定逻辑(从高到低): 1. Redis flag 白名单匹配 → unified 2. Redis flag 百分比命中 → unified 3. Redis flag 全局开启(100%)→ unified 4. 环境变量默认值 → legacy / unified 当 Redis 不可用时,自动降级到环境变量默认值,不影响业务。 """ def __init__( self, default_engine: str = ENGINE_LEGACY, redis_url: Optional[str] = None, refresh_interval: float = 30.0, store: Optional[FeatureFlagStore] = None, ) -> None: """ Args: default_engine: 环境变量默认的引擎名(legacy / unified) redis_url: Redis 连接 URL,传 None 时使用内存实现(测试用) refresh_interval: Redis flag 配置刷新间隔(秒) store: 直接传入 store 实例(测试用,优先级高于 redis_url) """ self._default_engine = default_engine.lower() if default_engine else ENGINE_LEGACY if self._default_engine not in VALID_ENGINES: logger.warning( "Invalid default engine '%s', fallback to '%s'", self._default_engine, ENGINE_LEGACY, ) self._default_engine = ENGINE_LEGACY if store is not None: self._store = store elif redis_url: self._store = RedisFeatureFlagStore(redis_url=redis_url) else: self._store = InMemoryFeatureFlagStore() logger.info("No Redis configured, using in-memory feature flag store") self._refresh_interval = refresh_interval self._lock = threading.Lock() self._cached_config: Optional[FeatureFlagConfig] = None self._last_refresh: float = 0.0 def _maybe_refresh(self) -> None: """惰性刷新配置,超过刷新间隔时从存储重新读取。""" import time now = time.time() if now - self._last_refresh < self._refresh_interval: return try: config = self._store.get(FLAG_RENDER_ENGINE) with self._lock: self._cached_config = config self._last_refresh = now except Exception as exc: logger.warning("Failed to refresh render engine flag: %s", exc) # 刷新失败时保留旧缓存,不中断业务 if self._cached_config is None: # 首次就读失败,设一个默认值 with self._lock: self._cached_config = FeatureFlagConfig(name=FLAG_RENDER_ENGINE) self._last_refresh = now def _get_config(self) -> FeatureFlagConfig: """获取当前 flag 配置(带缓存)。""" if self._cached_config is None: self._maybe_refresh() else: self._maybe_refresh() return self._cached_config or FeatureFlagConfig(name=FLAG_RENDER_ENGINE) def get_engine(self, user_id: Optional[str] = None) -> str: """获取当前应该使用的渲染引擎。 Args: user_id: 用户ID,用于白名单匹配和百分比哈希。 传 None 时只看全局开关。 Returns: "legacy" 或 "unified" """ config = self._get_config() # 全局关闭 → 用默认值 if not config.enabled: return self._default_engine # 白名单匹配 / 百分比命中 → unified if config.is_active(user_id): return ENGINE_UNIFIED # 未命中灰度 → 用默认值 return self._default_engine def should_use_unified(self, user_id: Optional[str] = None) -> bool: """便捷方法:是否应该使用统一渲染引擎。""" return self.get_engine(user_id) == ENGINE_UNIFIED def force_refresh(self) -> None: """强制立即刷新配置(用于管理接口修改后立即生效)。""" self._last_refresh = 0.0 if isinstance(self._store, RedisFeatureFlagStore): self._store.invalidate_cache(FLAG_RENDER_ENGINE) self._maybe_refresh() def get_config_snapshot(self) -> dict: """获取当前配置快照(用于管理接口展示)。""" config = self._get_config() return { "flag_name": FLAG_RENDER_ENGINE, "default_engine": self._default_engine, "enabled": config.enabled, "percentage": config.percentage, "whitelist": sorted(config.whitelist), "refresh_interval": self._refresh_interval, "last_refresh": self._last_refresh, } def set_flag(self, config: FeatureFlagConfig) -> None: """设置 flag 配置(管理接口用)。""" config.name = FLAG_RENDER_ENGINE self._store.set(config) self.force_refresh() # 全局单例 _resolver: Optional[RenderEngineResolver] = None _resolver_lock = threading.Lock() def get_render_engine_resolver() -> RenderEngineResolver: """获取全局单例(基于 worker 配置)。""" global _resolver if _resolver is not None: return _resolver with _resolver_lock: if _resolver is not None: return _resolver try: from worker_app.core.config import get_settings settings = get_settings() redis_url = getattr(settings, "redis_url", None) or getattr(settings, "broker_url", None) default = getattr(settings, "render_engine", ENGINE_LEGACY) _resolver = RenderEngineResolver( default_engine=default, redis_url=redis_url, ) logger.info( "RenderEngineResolver initialized: default=%s, redis=%s", default, bool(redis_url), ) except Exception as exc: logger.warning("Failed to init RenderEngineResolver from settings: %s", exc) _resolver = RenderEngineResolver(default_engine=ENGINE_LEGACY) return _resolver