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