rates/backend/scheduler.py
jtricerolph e781b1e8b9 Scraper: force-reset endpoint + watchdog to prevent stuck lock
Root cause: threading.Lock held indefinitely when Playwright browser
hangs inside run_in_executor (finally never fires from the async side).

Fixes:
- _acquire_scrape_lock/_release_scrape_lock track monotonic timestamp
- POST /competitors/scrape/reset force-releases the lock and marks any
  running batch as interrupted (queue rows stay intact for retry)
- GET /competitors/status now includes lock_held_seconds
- APScheduler watchdog job every 30 min auto-releases if held >3h
- Settings → Scraper Proxy tab shows live lock status (green/amber)
  with a Force Reset button requiring confirmation

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-07-09 16:56:35 +00:00

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