""" Filter Engine Main orchestrator for content filtering pipeline. """ import logging import traceback from typing import List, Dict, Any, Optional from datetime import datetime from pathlib import Path from concurrent.futures import ThreadPoolExecutor, as_completed from .config import FilterConfig from .cache import FilterCache from .models import FilterResult, ProcessingStatus, AIAnalysisResult from .registry import discover_modules, get_registered_stages logger = logging.getLogger(__name__) class FilterEngine: """ Main filter pipeline orchestrator. Coordinates multi-stage content filtering with intelligent caching. Compatible with user preferences and filterset selections. """ _instance = None def __init__( self, config_file: str = 'filter_config.json', filtersets_file: str = 'filtersets.json' ): """ Initialize filter engine. Args: config_file: Path to filter_config.json filtersets_file: Path to filtersets.json """ self.config = FilterConfig(config_file, filtersets_file) self.cache = FilterCache(self.config.get_cache_dir()) # Lazy-loaded stages (will be imported when AI is enabled) self._stages = None logger.info("FilterEngine initialized") logger.info(f"Configuration: {self.config.get_config_summary()}") @classmethod def get_instance(cls) -> 'FilterEngine': """Get singleton instance of FilterEngine""" if cls._instance is None: cls._instance = cls() return cls._instance def _init_stages(self): """Initialize pipeline stages from the registry (lazy loading).""" if self._stages is not None: return # Import built-ins and any configured extension modules for decorator # side effects. This keeps engine orchestration independent of concrete # stage classes and gives plugins a zero-core-edit registration path. discover_modules([ 'filter_pipeline.stages.categorizer', 'filter_pipeline.stages.moderator', 'filter_pipeline.stages.filter', 'filter_pipeline.stages.ranker', 'filter_pipeline.stages.plugins', 'filter_pipeline.stages.comment_filter', 'filter_pipeline.plugins.keyword', 'filter_pipeline.plugins.quality', *self.config.config.get('pipeline', {}).get('stage_modules', []), *self.config.config.get('plugins', {}).get('modules', []), ]) self._stages = { name: stage_cls(self.config, self.cache) for name, stage_cls in get_registered_stages().items() } logger.info( "Initialized %s registered pipeline stages: %s", len(self._stages), ', '.join(sorted(self._stages.keys())) ) def apply_filterset( self, posts: List[Dict[str, Any]], filterset_name: str = 'no_filter', use_cache: bool = True ) -> List[Dict[str, Any]]: """ Apply filterset to posts (compatible with user preferences). This is the main public API used by app.py when loading user feeds. Args: posts: List of post dictionaries filterset_name: Name of filterset from user settings (e.g., 'safe_content') use_cache: Whether to use cached results Returns: List of posts that passed the filter, with score and metadata added """ if not posts: return [] # Validate filterset exists filterset = self.config.get_filterset(filterset_name) if not filterset: logger.warning(f"Filterset '{filterset_name}' not found, using 'no_filter'") filterset_name = 'no_filter' logger.info(f"Applying filterset '{filterset_name}' to {len(posts)} posts") # Check if we have cached filterset results if use_cache and self.config.is_cache_enabled(): filterset_version = self.config.get_filterset_version(filterset_name) cached_results = self.cache.get_filterset_results( filterset_name, filterset_version, self.config.get_cache_ttl_hours() ) if cached_results: # Filter posts using cached results filtered_posts = [] for post in posts: post_uuid = post.get('uuid') if post_uuid in cached_results: result = cached_results[post_uuid] if result.passed: # Add filter metadata to post post['_filter_score'] = result.score post['_filter_categories'] = result.categories post['_filter_tags'] = result.tags filtered_posts.append(post) logger.info(f"Cache hit: {len(filtered_posts)}/{len(posts)} posts passed filter") return filtered_posts # Cache miss or disabled - process posts through pipeline results = self.process_batch(posts, filterset_name) # Save to filterset cache if self.config.is_cache_enabled(): filterset_version = self.config.get_filterset_version(filterset_name) results_dict = {r.post_uuid: r for r in results} self.cache.set_filterset_results(filterset_name, filterset_version, results_dict) # Build filtered post list filtered_posts = [] results_by_uuid = {r.post_uuid: r for r in results} for post in posts: post_uuid = post.get('uuid') result = results_by_uuid.get(post_uuid) if result and result.passed: # Add filter metadata to post post['_filter_score'] = result.score post['_filter_categories'] = result.categories post['_filter_tags'] = result.tags filtered_posts.append(post) logger.info(f"Processed: {len(filtered_posts)}/{len(posts)} posts passed filter") return filtered_posts def process_batch( self, posts: List[Dict[str, Any]], filterset_name: str = 'no_filter' ) -> List[FilterResult]: """ Process batch of posts through pipeline. Args: posts: List of post dictionaries filterset_name: Name of filterset to apply Returns: List of FilterResults for each post """ if not posts: return [] # Special case: no_filter passes everything with default scores if filterset_name == 'no_filter': return self._process_no_filter(posts) # Initialize stages (registry-driven). This must happen regardless of # whether AI is enabled: offline filtersets (rules/plugins/ranker) still # need their stages instantiated to run. self._init_stages() # Get pipeline stages for this filterset stage_names = self._get_stages_for_filterset(filterset_name) # If AI is disabled but the filterset's stages require AI, do NOT silently # pass everything as no_filter. Pass the posts through (so the feed is not # blanked) but mark every result as FAILED with an explicit error so the # degradation is observable, not silent. Filtersets whose stages are all # offline (filter/plugins/ranker/comment_filter) still run normally. if not self.config.is_ai_enabled() and self._stages_need_ai(stage_names): logger.warning( f"AI disabled but filterset '{filterset_name}' requires AI stages " f"({stage_names}) - passing posts through unfiltered with FAILED status" ) return self._process_ai_disabled(filterset_name, posts) # Process posts (parallel or sequential based on config) if self.config.is_parallel_enabled(): results = self._process_batch_parallel(posts, filterset_name, stage_names) else: results = self._process_batch_sequential(posts, filterset_name, stage_names) return results def _stages_need_ai(self, stage_names: List[str]) -> bool: """Return True if any named stage class declares ``requires_ai``.""" from .registry import get_stage_class for name in stage_names: stage_cls = get_stage_class(name) if stage_cls is not None and getattr(stage_cls, 'requires_ai', False): return True return False def _process_no_filter(self, posts: List[Dict[str, Any]]) -> List[FilterResult]: """Process posts with no_filter (all pass with default scores)""" results = [] for post in posts: result = FilterResult( post_uuid=post.get('uuid', ''), passed=True, score=0.5, # Neutral score categories=[], tags=[], filterset_name='no_filter', processed_at=datetime.now(), status=ProcessingStatus.COMPLETED ) results.append(result) return results def _process_ai_disabled(self, filterset_name: str, posts: List[Dict[str, Any]]) -> List[FilterResult]: """Pass posts through unfiltered when the requested filterset needs AI but AI is disabled. Unlike no_filter, every result is marked FAILED with an explicit error so the degradation is observable rather than silent. """ results = [] for post in posts: result = FilterResult( post_uuid=post.get('uuid', ''), passed=True, # do not blank the feed score=0.5, # neutral score categories=[], tags=[], filterset_name=filterset_name, processed_at=datetime.now(), status=ProcessingStatus.FAILED, error=f"AI disabled: filterset '{filterset_name}' requires AI; passed through unfiltered" ) results.append(result) return results def _get_stages_for_filterset(self, filterset_name: str) -> List[str]: """Get pipeline stages to run for a filterset""" filterset = self.config.get_filterset(filterset_name) # Check if filterset specifies custom stages if filterset and 'pipeline_stages' in filterset: return filterset['pipeline_stages'] # Use default stages return self.config.get_default_stages() def _process_batch_parallel( self, posts: List[Dict[str, Any]], filterset_name: str, stage_names: List[str] ) -> List[FilterResult]: """Process posts in parallel""" results = [None] * len(posts) workers = self.config.get_parallel_workers() def process_single_post(idx_post): idx, post = idx_post try: result = self._process_single_post(post, filterset_name, stage_names) return idx, result except Exception as e: logger.error(f"Error processing post {idx}: {e}") logger.error(traceback.format_exc()) # Return failed result return idx, FilterResult( post_uuid=post.get('uuid', ''), passed=False, score=0.0, filterset_name=filterset_name, processed_at=datetime.now(), status=ProcessingStatus.FAILED, error=str(e) ) with ThreadPoolExecutor(max_workers=workers) as executor: futures = {executor.submit(process_single_post, (i, post)): i for i, post in enumerate(posts)} for future in as_completed(futures): idx, result = future.result() results[idx] = result return results def _process_batch_sequential( self, posts: List[Dict[str, Any]], filterset_name: str, stage_names: List[str] ) -> List[FilterResult]: """Process posts sequentially""" results = [] for post in posts: try: result = self._process_single_post(post, filterset_name, stage_names) results.append(result) except Exception as e: logger.error(f"Error processing post: {e}") results.append(FilterResult( post_uuid=post.get('uuid', ''), passed=False, score=0.0, filterset_name=filterset_name, processed_at=datetime.now(), status=ProcessingStatus.FAILED, error=str(e) )) return results def _process_single_post( self, post: Dict[str, Any], filterset_name: str, stage_names: List[str] ) -> FilterResult: """ Process single post through pipeline stages. Stages are run sequentially: Categorizer → Moderator → Filter → Ranker """ # Initialize result result = FilterResult( post_uuid=post.get('uuid', ''), passed=True, # Start as passed, stages can reject score=0.5, # Default score filterset_name=filterset_name, processed_at=datetime.now(), status=ProcessingStatus.PROCESSING ) # Run each stage for stage_name in stage_names: if stage_name not in self._stages: logger.warning(f"Stage '{stage_name}' not found, skipping") continue stage = self._stages[stage_name] if not stage.is_enabled(): logger.debug(f"Stage '{stage_name}' disabled, skipping") continue # Process through stage try: result = stage.process(post, result) # If post was rejected by this stage, stop processing if not result.passed: logger.debug(f"Post {post.get('uuid', '')} rejected by {stage_name}") break except Exception as e: logger.error(f"Error in stage '{stage_name}': {e}") result.status = ProcessingStatus.FAILED result.error = f"{stage_name}: {str(e)}" result.passed = False break # Mark as completed if not failed if result.status != ProcessingStatus.FAILED: result.status = ProcessingStatus.COMPLETED return result # ===== Utility Methods ===== def get_available_filtersets(self) -> List[str]: """Get list of available filterset names (for user settings UI)""" return self.config.get_filterset_names() def get_filterset_description(self, name: str) -> Optional[str]: """Get description of a filterset (for user settings UI)""" filterset = self.config.get_filterset(name) return filterset.get('description') if filterset else None def invalidate_filterset_cache(self, filterset_name: str): """Invalidate cache for a filterset (when definition changes)""" self.cache.invalidate_filterset(filterset_name) logger.info(f"Invalidated cache for filterset '{filterset_name}'") def get_cache_stats(self) -> Dict[str, Any]: """Get cache statistics""" return self.cache.get_cache_stats() def reload_config(self): """Reload configuration from disk""" self.config.reload() self._stages = None # Force re-initialization of stages logger.info("Configuration reloaded") def filter_comments( self, comments: List[Dict[str, Any]], filterset_name: str = 'no_filter' ) -> List[Dict[str, Any]]: """Filter a post's flat comment list according to a filterset. Comment filtering is tree-shaped (per post) and lives in the registered ``comment_filter`` stage rather than the per-post stage pipeline. The caller (API endpoint) builds the tree from the returned flat list via ``PostService.build_comment_tree``. Fails open: if the ``comment_filter`` stage is not registered, the comments are returned unchanged so a missing stage never blanks them. """ if not comments: return [] self._init_stages() comment_stage = self._stages.get('comment_filter') if not comment_stage: logger.warning("comment_filter stage not registered; returning comments unfiltered") return comments return comment_stage.filter_comments(comments, filterset_name)