NewsArchiverV2/rss_processor.py
Jarian Cottingham a7936b8b11 fix: OPE hardening - logging, Docker, deps, scheduler, cleanup
- #12: Remove duplicate logging.basicConfig() from 10 modules
- #15: Remove redundant import re in rebuild_database.py
- #17: rglob('*') → rglob('*.html/json/txt/xml/md') for speed
- #18: Dockerfile individual COPY → glob COPY *.py/*.json + .dockerignore
- #19: Remove deprecated docker-compose version field
- #20: Pin requirements.txt versions (flask, requests, etc.)
- #22: SIGALRM → threading.Timer for multi-threaded safety
- #23: AP regex parsing → BeautifulSoup selectors
2026-07-05 04:14:05 +00:00

455 lines
14 KiB
Python

#!/usr/bin/env python3
"""RSS Feed Processor for NewsArchiver - Phase 2.1"""
import argparse
import json
import logging
import sys
import time
import urllib.error
import urllib.request
from datetime import datetime
from pathlib import Path
from typing import List, Optional
try:
import feedparser
from feedparser import FeedParserDict
except ImportError:
print("ERROR: feedparser is required. Install with: pip install feedparser")
sys.exit(1)
try:
import sqlite3
except ImportError:
print("ERROR: sqlite3 is required (should be built-in)")
sys.exit(1)
SCRIPT_DIR = Path(__file__).parent
ARCHIVE_DIR = SCRIPT_DIR / 'archival_data'
ARCHIVE_DIR.mkdir(exist_ok=True)
logger = logging.getLogger(__name__)
def init_db(db_path: Path = ARCHIVE_DIR / 'cache.db') -> None:
"""Initialize SQLite database for caching processed articles."""
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
cursor.execute('''
CREATE TABLE IF NOT EXISTS articles (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source_name TEXT NOT NULL,
article_url TEXT NOT NULL UNIQUE,
article_guid TEXT,
title TEXT,
author TEXT,
publish_date TEXT,
content_text TEXT,
content_html TEXT,
archive_file_path TEXT,
metadata_file_path TEXT,
status TEXT DEFAULT 'pending',
error_message TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
''')
cursor.execute('''
CREATE TABLE IF NOT EXISTS processing_log (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source_name TEXT,
timestamp TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
action TEXT,
status TEXT,
message TEXT
)
''')
cursor.execute('''
CREATE INDEX IF NOT EXISTS idx_articles_source_url
ON articles(source_name, article_url)
''')
cursor.execute('''
CREATE INDEX IF NOT EXISTS idx_articles_status
ON articles(status)
''')
conn.commit()
conn.close()
logger.debug("Database initialized: %s", db_path)
def get_db_connection(db_path: Path = ARCHIVE_DIR / 'cache.db'):
"""Get database connection."""
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
return conn
def is_duplicate(article_url: str, source_name: str, db_path: Path = ARCHIVE_DIR / 'cache.db') -> bool:
"""Check if article already in cache."""
conn = get_db_connection(db_path)
cursor = conn.cursor()
cursor.execute(
'SELECT 1 FROM articles WHERE source_name = ? AND article_url = ?',
(source_name, article_url)
)
result = cursor.fetchone()
conn.close()
return result is not None
def add_to_cache(
article_url: str,
source_name: str,
timestamp: datetime,
article_guid: str = None,
title: str = None,
author: str = None,
publish_date: str = None,
db_path: Path = ARCHIVE_DIR / 'cache.db'
) -> bool:
"""Add article to cache."""
conn = get_db_connection(db_path)
cursor = conn.cursor()
try:
cursor.execute('''
INSERT OR IGNORE INTO articles
(source_name, article_url, article_guid, title, author, publish_date, status)
VALUES (?, ?, ?, ?, ?, ?, ?)
''', (
source_name,
article_url,
article_guid,
title,
author,
publish_date,
'pending'
))
conn.commit()
conn.close()
return True
except sqlite3.IntegrityError:
conn.close()
return False
def save_article(
source_name: str,
article_data: dict,
db_path: Path = ARCHIVE_DIR / 'cache.db'
) -> str:
"""Save article to storage and update cache."""
conn = get_db_connection(db_path)
cursor = conn.cursor()
article_url = article_data.get('link', '')
article_guid = article_data.get('id', article_url)
title = article_data.get('title', '')
author = article_data.get('author', '')
publish_date = article_data.get('published', '')
summary = article_data.get('summary', '')
cursor.execute('''
INSERT OR REPLACE INTO articles
(source_name, article_url, article_guid, title, author, publish_date, content_text, status)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
''', (
source_name,
article_url,
article_guid,
title,
author,
publish_date,
summary,
'archived'
))
conn.commit()
conn.close()
return article_url
def fetch_rss_feed(rss_url: str, timeout: int = 30) -> FeedParserDict:
"""Fetch and parse RSS feed."""
logger.info("Fetching RSS feed: %s", rss_url[:50] + "..." if len(rss_url) > 50 else rss_url)
try:
req = urllib.request.Request(
rss_url,
headers={'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'}
)
with urllib.request.urlopen(req, timeout=timeout) as response:
feed_content = response.read()
feed = feedparser.parse(feed_content)
if feed.bozo:
logger.warning("Feed parsing completed with warnings: %s", feed.bozo)
entry_count = len(feed.entries)
logger.info("Feed parsed: %d entries", entry_count)
return feed
except urllib.error.HTTPError as e:
logger.error("HTTP error fetching RSS feed %s: %s", rss_url, str(e.code))
raise
except urllib.error.URLError as e:
logger.error("URL error fetching RSS feed %s: %s", rss_url, str(e.reason))
raise
except Exception as e:
logger.error("Failed to fetch RSS feed %s: %s", rss_url, str(e))
raise
def decode_google_news_url(google_url: str, entry: dict = None) -> Optional[str]:
"""Decode Google News encrypted URL to actual article URL.
Google News RSS uses encrypted URLs like:
https://news.google.com/rss/articles/CBMioAFB... which need to be decoded.
Args:
google_url: Google News encrypted URL
entry: Full RSS entry for additional context
Returns:
Decoded article URL or None if not a Google News URL
"""
if 'news.google.com' not in google_url:
return google_url
try:
import base64
import urllib.parse
if '/rss/articles/' in google_url:
parts = google_url.split('/rss/articles/')
if len(parts) >= 2:
encoded = parts[1]
if encoded.startswith('CBM'):
encoded = encoded[3:]
padding = (4 - len(encoded) % 4) % 4
encoded += '=' * padding
try:
decoded = base64.urlsafe_b64decode(encoded).decode('utf-8')
logger.debug("Decoded Google News URL: %s -> %s", google_url[:60], decoded[:60])
return decoded
except Exception:
pass
if '/articles/' in google_url:
parts = google_url.split('/articles/')
if len(parts) >= 2:
encoded = parts[1]
if encoded.startswith('CBM'):
encoded = encoded[3:]
padding = (4 - len(encoded) % 4) % 4
encoded += '=' * padding
try:
decoded = base64.urlsafe_b64decode(encoded).decode('utf-8')
logger.debug("Decoded Google News URL: %s -> %s", google_url[:60], decoded[:60])
return decoded
except Exception:
pass
if 'url=' in google_url:
parsed = urllib.parse.urlparse(google_url)
params = urllib.parse.parse_qs(parsed.query)
if 'url' in params:
return params['url'][0]
logger.debug("Could not decode Google News URL: %s", google_url[:60])
if entry and 'source' in entry:
source = entry.get('source', {})
if isinstance(source, dict) and 'href' in source:
source_href = source['href']
logger.debug("Using source URL as fallback: %s", source_href)
return source_href
return google_url
except Exception as e:
logger.debug("Failed to decode Google News URL %s: %s", google_url[:60], str(e))
return google_url
def process_rss_feed(
rss_url: str,
source_name: str,
output_dir: Path,
db_path: Path = ARCHIVE_DIR / 'cache.db'
) -> List[dict]:
"""Process RSS feed and archive new articles."""
logger.info("Processing RSS feed for source: %s", source_name)
try:
feed = fetch_rss_feed(rss_url)
except Exception as e:
logger.error("Failed to fetch feed: %s", str(e))
return []
new_articles = []
skipped_count = 0
for entry in feed.entries:
article_url = entry.get('link', '')
article_url = decode_google_news_url(article_url, entry)
article_guid = entry.get('id', article_url)
if not article_url:
logger.warning("Skipping entry without URL")
continue
if is_duplicate(article_url, source_name, db_path):
logger.debug("Skipping duplicate: %s", article_url[:60])
skipped_count += 1
continue
title = entry.get('title', 'No Title')
author = entry.get('author', entry.get('authors', [{}])[0].get('name', '') if entry.get('authors') else '')
published = entry.get('published', entry.get('published_parsed', ''))
summary = entry.get('summary', entry.get('description', ''))
if published:
if hasattr(published, 'tm_year'):
publish_date = datetime(*published[:6]).isoformat()
else:
publish_date = published
else:
publish_date = datetime.now().isoformat()
article_data = {
'source_name': source_name,
'url': article_url,
'guid': article_guid,
'title': title,
'author': author,
'publish_date': publish_date,
'summary': summary,
'entry': entry
}
add_to_cache(
article_url=article_url,
source_name=source_name,
timestamp=datetime.now(),
article_guid=article_guid,
title=title,
author=author,
publish_date=publish_date,
db_path=db_path
)
new_articles.append(article_data)
logger.debug("New article: %s", title[:60])
logger.info("Processed %s: %d new, %d skipped", source_name, len(new_articles), skipped_count)
return new_articles
def process_all_feeds(
rss_feeds_path: Path = SCRIPT_DIR / 'rss_feeds.json',
output_dir: Path = ARCHIVE_DIR,
db_path: Path = ARCHIVE_DIR / 'cache.db'
) -> dict:
"""Process all RSS feeds from rss_feeds.json."""
if not rss_feeds_path.exists():
logger.error("RSS feeds file not found: %s", rss_feeds_path)
return {'success': False, 'error': 'File not found'}
with open(rss_feeds_path, 'r', encoding='utf-8') as f:
rss_feeds = json.load(f)
results = {
'total': 0,
'success': 0,
'failed': 0,
'new_articles': 0,
'errors': []
}
for source_name, feed_info in rss_feeds.items():
rss_url = feed_info.get('rss_url', '')
if not rss_url:
logger.warning("No RSS URL for source: %s", source_name)
continue
if feed_info.get('disabled', False):
logger.info("Skipping disabled source: %s (%s)", source_name,
feed_info.get('disable_reason', 'No reason provided'))
continue
results['total'] += 1
try:
articles = process_rss_feed(rss_url, source_name, output_dir, db_path)
results['success'] += 1
results['new_articles'] += len(articles)
logger.info("Completed %s: %d new articles", source_name, len(articles))
except Exception as e:
results['failed'] += 1
results['errors'].append({
'source': source_name,
'url': rss_url,
'error': str(e)
})
logger.error("Failed to process %s: %s", source_name, str(e))
return results
def main():
parser = argparse.ArgumentParser(description='RSS Feed Processor for NewsArchiver')
parser.add_argument('--url', help='Single RSS URL to process')
parser.add_argument('--source', help='Source name (required if --url provided)')
parser.add_argument('--all', action='store_true', help='Process all feeds from rss_feeds.json')
parser.add_argument('--rss-feeds', type=Path, default=SCRIPT_DIR / 'rss_feeds.json',
help='Path to RSS feeds JSON file')
parser.add_argument('--output', type=Path, default=ARCHIVE_DIR,
help='Output directory for archived content')
parser.add_argument('--verbose', action='store_true', help='Enable verbose logging')
args = parser.parse_args()
if args.verbose:
logger.setLevel(logging.DEBUG)
logger.info("=" * 60)
logger.info("RSS Feed Processor - Phase 2.1")
logger.info("=" * 60)
init_db()
if args.url:
if not args.source:
logger.error("Source name required when using --url")
return
process_rss_feed(args.url, args.source, args.output)
elif args.all:
results = process_all_feeds(args.rss_feeds, args.output)
print("\n" + "=" * 60)
print("PROCESSING COMPLETE")
print("=" * 60)
print(f"Total feeds: {results['total']}")
print(f"Successful: {results['success']}")
print(f"Failed: {results['failed']}")
print(f"New articles: {results['new_articles']}")
if results['errors']:
print("\nErrors:")
for error in results['errors']:
print(f" - {error['source']}: {error['error']}")
print("=" * 60)
else:
parser.print_help()
if __name__ == '__main__':
main()