""" Article processor for handling the processing of articles from the scraper directory. Manages batching, processing, and integration with the fact extraction system. """ import os import json import logging import time from datetime import datetime from typing import List, Tuple from fact_extractor import FactExtractor from cache_manager import CacheManager from metrics_collector import metrics_collector from config import BATCH_SIZE, CACHE_FILE logger = logging.getLogger(__name__) class ArticleProcessor: """Processes articles from the scraper directory and extracts facts.""" def __init__(self): self.fact_extractor = FactExtractor() self.cache_manager = CacheManager(CACHE_FILE) self.batch_size = BATCH_SIZE def find_unprocessed_articles(self, scraper_dir: str = "/app/articles") -> List[Tuple[str, str]]: """ Find all unprocessed articles in the scraper directory. Args: scraper_dir (str): Path to the scraper articles directory Returns: List of tuples (file_path, filename) """ unprocessed_articles = [] try: # Debug: Check if directory exists logger.info(f"Checking for articles in: {scraper_dir}") if not os.path.exists(scraper_dir): logger.warning(f"Scraper directory does not exist: {scraper_dir}") return [] logger.info(f"Directory exists, walking through files...") for root, dirs, files in os.walk(scraper_dir): for file in files: if file.endswith('.json'): file_path = os.path.join(root, file) if not self.cache_manager.is_processed(file_path): unprocessed_articles.append((file_path, file)) logger.info(f"Found {len(unprocessed_articles)} unprocessed articles") return unprocessed_articles except Exception as e: logger.error(f"Error finding unprocessed articles: {e}") return [] def process_article_file(self, file_path: str, filename: str) -> dict: """ Process a single article file and extract facts. Args: file_path (str): Path to the article file filename (str): Name of the article file Returns: Dictionary containing the extracted facts or None if failed """ try: with open(file_path, 'r', encoding='utf-8') as f: article_data = json.load(f) # Extract facts from the article facts = self.fact_extractor.extract_facts_from_article( article_data.get('original_content', ''), article_data.get('title', filename) ) # Add metadata facts['source'] = article_data.get('source', 'Unknown') facts['published'] = article_data.get('published', 'Unknown') facts['filename'] = filename facts['processed_at'] = datetime.now().isoformat() # Mark as processed in cache self.cache_manager.mark_processed(file_path) logger.info(f"Successfully processed article: {filename}") return facts except Exception as e: logger.error(f"Error processing article {filename}: {e}") metrics_collector.increment_articles_failed() return None def process_batch(self, articles_batch: List[Tuple[str, str]]) -> Tuple[int, int]: """ Process a batch of articles. Args: articles_batch (List): List of (file_path, filename) tuples Returns: Tuple of (successful_count, failed_count) """ successful = 0 failed = 0 logger.info(f"Processing batch of {len(articles_batch)} articles") for file_path, filename in articles_batch: try: facts = self.process_article_file(file_path, filename) if facts: successful += 1 metrics_collector.increment_articles_processed() else: failed += 1 metrics_collector.increment_articles_failed() except Exception as e: logger.error(f"Error processing batch item {filename}: {e}") failed += 1 metrics_collector.increment_articles_failed() logger.info(f"Batch completed: {successful} successful, {failed} failed") return successful, failed def process_all_articles(self, scraper_dir: str = "/scraper/articles") -> dict: """ Process all unprocessed articles in the scraper directory. Args: scraper_dir (str): Path to the scraper articles directory Returns: Dictionary with processing statistics """ metrics_collector.start_processing() start_time = time.time() # Find all unprocessed articles unprocessed_articles = self.find_unprocessed_articles(scraper_dir) if not unprocessed_articles: logger.info("No unprocessed articles found") metrics_collector.stop_processing() return { "total_processed": 0, "total_failed": 0, "duration": 0, "status": "no_new_articles" } logger.info(f"Starting to process {len(unprocessed_articles)} articles in batches of {self.batch_size}") total_processed = 0 total_failed = 0 # Process articles in batches for i in range(0, len(unprocessed_articles), self.batch_size): batch = unprocessed_articles[i:i + self.batch_size] successful, failed = self.process_batch(batch) total_processed += successful total_failed += failed # Add a small delay between batches to prevent overwhelming the system if i + self.batch_size < len(unprocessed_articles): time.sleep(0.1) end_time = time.time() duration = end_time - start_time metrics_collector.stop_processing() metrics_collector.record_processing_time(duration) stats = { "total_processed": total_processed, "total_failed": total_failed, "duration": duration, "batch_size": self.batch_size, "cache_stats": self.cache_manager.get_cache_stats(), "status": "completed" } logger.info(f"Processing completed in {duration:.2f} seconds") logger.info(f"Total processed: {total_processed}, Total failed: {total_failed}") return stats def process_new_articles(self, scraper_dir: str = "/scraper/articles") -> dict: """ Process only new articles (those that haven't been processed yet). This is designed for real-time processing of new articles. Args: scraper_dir (str): Path to the scraper articles directory Returns: Dictionary with processing statistics """ return self.process_all_articles(scraper_dir)