""" 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_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, ) 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") def shutdown_scheduler(): if scheduler.running: scheduler.shutdown()