NewsArchiverV2/storage_manager.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

860 lines
27 KiB
Python

#!/usr/bin/env python3
"""Storage Manager for NewsArchiver - Phase 2.3
Organizes archived data by newspaper, manages SQLite cache database,
stores both raw HTML and extracted content.
"""
import json
import logging
import os
import sqlite3
import sys
from datetime import datetime
from pathlib import Path
from typing import List, Optional
from content_extractor import ArticleData
try:
import feedgenerator
FEEDGENERATOR_AVAILABLE = True
except ImportError:
FEEDGENERATOR_AVAILABLE = False
SCRIPT_DIR = Path(__file__).parent.resolve()
ARCHIVE_DIR = Path(os.environ.get("ARCHIVE_DIR", str(SCRIPT_DIR / "archival_data"))).resolve()
ARCHIVE_DIR.mkdir(exist_ok=True)
logger = logging.getLogger(__name__)
DB_PATH = ARCHIVE_DIR / "cache.db"
WEBSITES_DIR = ARCHIVE_DIR / "websites"
def _get_db_connection() -> sqlite3.Connection:
"""Get database connection with row factory."""
conn = sqlite3.connect(DB_PATH, timeout=30.0, isolation_level=None)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA busy_timeout=30000")
return conn
def _init_database() -> None:
"""Initialize database schema."""
with _get_db_connection() as conn:
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,
extraction_method TEXT,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
""")
cursor.execute("""
CREATE TABLE IF NOT EXISTS article_archives (
article_url TEXT PRIMARY KEY,
source_name TEXT,
archive_file_path TEXT,
created_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 ON articles(source_name)"
)
cursor.execute(
"CREATE INDEX IF NOT EXISTS idx_articles_url ON articles(article_url)"
)
cursor.execute(
"CREATE INDEX IF NOT EXISTS idx_articles_status ON articles(status)"
)
# Add extraction_method column if it doesn't exist
try:
cursor.execute("ALTER TABLE articles ADD COLUMN extraction_method TEXT")
conn.commit()
except sqlite3.OperationalError:
pass
cursor.execute(
"CREATE INDEX IF NOT EXISTS idx_articles_extraction_method ON articles(extraction_method)"
)
cursor.execute(
"CREATE INDEX IF NOT EXISTS idx_article_archives_url ON article_archives(article_url)"
)
cursor.execute(
"CREATE INDEX IF NOT EXISTS idx_article_archives_source ON article_archives(source_name)"
)
logger.debug("Database initialized at %s", DB_PATH)
def _ensure_directory_structure(source_name: str, date_str: str) -> tuple:
"""Ensure directory structure exists for a source and date.
Args:
source_name: Newspaper source name
date_str: Date string in YYYY-MM-DD format
Returns:
Tuple of (html_dir, articles_dir) Path objects
"""
source_dir = WEBSITES_DIR / source_name
html_dir = source_dir / "html" / date_str
articles_dir = source_dir / "articles" / date_str
html_dir.mkdir(parents=True, exist_ok=True)
articles_dir.mkdir(parents=True, exist_ok=True)
return html_dir, articles_dir
def _get_next_file_index(html_dir: Path, articles_dir: Path) -> int:
"""Get next available file index for article naming.
Args:
html_dir: Directory containing HTML files
articles_dir: Directory containing JSON files
Returns:
Next available index (1-indexed)
"""
def get_max_index(directory: Path, extension: str) -> int:
max_idx = 0
if directory.exists():
for file in directory.glob(f"*{extension}"):
try:
name = file.stem
if name.startswith("article_"):
idx = int(name.replace("article_", ""))
max_idx = max(max_idx, idx)
except ValueError:
continue
return max_idx
html_idx = get_max_index(html_dir, ".html")
json_idx = get_max_index(articles_dir, ".json")
return max(html_idx, json_idx) + 1
def _log_processing(source_name: str, action: str, status: str, message: str) -> None:
"""Log processing action to database.
Args:
source_name: Newspaper source name
action: Action performed
status: Status of action
message: Log message
"""
try:
with _get_db_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"INSERT INTO processing_log (source_name, action, status, message) VALUES (?, ?, ?, ?)",
(source_name, action, status, message),
)
except Exception as e:
logger.error("Failed to log processing: %s", str(e))
def _save_archive_mapping(
article_url: str, source_name: str, archive_file_path: str
) -> None:
"""Save mapping between article URL and archive file path.
Args:
article_url: Article URL
source_name: Newspaper source name
archive_file_path: Path to archived HTML file
"""
try:
with _get_db_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"""
INSERT OR REPLACE INTO article_archives (article_url, source_name, archive_file_path)
VALUES (?, ?, ?)
""",
(article_url, source_name, archive_file_path),
)
logger.debug("Saved archive mapping: %s -> %s", article_url, archive_file_path)
except Exception as e:
logger.error("Failed to save archive mapping: %s", str(e))
def save_article(source_name: str, article_data: ArticleData) -> str:
"""Save article to storage and update cache.
Args:
source_name: Newspaper source name (e.g., 'reuters', 'bbc')
article_data: ArticleData object with article content
Returns:
Status message describing the result
"""
try:
_init_database()
publish_date = article_data.publish_date
if publish_date:
try:
date_obj = datetime.fromisoformat(publish_date.replace("Z", "+00:00"))
date_str = date_obj.strftime("%Y-%m-%d")
except (ValueError, AttributeError):
date_str = datetime.now().strftime("%Y-%m-%d")
else:
date_str = datetime.now().strftime("%Y-%m-%d")
html_dir, articles_dir = _ensure_directory_structure(source_name, date_str)
file_index = _get_next_file_index(html_dir, articles_dir)
file_prefix = f"article_{file_index:03d}"
archive_file_path = html_dir / f"{file_prefix}.html"
metadata_file_path = articles_dir / f"{file_prefix}.json"
if article_data.raw_html:
with open(archive_file_path, "w", encoding="utf-8") as f:
f.write(article_data.raw_html)
metadata = {
"id": file_index,
"source_name": source_name,
"url": article_data.url,
"title": article_data.title,
"author": article_data.author,
"publish_date": article_data.publish_date,
"content_text": article_data.content_text,
"content_html": article_data.content_html,
"tags": article_data.tags,
"language": article_data.language,
"extraction_method": article_data.extraction_method,
"archive_file": f"{file_prefix}.html",
"metadata_file": f"{file_prefix}.json",
"saved_at": datetime.now().isoformat(),
}
with open(metadata_file_path, "w", encoding="utf-8") as f:
json.dump(metadata, f, indent=2, ensure_ascii=False)
with _get_db_connection() as conn:
cursor = conn.cursor()
try:
cursor.execute(
"""
INSERT OR IGNORE INTO articles
(source_name, article_url, article_guid, title, author, publish_date,
content_text, content_html, archive_file_path, metadata_file_path, status, extraction_method)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
source_name,
article_data.url,
getattr(article_data, "guid", None),
article_data.title,
article_data.author,
article_data.publish_date,
article_data.content_text,
article_data.content_html,
str(archive_file_path.relative_to(ARCHIVE_DIR)),
str(metadata_file_path.relative_to(ARCHIVE_DIR)),
"archived" if not article_data.error else "failed",
article_data.extraction_method,
),
)
except sqlite3.IntegrityError:
cursor.execute(
"""
UPDATE articles
SET title = ?, author = ?, publish_date = ?,
content_text = ?, content_html = ?,
archive_file_path = ?, metadata_file_path = ?,
status = ?, extraction_method = ?, updated_at = CURRENT_TIMESTAMP
WHERE article_url = ? AND source_name = ?
""",
(
article_data.title,
article_data.author,
article_data.publish_date,
article_data.content_text,
article_data.content_html,
str(archive_file_path.relative_to(ARCHIVE_DIR)),
str(metadata_file_path.relative_to(ARCHIVE_DIR)),
"archived" if not article_data.error else "failed",
article_data.extraction_method,
article_data.url,
source_name,
),
)
_save_archive_mapping(
article_data.url,
source_name,
str(archive_file_path.relative_to(ARCHIVE_DIR)),
)
_log_processing(
source_name,
"save_article",
"success",
f"Saved article: {article_data.url} -> {metadata_file_path.name}",
)
logger.info(
"Article saved: %s -> %s", article_data.url, metadata_file_path.name
)
return f"Article saved: {metadata_file_path.name}"
except Exception as e:
error_msg = f"Failed to save article: {str(e)}"
logger.error(error_msg)
_log_processing(source_name, "save_article", "error", error_msg)
return error_msg
def get_article(source_name: str, article_id: int) -> Optional[ArticleData]:
"""Retrieve article from storage.
Args:
source_name: Newspaper source name
article_id: Database ID of article
Returns:
ArticleData object or None if not found
"""
try:
_init_database()
with _get_db_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"""
SELECT * FROM articles
WHERE id = ? AND source_name = ?
""",
(article_id, source_name),
)
row = cursor.fetchone()
if not row:
return None
archive_path = (
Path(row["archive_file_path"]) if row["archive_file_path"] else None
)
# Handle both absolute and relative paths
if archive_path:
if not archive_path.is_absolute():
archive_path = ARCHIVE_DIR / archive_path
# Return relative path for web interface
archive_file_path = str(archive_path.relative_to(ARCHIVE_DIR))
else:
archive_file_path = None
archive_content = None
if archive_path and archive_path.exists():
archive_content = archive_path.read_text(encoding="utf-8")
article = ArticleData(
url=row["article_url"],
title=row["title"],
author=row["author"],
publish_date=row["publish_date"],
content_text=row["content_text"],
content_html=row["content_html"],
raw_html=archive_content,
archive_file_path=archive_file_path,
tags=None,
language=None,
metadata=None,
extraction_method=None,
error=row["error_message"],
guid=row["article_guid"],
id=row["id"],
source_name=source_name,
)
return article
except Exception as e:
logger.error("Failed to get article %d: %s", article_id, str(e))
return None
def get_articles_by_source(
source_name: str, limit: int = 50, offset: int = 0
) -> List[ArticleData]:
"""Get paginated articles for a source.
Args:
source_name: Newspaper source name
limit: Maximum number of articles to return
offset: Number of articles to skip
Returns:
List of ArticleData objects
"""
try:
_init_database()
with _get_db_connection() as conn:
cursor = conn.cursor()
cursor.execute(
"""
SELECT * FROM articles
WHERE source_name = ?
ORDER BY publish_date DESC, created_at DESC
LIMIT ? OFFSET ?
""",
(source_name, limit, offset),
)
rows = cursor.fetchall()
articles = []
for row in rows:
archive_path = (
Path(row["archive_file_path"]) if row["archive_file_path"] else None
)
# Handle both absolute and relative paths
if archive_path:
if not archive_path.is_absolute():
archive_path = ARCHIVE_DIR / archive_path
# Return relative path for web interface
archive_file_path = str(archive_path.relative_to(ARCHIVE_DIR))
else:
archive_file_path = None
archive_content = None
if archive_path and archive_path.exists():
archive_content = archive_path.read_text(encoding="utf-8")
article = ArticleData(
url=row["article_url"],
title=row["title"],
author=row["author"],
publish_date=row["publish_date"],
content_text=row["content_text"],
content_html=row["content_html"],
raw_html=archive_content,
archive_file_path=archive_file_path,
tags=None,
language=None,
metadata=None,
extraction_method=None,
error=row["error_message"],
guid=row["article_guid"],
id=row["id"],
source_name=source_name,
)
articles.append(article)
return articles
except Exception as e:
logger.error("Failed to get articles for %s: %s", source_name, str(e))
return []
def update_article_status(
source_name: str, article_url: str, status: str, error: str = None
) -> None:
"""Update article status in cache.
Args:
source_name: Newspaper source name
article_url: Article URL
status: New status (pending, archived, failed)
error: Error message if status is failed
"""
try:
_init_database()
conn = _get_db_connection()
cursor = conn.cursor()
cursor.execute(
"""
UPDATE articles
SET status = ?, error_message = ?, updated_at = CURRENT_TIMESTAMP
WHERE article_url = ? AND source_name = ?
""",
(status, error, article_url, source_name),
)
conn.commit()
conn.close()
if error:
_log_processing(
source_name,
"update_status",
"error",
f"Updated {article_url} status to {status}: {error}",
)
else:
_log_processing(
source_name,
"update_status",
"success",
f"Updated {article_url} status to {status}",
)
logger.info("Updated article status: %s -> %s", article_url, status)
except Exception as e:
logger.error("Failed to update article status: %s", str(e))
def get_source_stats(source_name: str) -> dict:
"""Get statistics for a news source.
Args:
source_name: Newspaper source name
Returns:
Dictionary with source statistics
"""
try:
_init_database()
conn = _get_db_connection()
cursor = conn.cursor()
cursor.execute(
"""
SELECT
COUNT(*) as total,
SUM(CASE WHEN status = 'archived' THEN 1 ELSE 0 END) as archived,
SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END) as failed,
SUM(CASE WHEN status = 'pending' THEN 1 ELSE 0 END) as pending,
MIN(created_at) as first_archived,
MAX(created_at) as last_archived,
MAX(publish_date) as latest_article_date
FROM articles
WHERE source_name = ?
""",
(source_name,),
)
row = cursor.fetchone()
conn.close()
stats = {
"source_name": source_name,
"total_articles": row["total"] or 0,
"archived": row["archived"] or 0,
"failed": row["failed"] or 0,
"pending": row["pending"] or 0,
"first_archived": row["first_archived"],
"last_archived": row["last_archived"],
"latest_article_date": row["latest_article_date"],
}
_log_processing(
source_name,
"get_stats",
"success",
f"Stats: {stats['total_articles']} total, {stats['archived']} archived, {stats['failed']} failed",
)
return stats
except Exception as e:
logger.error("Failed to get stats for %s: %s", source_name, str(e))
return {
"source_name": source_name,
"total_articles": 0,
"archived": 0,
"failed": 0,
"pending": 0,
}
def get_all_sources() -> List[str]:
"""Get list of all sources in the archive.
Returns:
List of source names
"""
try:
_init_database()
conn = _get_db_connection()
cursor = conn.cursor()
cursor.execute("""
SELECT DISTINCT source_name FROM articles ORDER BY source_name
""")
sources = [row["source_name"] for row in cursor.fetchall()]
conn.close()
return sources
except Exception as e:
logger.error("Failed to get sources: %s", str(e))
return []
def get_latest_articles(limit: int = 50) -> List[ArticleData]:
"""Get latest articles across all sources.
Args:
limit: Maximum number of articles to return
Returns:
List of ArticleData objects ordered by publish_date DESC
"""
try:
_init_database()
conn = _get_db_connection()
cursor = conn.cursor()
cursor.execute(
"""
SELECT * FROM articles
ORDER BY publish_date DESC, created_at DESC
LIMIT ?
""",
(limit,),
)
rows = cursor.fetchall()
conn.close()
articles = []
for row in rows:
archive_path = (
Path(row["archive_file_path"]) if row["archive_file_path"] else None
)
# Handle both absolute and relative paths
if archive_path:
if not archive_path.is_absolute():
archive_path = ARCHIVE_DIR / archive_path
# Return relative path for web interface
archive_file_path = str(archive_path.relative_to(ARCHIVE_DIR))
else:
archive_file_path = None
archive_content = None
if archive_path and archive_path.exists():
archive_content = archive_path.read_text(encoding="utf-8")
article = ArticleData(
url=row["article_url"],
title=row["title"],
author=row["author"],
publish_date=row["publish_date"],
content_text=row["content_text"],
content_html=row["content_html"],
raw_html=archive_content,
archive_file_path=archive_file_path,
tags=None,
language=None,
metadata=None,
extraction_method=None,
error=row["error_message"],
guid=row["article_guid"],
id=row["id"],
source_name=row["source_name"],
)
articles.append(article)
return articles
except Exception as e:
logger.error("Failed to get latest articles: %s", str(e))
return []
def get_source_directory(source_name: str) -> Path:
"""Get the directory path for a source.
Args:
source_name: Newspaper source name
Returns:
Path to source directory
"""
return WEBSITES_DIR / source_name
def get_archive_file_path_from_db(
article_url: str, source_name: str = None
) -> Optional[str]:
"""Get archive file path from database mapping.
Args:
article_url: Article URL
source_name: Newspaper source name (optional, for filtering)
Returns:
Archive file path if found, None otherwise
"""
try:
_init_database()
conn = _get_db_connection()
cursor = conn.cursor()
if source_name:
cursor.execute(
"""
SELECT archive_file_path FROM article_archives
WHERE article_url = ? AND source_name = ?
""",
(article_url, source_name),
)
else:
cursor.execute(
"""
SELECT archive_file_path FROM article_archives
WHERE article_url = ?
""",
(article_url,),
)
row = cursor.fetchone()
conn.close()
return row["archive_file_path"] if row else None
except Exception as e:
logger.error("Failed to get archive file path from DB: %s", str(e))
return None
def get_daily_articles(source_name: str, date_str: str) -> List[ArticleData]:
"""Get articles for a specific date.
Args:
source_name: Newspaper source name
date_str: Date string in YYYY-MM-DD format
Returns:
List of ArticleData objects
"""
try:
source_dir = get_source_directory(source_name)
articles_dir = source_dir / "articles" / date_str
if not articles_dir.exists():
return []
articles = []
for json_file in sorted(articles_dir.glob("article_*.json")):
try:
with open(json_file, "r", encoding="utf-8") as f:
metadata = json.load(f)
archive_filename = metadata.get("archive_file", "")
archive_path = (
source_dir / "html" / date_str / archive_filename
if archive_filename
else None
)
archive_content = None
if archive_path and archive_path.exists():
archive_content = archive_path.read_text(encoding="utf-8")
article = ArticleData(
url=metadata.get("url", ""),
title=metadata.get("title"),
author=metadata.get("author"),
publish_date=metadata.get("publish_date"),
content_text=metadata.get("content_text"),
content_html=metadata.get("content_html"),
raw_html=archive_content,
# Handle both absolute and relative paths
archive_file_path=str(archive_path.relative_to(ARCHIVE_DIR))
if archive_path
else None,
tags=metadata.get("tags"),
language=metadata.get("language"),
metadata=metadata,
extraction_method=metadata.get("extraction_method"),
source_name=source_name,
)
articles.append(article)
except Exception as e:
logger.warning("Failed to load article from %s: %s", json_file, str(e))
continue
return articles
except Exception as e:
logger.error(
"Failed to get daily articles for %s on %s: %s",
source_name,
date_str,
str(e),
)
return []
def initialize_storage() -> None:
"""Initialize the storage system.
Creates database, directory structure, and logs initialization.
"""
logger.info("Initializing storage system...")
_init_database()
(WEBSITES_DIR / "sample").mkdir(parents=True, exist_ok=True)
(WEBSITES_DIR / "sample" / "html").mkdir(exist_ok=True)
(WEBSITES_DIR / "sample" / "articles").mkdir(exist_ok=True)
_log_processing("system", "initialize", "success", "Storage system initialized")
logger.info("Storage system initialized at %s", ARCHIVE_DIR)
if __name__ == "__main__":
initialize_storage()
sources = get_all_sources()
print(f"Sources in archive: {sources}")
if sources:
for source in sources:
stats = get_source_stats(source)
print(f"\n{source}:")
print(f" Total: {stats['total_articles']}")
print(f" Archived: {stats['archived']}")
print(f" Failed: {stats['failed']}")