NewsArchiverV2/scheduler.py

200 lines
5.8 KiB
Python

#!/usr/bin/env python3
"""Scheduler for NewsArchiver - Phase 4
Background scheduler using APScheduler to automate
daily archiving of news sources.
"""
import atexit
import logging
import signal
import sys
import time
from datetime import datetime
from pathlib import Path
try:
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.interval import IntervalTrigger
except ImportError:
print("ERROR: APScheduler is required. Install with: pip install apscheduler")
sys.exit(1)
try:
from archive_engine import archive_all_sources
except ImportError:
print("ERROR: archive_engine is required")
sys.exit(1)
SCRIPT_DIR = Path(__file__).parent
ARCHIVE_DIR = SCRIPT_DIR / 'archival_data'
ARCHIVE_DIR.mkdir(exist_ok=True)
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s',
handlers=[
logging.StreamHandler(sys.stdout),
logging.FileHandler(ARCHIVE_DIR / 'processing.log', encoding='utf-8')
]
)
logger = logging.getLogger(__name__)
scheduler = BackgroundScheduler()
# Timeout configuration
MAX_RUN_TIME_SECONDS = 3600 # 1 hour
start_time = None
def timeout_handler(signum, frame):
"""Handle timeout signal - exit gracefully after current download completes."""
logger.warning("Timeout reached (%d seconds). Will exit after current download completes.", MAX_RUN_TIME_SECONDS)
raise SystemExit(0)
def check_timeout() -> bool:
"""Check if timeout has been reached.
Returns:
True if timeout reached, False otherwise
"""
global start_time
elapsed = (datetime.now() - start_time).total_seconds()
if elapsed >= MAX_RUN_TIME_SECONDS:
logger.warning("Maximum runtime of %d seconds reached (%d seconds elapsed)", MAX_RUN_TIME_SECONDS, int(elapsed))
return True
return False
def scheduled_archive() -> None:
"""Run archiving for all sources."""
global start_time
start_time = datetime.now()
logger.info("=" * 60)
logger.info("Starting scheduled archive run")
logger.info("=" * 60)
signal.signal(signal.SIGALRM, timeout_handler)
signal.alarm(MAX_RUN_TIME_SECONDS)
try:
results = archive_all_sources(
rss_feeds_path=SCRIPT_DIR / 'rss_feeds.json',
output_dir=ARCHIVE_DIR,
dry_run=False
)
if results['success']:
logger.info("Scheduled archive completed successfully")
logger.info("Sources processed: %d", results.get('sources_processed', 0))
logger.info("Total articles archived: %d", results.get('total_articles_archived', 0))
else:
logger.error("Scheduled archive failed: %s", results.get('error', 'Unknown error'))
except SystemExit as e:
logger.info("Scheduler exiting due to timeout")
raise e
except Exception as e:
logger.error("Scheduled archive failed with exception: %s", str(e))
finally:
signal.alarm(0)
def start_scheduler(interval_minutes: int = 60) -> BackgroundScheduler:
"""Start the background scheduler.
Args:
interval_minutes: Interval between archive runs in minutes
Returns:
The scheduler instance
"""
scheduler.add_job(
func=scheduled_archive,
trigger=IntervalTrigger(minutes=interval_minutes),
id='archive_news',
replace_existing=True,
misfire_grace_time=60,
coalesce=True
)
scheduler.start()
logger.info("Scheduler started with %d minute interval", interval_minutes)
logger.info("Running initial archive immediately...")
scheduled_archive()
return scheduler
def run_once() -> None:
"""Run archiving once (for CLI --run flag)."""
global start_time
start_time = datetime.now()
logger.info("=" * 60)
logger.info("Running one-time archive")
logger.info("=" * 60)
signal.signal(signal.SIGALRM, timeout_handler)
signal.alarm(MAX_RUN_TIME_SECONDS)
try:
scheduled_archive()
logger.info("=" * 60)
logger.info("One-time archive completed")
logger.info("=" * 60)
except SystemExit as e:
logger.info("Archiver exiting due to timeout")
raise e
finally:
signal.alarm(0)
def stop_scheduler() -> None:
"""Stop the scheduler gracefully."""
if scheduler.running:
scheduler.shutdown()
logger.info("Scheduler stopped")
atexit.register(lambda: stop_scheduler())
if __name__ == '__main__':
import argparse
parser = argparse.ArgumentParser(description='NewsArchiver - Scheduler')
parser.add_argument('--run', action='store_true', help='Run archiving once')
parser.add_argument('--serve', action='store_true', help='Start web server')
parser.add_argument('--interval', type=int, default=60, help='Scheduler interval in minutes (default: 60)')
parser.add_argument('--host', default='0.0.0.0', help='Host for web server')
parser.add_argument('--port', type=int, default=5000, help='Port for web server')
args = parser.parse_args()
if args.run:
run_once()
stop_scheduler()
elif args.serve:
from web_interface import app
logger.info("Starting web server on %s:%d", args.host, args.port)
try:
app.run(host=args.host, port=args.port)
except Exception as e:
logger.error("Web server error: %s", str(e))
sys.exit(1)
else:
start_scheduler(args.interval)
logger.info("Press Ctrl+C to stop")
try:
while True:
time.sleep(1)
if check_timeout():
logger.info("Maximum runtime reached. Exiting.")
stop_scheduler()
sys.exit(0)
except (KeyboardInterrupt, SystemExit):
stop_scheduler()