import newspaper from newspaper import Config import json import feedparser import time import os import logging import random from datetime import datetime from selenium import webdriver from selenium.webdriver.firefox.options import Options as FirefoxOptions from selenium.webdriver.common.by import By from selenium.webdriver.support.ui import WebDriverWait from selenium.webdriver.support import expected_conditions as EC from concurrent.futures import ThreadPoolExecutor, as_completed, TimeoutError import nltk from nltk.downloader import Downloader # Setup logging with timestamps logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s', datefmt='%Y-%m-%d %H:%M:%S' ) logger = logging.getLogger(__name__) FEED_FILE = os.getenv("FEED_FILE", "./rss_short_feed.json") MAX_FEED_WORKERS = int(os.getenv("MAX_FEED_WORKERS", "10")) MAX_ARTICLE_WORKERS = int(os.getenv("MAX_ARTICLE_WORKERS", "10")) BATCH_SIZE = int(os.getenv("BATCH_SIZE", "50")) # Batch processing size # Rotating User-Agents to bypass bot detection (Reuters, etc.) USER_AGENTS = [ "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "Mozilla/5.0 (Macintosh; Intel Mac OS X 14_5) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.5 Safari/605.1.15", "Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:134.0) Gecko/20100101 Firefox/134.0", "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", "Mozilla/5.0 (Macintosh; Intel Mac OS X 14_5) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/131.0.0.0 Safari/537.36", ] def get_random_ua(): return random.choice(USER_AGENTS) # Ensure necessary NLTK resources are downloaded d = Downloader() if not d.is_installed("punkt_tab"): nltk.download("punkt_tab") articles = [] # Enhanced cache system def load_processed_cache(): """Load the processed articles cache with enhanced tracking""" cache_path = "articles/processed_articles_cache.json" try: if os.path.exists(cache_path): with open(cache_path, 'r', encoding='utf-8') as f: return json.load(f) else: return {} except Exception as e: logger.error(f"Error loading cache: {e}") return {} def save_processed_cache(cache_data): """Save the processed articles cache with enhanced tracking""" cache_path = "articles/processed_articles_cache.json" try: with open(cache_path, 'w', encoding='utf-8') as f: json.dump(cache_data, f, indent=2, ensure_ascii=False) logger.info(f"Cache saved with {len(cache_data)} entries") except Exception as e: logger.error(f"Error saving cache: {e}") def is_article_processed(article_path, cache_data): """Check if an article has been processed""" return article_path in cache_data def mark_article_processed(article_path, status="completed", embedding_status="pending"): """Mark an article as processed with detailed status tracking""" cache_data = load_processed_cache() cache_data[article_path] = { "processed_date": datetime.now().isoformat(), "status": status, "embedding_status": embedding_status, "last_updated": datetime.now().isoformat() } save_processed_cache(cache_data) def get_processing_progress(): """Get overall processing progress""" cache_data = load_processed_cache() total_articles = len(cache_data) completed_articles = sum(1 for data in cache_data.values() if data.get('status') == 'completed') embedded_articles = sum(1 for data in cache_data.values() if data.get('embedding_status') == 'completed') return { "total_articles": total_articles, "completed_articles": completed_articles, "embedded_articles": embedded_articles, "completion_rate": (completed_articles / total_articles * 100) if total_articles > 0 else 0 } def load_rss_feed_sources(feed_file=FEED_FILE): """ Loads the RSS feed sources from a JSON file. """ logger.info(f"Loading RSS feed sources from {feed_file}...") try: with open(feed_file, "r", encoding="utf-8") as f: data = json.load(f) logger.info(f"Successfully loaded {feed_file}") return data except FileNotFoundError: logger.error(f"{feed_file} not found, returning empty list.") return [] except json.JSONDecodeError: logger.error(f"Error decoding {feed_file} , returning empty list.") return [] def mine_all_articles(rss_feed_sources, limit=None): """ Mines all articles from the given RSS feed sources. Returns a list of (site, title, link) tuples. """ all_links = [] # Check if rss_feed_sources is a valid dict with rss_feeds key if not isinstance(rss_feed_sources, dict): logger.warning(f"rss_feed_sources is not a dict, it's {type(rss_feed_sources)}") return all_links if "rss_feeds" not in rss_feed_sources: logger.warning("rss_feeds key not found in rss_feed_sources") return all_links sources = rss_feed_sources["rss_feeds"] # Parse RSS feeds in parallel for better performance def parse_single_feed(site, data): """Parse a single RSS feed with timeout and error handling""" try: logger.info(f"Parsing RSS feed: {data['rss_url']}") # Add more aggressive timeout settings with fallback # Use a wrapper to ensure we don't hang indefinitely import signal def timeout_handler(signum, frame): raise TimeoutError(f"Timeout parsing feed: {site}") # Set up signal-based timeout (this is a fallback for truly hanging requests) old_handler = signal.signal(signal.SIGALRM, timeout_handler) signal.alarm(10) # 10 second alarm feed = feedparser.parse(data["rss_url"], timeout=8) # 8 second timeout signal.alarm(0) # Cancel the alarm signal.signal(signal.SIGALRM, old_handler) feed_entries = feed.entries[:limit] if limit else feed.entries entries = [] for entry in feed_entries: if "link" in entry and "title" in entry: entries.append((site, entry.title, entry.link)) return entries except TimeoutError as e: logger.error(f"Timeout parsing RSS feed: {site} Error: {str(e)}") return [] except Exception as e: error_str = str(e).lower() # Handle specific network connection issues if "remote end closed connection" in error_str or "connection closed" in error_str: logger.warning(f"Network connection closed by remote end for feed: {site} - {str(e)}") logger.info(f"Skipping problematic feed: {site}") return [] elif "timeout" in error_str: logger.error(f"Timeout parsing RSS feed: {site} Error: {str(e)}") return [] else: logger.error(f"Error parsing RSS feed: {site} Error: {str(e)}") return [] # Use ThreadPoolExecutor for parallel RSS feed parsing from concurrent.futures import ThreadPoolExecutor, as_completed max_workers = min(10, len(sources)) # Limit concurrent workers with ThreadPoolExecutor(max_workers=max_workers) as executor: # Submit all feed parsing tasks future_to_site = { executor.submit(parse_single_feed, site, data): site for site, data in sources.items() } # Collect results as they complete for future in as_completed(future_to_site, timeout=30): # 30 second overall timeout try: entries = future.result() all_links.extend(entries) except Exception as e: site = future_to_site[future] logger.error(f"Error processing feed for {site}: {str(e)}") return all_links def generate_filename_from_url(url): """ Generates a filename from the given URL by replacing slashes with underscores. """ # Use only the last part of the URL or replace slashes filename = url.replace("https://", "").replace("http://", "").replace("/", "_") # Sanitize filename to remove/replace invalid characters import re filename = re.sub(r'[\\/*?:"<>|]', "_", filename) return filename def generate_safe_filename(name): # Remove/replace characters not allowed in filenames import re safe = re.sub(r'[\\/*?:"<>|]', "_", name) return safe def save_article_to_file(article, filename, source="Unfiltered"): """ Saves the given article text to a file with the specified filename. """ # articles dir should already be there # os.makedirs("articles", exist_ok=True) outputDir = "articles/" + source os.makedirs(outputDir, exist_ok=True) if source else None # Sanitize filename: use only the last part of the URL or replace slashes safe_filename = generate_filename_from_url(filename) file_path = os.path.join(outputDir, safe_filename) # Save the source as the first line in the file for later retrieval with open(file_path, "w", encoding="utf-8") as f: f.write(f"SOURCE:{source}\n") f.write(article) # Only log when a new file is actually created (not cached) logger.info(f"New article saved: {safe_filename} from {source}") def get_article_with_selenium(url): """ Gets article text using Selenium Firefox driver with proper error handling, cleanup, and bot-detection evasion. """ driver = None try: # Configure Firefox options with bot-detection evasion options = FirefoxOptions() options.add_argument("--headless") options.set_preference("dom.ipc.processCount", 1) options.set_preference("general.useragent.override", get_random_ua()) options.set_preference("permissions.default.image", 2) # Block images for speed options.set_preference("dom.webnotifications.enabled", False) # Try to initialize driver with explicit path to Firefox try: driver = webdriver.Firefox(options=options) except Exception as e: # If that fails, try with explicit Firefox path if "binary is not a firefox executable" in str(e).lower(): logger.info("Attempting to use Firefox at /usr/bin/firefox") options.binary_location = "/usr/bin/firefox" driver = webdriver.Firefox(options=options) else: raise e driver.set_page_load_timeout(30) # 30 seconds timeout # Navigate to URL driver.get(url) # Wait for page to load (explicit wait instead of sleep) try: WebDriverWait(driver, 15).until( EC.presence_of_element_located((By.TAG_NAME, "body")) ) except: pass # Continue even if wait times out time.sleep(random.uniform(1, 3)) # Random wait to mimic human behavior html = driver.page_source # Parse with Newspaper4k article = newspaper.article(url, input_html=html, language="en") article.nlp() logger.info(f"Successfully extracted article with Selenium from {url}") return article.text except Exception as e: # Check if this is a Firefox binary not found error error_str = str(e).lower() if "binary is not a firefox executable" in error_str or "firefox" in error_str: logger.error(f"Firefox not found or not properly configured for {url}: {str(e)}") logger.error("Firefox is installed at /usr/bin/firefox but may not be accessible. Check PATH or permissions.") else: logger.error(f"Selenium failed for {url}: {str(e)}") return "" finally: # Always quit the driver if driver: try: driver.quit() except: pass # Ignore errors in cleanup def get_article_with_playwright(url): """ Gets article text using Playwright with proper bot-detection evasion. """ try: from playwright.sync_api import sync_playwright with sync_playwright() as p: # Use Chromium with full browser context for UA spoofing browser = p.chromium.launch(headless=True, timeout=30000) context = browser.new_context( user_agent=get_random_ua(), viewport={"width": 1920, "height": 1080}, locale="en-US", timezone_id="America/New_York", ) page = context.new_page() # Additional headers for legitimacy page.set_extra_http_headers({ "Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", "Accept-Language": "en-US,en;q=0.9", "Accept-Encoding": "gzip, deflate, br", "Connection": "keep-alive", "Upgrade-Insecure-Requests": "1", }) page.goto(url, wait_until="domcontentloaded", timeout=30000) # Wait for content to load time.sleep(random.uniform(2, 4)) html = page.content() context.close() browser.close() # Parse with Newspaper4k article = newspaper.article(url, input_html=html, language="en") article.nlp() logger.info(f"Successfully extracted article with Playwright from {url}") return article.text except Exception as e: logger.error(f"Playwright failed for {url}: {str(e)}") return "" def pull_article(link, source, title=None, save_to_file=True): """ Pulls an article from a given link with fallback mechanisms. """ filename = title if title else link safe_filename = generate_filename_from_url(filename) # Check if already cached if os.path.exists(os.path.join("articles", source, safe_filename)): logger.info(f"Article already cached: {filename}") with open( os.path.join("articles", source, safe_filename), "r", encoding="utf-8" ) as f: return f.read() # Random delay before fetching to avoid rate-limiting / bot detection time.sleep(random.uniform(0.5, 2)) text = "" try: # Try newspaper4k first with proper User-Agent to bypass bot detection ua = get_random_ua() config = Config(browser_user_agent=ua) article = newspaper.article(link, browser_user_agent=ua) article.download() article.parse() text = article.text if not text or len(text) < 200: raise ValueError( "\tArticle text too short, falling back to Playwright/Selenium." ) logger.info(f"Successfully pulled article with newspaper4k from {link}") except Exception as e: logger.warning( f"newspaper4k extraction failed for {link}: {e}, falling back to Playwright." ) try: text = get_article_with_playwright(link) logger.info(f"Successfully pulled article from {link} with Playwright") if not text or len(text) < 200: logger.warning(f"Playwright article too short, falling back to Selenium.") # Fallback to Selenium with better error handling text = get_article_with_selenium(link) logger.info(f"Successfully pulled article from {link} with Selenium") except Exception as e: logger.error(f"Playwright failed for {link}: {e}") # Fallback to Selenium try: text = get_article_with_selenium(link) logger.info(f"Successfully pulled article from {link} with Selenium") except Exception as e: logger.error(f"Selenium failed for {link}: {e}") return "" if save_to_file: save_article_to_file(text, filename, source) return text def is_article_downloaded(source, title, link): """ Check if an article is already downloaded by checking if the file exists. """ # Generate the same filename that would be used for saving filename = title if title else link safe_filename = generate_filename_from_url(filename) # Check if file exists in the articles directory file_path = os.path.join("articles", source, safe_filename) return os.path.exists(file_path) def safe_pull_articles(article_list): """ Safely pull articles with improved error handling and increased parallelism. """ if not article_list: return [], [] # Filter out articles that are already downloaded filtered_article_list = [] total_articles = len(article_list) for source, title, link in article_list: if not is_article_downloaded(source, title, link): filtered_article_list.append((source, title, link)) else: logger.info(f"Skipping already downloaded article: {title[:50]}... from {source}") logger.info(f"Filtered out {total_articles - len(filtered_article_list)} articles that were already downloaded") logger.info(f"Processing {len(filtered_article_list)} remaining articles") if not filtered_article_list: logger.info("No new articles to process") return [], [] results = [] errors = [] # Use ThreadPoolExecutor for parallel article pulling with higher concurrency # Use configurable worker setting max_workers = min(MAX_ARTICLE_WORKERS, len(filtered_article_list)) # Cap at configured workers, but don't exceed article count batch_size = max(1, min(20, len(filtered_article_list) // 4)) # Dynamic batch size logger.info(f"Starting parallel article pulling with {max_workers} workers and batch size {batch_size}") # Process all articles in parallel with proper error handling with ThreadPoolExecutor(max_workers=max_workers) as executor: # Submit all tasks at once for maximum parallelism futures = [ executor.submit(pull_article, link, source, title) for source, title, link in filtered_article_list ] # Collect results as they complete for i, future in enumerate(as_completed(futures, timeout=300)): # 5 minute timeout total try: result = future.result(timeout=120) # 2 minute timeout per article if result: # Only count non-empty results results.append(result) # Log progress every 100 articles with proper batch information if (i + 1) % 100 == 0: logger.info(f"Processed {i + 1} articles out of {len(filtered_article_list)}") except Exception as e: logger.error(f"Error processing article: {e}") errors.append(e) return results, errors def main(): """ Main scraping loop. """ while True: logger.info("=========================================") logger.info("Starting new scraping iteration...") try: # Pull the RSS feed sources from the JSON file rss_feed_sources = load_rss_feed_sources() # Mine all articles from the RSS feed sources rss_feed_links = mine_all_articles(rss_feed_sources) # Randomize the order of the links to help with load balancing import random random.shuffle(rss_feed_links) logger.info(f"Found {len(rss_feed_links)} articles to process") if not rss_feed_links: logger.info("No articles found, sleeping for 15 minutes") time.sleep(15 * 60) continue # Process articles with better error handling and resource management results, errors = safe_pull_articles(rss_feed_links) logger.info(f"Attempted to Pull {len(results)} articles in parallel.") logger.info( f"Encountered {len(errors)} errors during article pulling. " + "Outputting errors to a local file." ) # Output errors to a local file if errors: with open("errors.txt", "w", encoding="utf-8") as f: for error in errors: f.write(str(error) + "\n") logger.info(f"Errors logged to errors.txt") # Print all results to a log file with open("results.txt", "w", encoding="utf-8") as f: for result in results: f.write(result + "\n") logger.info("All articles pulled successfully.") # Log processing progress progress = get_processing_progress() logger.info(f"Processing progress - Total: {progress['total_articles']}, " f"Completed: {progress['completed_articles']}, " f"Embedded: {progress['embedded_articles']}, " f"Completion rate: {progress['completion_rate']:.1f}%") except Exception as e: logger.error(f"Major error in main loop: {e}") # Continue to next iteration even if there's a major error # Sleep for a while before the next iteration logger.info("Sleeping for 15 minutes before the next iteration...") time.sleep(15 * 60) # Sleep for 15 minutes if __name__ == "__main__": main()