""" APScheduler configuration for rate monitor jobs """ import logging from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger from sqlalchemy import text from database import SyncSessionLocal logger = logging.getLogger(__name__) scheduler = AsyncIOScheduler( job_defaults={ 'misfire_grace_time': 3600, 'coalesce': True, } ) def get_config_value(key: str, default: str = None) -> str: db = SyncSessionLocal() try: result = db.execute( text("SELECT config_value FROM system_config WHERE config_key = :key"), {"key": key} ) row = result.fetchone() if row and row.config_value: return row.config_value return default except Exception as e: logger.error(f"Error getting config {key}: {e}") return default finally: db.close() def is_sync_enabled(source: str) -> bool: value = get_config_value(f"sync_{source}_enabled") if value: return value.lower() in ('true', '1', 'yes', 'enabled') return False def get_sync_time(source: str, default_hour: int = 5, default_minute: int = 0) -> tuple: time_str = get_config_value(f"sync_{source}_time") if time_str: try: parts = time_str.split(':') return (int(parts[0]), int(parts[1])) except (ValueError, IndexError): pass return (default_hour, default_minute) async def run_scheduled_direct_scrape(): from jobs.scrape_direct_rates import run_scrape_all_direct import asyncio loop = asyncio.get_event_loop() await loop.run_in_executor(None, run_scrape_all_direct) async def run_short_newbook_rescrape(): """Fetch next-30-day NewBook rates — lightweight complement to the full nightly run.""" from jobs.fetch_current_rates import run_fetch_current_rates await run_fetch_current_rates(horizon_days=30) def _compute_rescrape_times(base_hour: int, base_minute: int, interval_hours: int) -> list: """Return (hour, minute) tuples evenly spaced around the clock, excluding the base (main) run.""" return [ ((base_hour + offset) % 24, base_minute) for offset in range(interval_hours, 24, interval_hours) ] def apply_newbook_rescrape_schedule(): """Read config and (re)register 30-day NewBook rescrape jobs without a scheduler restart.""" for job in list(scheduler.get_jobs()): if job.id.startswith('nb_rescrape_'): scheduler.remove_job(job.id) interval_str = get_config_value('newbook_rescrape_interval_hours', '0') try: interval_hours = int(interval_str or '0') except ValueError: interval_hours = 0 if interval_hours not in (2, 4, 6, 12): logger.info("NewBook 30-day rescrape disabled") return rates_hour, rates_minute = get_sync_time('newbook_current_rates', 5, 20) times = _compute_rescrape_times(rates_hour, rates_minute, interval_hours) for i, (h, m) in enumerate(times): scheduler.add_job( run_short_newbook_rescrape, CronTrigger(hour=h, minute=m), id=f'nb_rescrape_{i}', replace_existing=True, ) logger.info( f"NewBook rescrape: {len(times)} extra run(s) at {interval_hours}h intervals " f"(base {rates_hour:02d}:{rates_minute:02d}): " f"{[f'{h:02d}:{m:02d}' for h, m in times]}" ) async def run_scheduled_parity_check(): from jobs.check_rate_parity import run_parity_check import asyncio loop = asyncio.get_event_loop() await loop.run_in_executor(None, run_parity_check) async def run_scrape_watchdog(): """Auto-release the scrape lock if held for >3 hours (hung Playwright browser).""" from services.booking_scraper import get_lock_status, force_reset_scraper status = get_lock_status() held = status.get("held_seconds") if held and held > 3 * 3600: logger.warning(f"Scrape watchdog: lock held for {held}s — force releasing") db = SyncSessionLocal() try: force_reset_scraper(db) except Exception as e: logger.error(f"Scrape watchdog reset failed: {e}") finally: db.close() async def run_scheduled_booking_scrape_async(): from jobs.scrape_booking_rates import run_scheduled_booking_scrape import asyncio loop = asyncio.get_event_loop() await loop.run_in_executor(None, run_scheduled_booking_scrape) async def run_scheduled_fetch_current_rates(): if is_sync_enabled("newbook_current_rates"): # run_fetch_current_rates is a coroutine — await it directly; # run_in_executor would return the coroutine object unawaited from jobs.fetch_current_rates import run_fetch_current_rates from jobs.sync_occupancy import run_sync_occupancy try: await run_sync_occupancy() except Exception as e: logger.warning(f"Occupancy sync failed (continuing with rates): {e}") await run_fetch_current_rates() else: logger.debug("Newbook current rates sync skipped (disabled)") def start_scheduler(): # Booking.com scrape — daily at configurable time (default 05:30) scrape_time = get_config_value('booking_scraper_daily_time', '05:30') try: h, m = scrape_time.split(':') scrape_hour, scrape_minute = int(h), int(m) except Exception: scrape_hour, scrape_minute = 5, 30 scheduler.add_job( run_scheduled_booking_scrape_async, CronTrigger(hour=scrape_hour, minute=scrape_minute), id='booking_scrape', replace_existing=True, ) # Newbook current rates fetch — daily at 05:20 rates_hour, rates_minute = get_sync_time('newbook_current_rates', 5, 20) scheduler.add_job( run_scheduled_fetch_current_rates, CronTrigger(hour=rates_hour, minute=rates_minute), id='fetch_current_rates', replace_existing=True, ) # Direct booking engine scrape — daily at 06:00 scheduler.add_job( run_scheduled_direct_scrape, CronTrigger(hour=6, minute=0), id='scrape_direct_rates', replace_existing=True, ) # Rate parity check — daily at 06:45, after Newbook fetch + Booking scrape scheduler.add_job( run_scheduled_parity_check, CronTrigger(hour=6, minute=45), id='parity_check', replace_existing=True, ) # Scrape lock watchdog — every 30 min; force-releases if held >3 hours from apscheduler.triggers.interval import IntervalTrigger scheduler.add_job( run_scrape_watchdog, IntervalTrigger(minutes=30), id='scrape_watchdog', replace_existing=True, ) scheduler.start() logger.info(f"Scheduler started: booking scrape at {scrape_hour:02d}:{scrape_minute:02d}, rates fetch at {rates_hour:02d}:{rates_minute:02d}, direct scrape at 06:00, parity check at 06:45, watchdog every 30m") # 30-day NewBook rescrape — optional intraday refresh, interval from config apply_newbook_rescrape_schedule() def shutdown_scheduler(): if scheduler.running: scheduler.shutdown()