From 492b30a7134296b8053ded423704f9e4b369dec5 Mon Sep 17 00:00:00 2001 From: Jarian Cottingham Date: Mon, 2 Feb 2026 10:35:55 -0600 Subject: [PATCH] Fix cache manager to properly handle empty cache files and add diagnostic logging --- .../processed_articles_cache.json | 0 ai_processor/article_processor.py | 170 +++++++++++------- ai_processor/cache_manager.py | 128 +++++++++---- ai_processor/processed_articles_cache.json | 5 - 4 files changed, 199 insertions(+), 104 deletions(-) create mode 100644 ai_processor/ai_processor/processed_articles_cache.json delete mode 100644 ai_processor/processed_articles_cache.json diff --git a/ai_processor/ai_processor/processed_articles_cache.json b/ai_processor/ai_processor/processed_articles_cache.json new file mode 100644 index 0000000..e69de29 diff --git a/ai_processor/article_processor.py b/ai_processor/article_processor.py index 4b8ec5c..eacf16b 100644 --- a/ai_processor/article_processor.py +++ b/ai_processor/article_processor.py @@ -3,40 +3,43 @@ Article processor for handling the processing of articles from the scraper direc Manages batching, processing, and integration with the fact extraction system. """ -import os import json import logging +import os 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 +from fact_extractor import FactExtractor +from metrics_collector import metrics_collector 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 = "../scraper/articles") -> List[Tuple[str, str]]: + + def find_unprocessed_articles( + self, scraper_dir: str = "../scraper/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}") @@ -47,92 +50,137 @@ class ArticleProcessor: "../scraper/articles", "/scraper/articles", "/app/articles", - "/articles" + "/articles", ] for alt_path in alternative_paths: if os.path.exists(alt_path): - logger.info(f"Found articles directory at alternative path: {alt_path}") + logger.info( + f"Found articles directory at alternative path: {alt_path}" + ) scraper_dir = alt_path break else: return [] - + + # Log cache state before scanning + cache_stats = self.cache_manager.get_cache_stats() + logger.info( + f"Cache state before scanning: {cache_stats['processed_files']} files marked as processed" + ) + logger.info(f"Directory exists, walking through files...") file_count = 0 + article_file_count = 0 + already_processed_count = 0 + extension_counts = {} + for root, dirs, files in os.walk(scraper_dir): for file in files: file_count += 1 + + # Track file extensions for debugging + ext = ( + os.path.splitext(file)[1].lower() + if "." in file + else "(no extension)" + ) + extension_counts[ext] = extension_counts.get(ext, 0) + 1 + # Check for JSON files (expected format) - if file.endswith('.json'): + if file.endswith(".json"): + article_file_count += 1 file_path = os.path.join(root, file) if not self.cache_manager.is_processed(file_path): unprocessed_articles.append((file_path, file)) + else: + already_processed_count += 1 # Also check for text files (fallback for different formats) - elif file.endswith(('.txt', '.md')): + elif file.endswith((".txt", ".md")): + article_file_count += 1 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"Scanned {file_count} files, found {len(unprocessed_articles)} unprocessed articles") + else: + already_processed_count += 1 + + # Log detailed breakdown + logger.info(f"File extension breakdown: {extension_counts}") + logger.info( + f"Scanned {file_count} total files, {article_file_count} are article files (.json/.txt/.md)" + ) + logger.info( + f"Already processed (in cache): {already_processed_count}, Unprocessed: {len(unprocessed_articles)}" + ) + + if ( + article_file_count > 0 + and len(unprocessed_articles) == 0 + and already_processed_count == article_file_count + ): + logger.warning( + f"All {article_file_count} article files are marked as processed in cache. " + f"If cache should be empty, check cache file: {self.cache_manager.cache_file}" + ) + 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: + 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) + 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() - + 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) @@ -146,27 +194,27 @@ class ArticleProcessor: 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() @@ -174,53 +222,55 @@ class ArticleProcessor: "total_processed": 0, "total_failed": 0, "duration": 0, - "status": "no_new_articles" + "status": "no_new_articles", } - - logger.info(f"Starting to process {len(unprocessed_articles)} articles in batches of {self.batch_size}") - + + 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] + 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" + "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 """ diff --git a/ai_processor/cache_manager.py b/ai_processor/cache_manager.py index bb26a18..364eca1 100644 --- a/ai_processor/cache_manager.py +++ b/ai_processor/cache_manager.py @@ -3,8 +3,8 @@ Cache manager for tracking processed articles to avoid reprocessing. """ import json -import os import logging +import os from datetime import datetime from typing import Dict, List, Optional @@ -12,59 +12,109 @@ from config import CACHE_FILE logger = logging.getLogger(__name__) + class CacheManager: """Manages caching of processed articles to prevent duplicate processing.""" - + def __init__(self, cache_file: str = CACHE_FILE): self.cache_file = cache_file self.cache = self._load_cache() - + def _load_cache(self) -> Dict: """Load cache from file.""" + empty_cache = { + "cache_version": "1.0", + "created": datetime.now().isoformat(), + "processed_files": {}, + } + try: - if os.path.exists(self.cache_file): - with open(self.cache_file, 'r', encoding='utf-8') as f: - cache_data = json.load(f) - # Ensure the cache has the correct structure - if "processed_files" not in cache_data: - cache_data["processed_files"] = {} - return cache_data - else: - # Create empty cache file if it doesn't exist - cache_data = { - "cache_version": "1.0", - "created": datetime.now().isoformat(), - "processed_files": {} - } - # Set the cache attribute directly - self.cache = cache_data - self._save_cache() - return cache_data + if not os.path.exists(self.cache_file): + logger.info( + f"Cache file does not exist: {self.cache_file}. Creating new empty cache." + ) + return empty_cache + + # Check if the file is empty (0 bytes) + file_size = os.path.getsize(self.cache_file) + if file_size == 0: + logger.warning( + f"Cache file is empty (0 bytes): {self.cache_file}. Treating as fresh cache." + ) + return empty_cache + + with open(self.cache_file, "r", encoding="utf-8") as f: + content = f.read().strip() + + # Check if content is empty or just whitespace + if not content: + logger.warning( + f"Cache file contains only whitespace: {self.cache_file}. Treating as fresh cache." + ) + return empty_cache + + cache_data = json.loads(content) + + # Handle case where JSON loaded as None, list, or other non-dict type + if not isinstance(cache_data, dict): + logger.warning( + f"Cache file contains invalid data type ({type(cache_data).__name__}). Treating as fresh cache." + ) + return empty_cache + + # Ensure the cache has the correct structure + if ( + "processed_files" not in cache_data + or cache_data["processed_files"] is None + ): + cache_data["processed_files"] = {} + + # Validate that processed_files is a dict + if not isinstance(cache_data["processed_files"], dict): + logger.warning( + f"processed_files is not a dict ({type(cache_data['processed_files']).__name__}). Resetting to empty." + ) + cache_data["processed_files"] = {} + + logger.info( + f"Loaded cache with {len(cache_data['processed_files'])} processed files" + ) + return cache_data + + except json.JSONDecodeError as e: + logger.warning( + f"Cache file contains invalid JSON: {e}. Treating as fresh cache." + ) + return empty_cache except Exception as e: logger.error(f"Error loading cache file {self.cache_file}: {e}") logger.info("Creating new empty cache due to load error") - # Return empty cache on error - return { - "cache_version": "1.0", - "created": datetime.now().isoformat(), - "processed_files": {} - } - + return empty_cache + def _save_cache(self) -> None: """Save cache to file.""" try: # Create directory if it doesn't exist os.makedirs(os.path.dirname(self.cache_file), exist_ok=True) - - with open(self.cache_file, 'w', encoding='utf-8') as f: + + with open(self.cache_file, "w", encoding="utf-8") as f: json.dump(self.cache, f, indent=2, ensure_ascii=False) except Exception as e: logger.error(f"Error saving cache file {self.cache_file}: {e}") - + def is_processed(self, file_path: str) -> bool: """Check if a file has been processed.""" - return file_path in self.cache.get("processed_files", {}) - + processed_files = self.cache.get("processed_files") + + # Handle case where processed_files is None or not a dict + if not isinstance(processed_files, dict): + logger.warning( + f"processed_files is not a valid dict (got {type(processed_files).__name__}). Returning False." + ) + return False + + return file_path in processed_files + def mark_processed(self, file_path: str, status: str = "processed") -> None: """Mark a file as processed.""" # Ensure we're using the correct cache structure @@ -73,24 +123,24 @@ class CacheManager: self.cache["processed_files"][file_path] = { "processed_date": datetime.now().isoformat(), "status": status, - "last_updated": datetime.now().isoformat() + "last_updated": datetime.now().isoformat(), } self._save_cache() - + def get_processed_files(self) -> List[str]: """Get list of all processed files.""" return list(self.cache.get("processed_files", {}).keys()) - + def get_cache_stats(self) -> Dict: """Get cache statistics.""" return { "total_files": len(self.cache.get("processed_files", {})), "processed_files": len(self.cache.get("processed_files", {})), - "cache_file": self.cache_file + "cache_file": self.cache_file, } - + def clear_cache(self) -> None: """Clear the entire cache.""" self.cache = {} self._save_cache() - logger.info("Cache cleared") \ No newline at end of file + logger.info("Cache cleared") diff --git a/ai_processor/processed_articles_cache.json b/ai_processor/processed_articles_cache.json deleted file mode 100644 index 13cf650..0000000 --- a/ai_processor/processed_articles_cache.json +++ /dev/null @@ -1,5 +0,0 @@ -{ - "cache_version": "1.0", - "created": "2026-02-02T10:02:48.000Z", - "processed_files": {} -} \ No newline at end of file