rates/backend/scheduler.py
jtricerolph 078cb47b16 Add configurable 30-day NewBook rate rescrape on 2/4/6/12h intervals
Intraday rescrape jobs are distributed evenly between the main nightly run
(05:20) and cover only the next 30 days — lightweight complement to the
full 720-day nightly sweep. Interval is configurable from the Newbook tab
in Settings and takes effect immediately without a container restart.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-15 09:44:18 +00:00

215 lines
7.1 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_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()