#!/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 os import sys import threading 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.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__) scheduler = BackgroundScheduler() # Timeout configuration MAX_RUN_TIME_SECONDS = 3600 # 1 hour start_time = None _timeout_timer = None def _timeout_checker(): """Daemon thread that raises SystemExit when timeout is reached.""" elapsed = (datetime.now() - start_time).total_seconds() remaining = MAX_RUN_TIME_SECONDS - elapsed if remaining > 0: 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 _start_timeout_thread(): """Start a daemon thread that will raise SystemExit after MAX_RUN_TIME_SECONDS.""" global _timeout_timer _timeout_timer = threading.Timer(MAX_RUN_TIME_SECONDS, _timeout_checker) _timeout_timer.daemon = True _timeout_timer.start() def _cancel_timeout_thread(): """Cancel the timeout thread.""" global _timeout_timer if _timeout_timer: _timeout_timer.cancel() _timeout_timer = None 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) _start_timeout_thread() 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: _cancel_timeout_thread() 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) _start_timeout_thread() 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: _cancel_timeout_thread() 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()