rates/backend/scheduler.py
jtricerolph c924d793e8 Fix scheduled Newbook sync — await the coroutine instead of run_in_executor
run_fetch_current_rates is async; run_in_executor called it in a thread
and discarded the coroutine, so the daily scheduled sync silently did
nothing.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-05 14:09:40 +00:00

120 lines
3.6 KiB
Python

"""
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
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()