""" Competitor Rates API endpoints Booking.com rate scraping, hotel management, and competitor comparison """ from typing import Optional, List, Dict, Any from datetime import date, datetime, timedelta from fastapi import APIRouter, Depends, HTTPException, BackgroundTasks from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import text from pydantic import BaseModel import logging from database import get_db, SyncSessionLocal from auth import get_current_user from services import proxy as proxy_util router = APIRouter() logger = logging.getLogger(__name__) # ============================================ # REQUEST/RESPONSE MODELS # ============================================ class ScrapeRequest(BaseModel): from_date: str to_date: Optional[str] = None class LocationConfigRequest(BaseModel): location_name: str = '' pages_to_scrape: int = 2 adults: int = 2 dest_id: Optional[str] = None # Booking.com numeric destination id (pins the search) location_search_url: Optional[str] = None # Pasted Booking.com search URL (pins destination) class HotelTierUpdate(BaseModel): tier: str # 'own', 'competitor', 'market' display_order: Optional[int] = None class HotelResponse(BaseModel): id: int booking_com_id: str name: str booking_com_url: Optional[str] star_rating: Optional[float] review_score: Optional[float] review_count: Optional[int] tier: str display_order: int notes: Optional[str] first_seen_at: Optional[datetime] last_seen_at: Optional[datetime] direct_hotel_id: Optional[int] = None class RateResponse(BaseModel): rate_date: str hotel_id: int hotel_name: str tier: str star_rating: Optional[float] review_score: Optional[float] availability_status: str rate_gross: Optional[float] room_type: Optional[str] breakfast_included: Optional[bool] free_cancellation: Optional[bool] no_prepayment: Optional[bool] rooms_left: Optional[int] scraped_at: Optional[datetime] class ScraperStatusResponse(BaseModel): enabled: bool backend: str location_configured: bool location_name: Optional[str] last_scrape: Optional[dict] # ============================================ # SCRAPER STATUS & CONFIGURATION # ============================================ @router.get("/status", response_model=ScraperStatusResponse) async def get_scraper_status( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Get current scraper status and configuration.""" # Get config values config_result = await db.execute( text(""" SELECT config_key, config_value FROM system_config WHERE config_key IN ( 'booking_scraper_enabled', 'booking_scraper_backend' ) """) ) config = {row.config_key: row.config_value for row in config_result.fetchall()} # Get location config location_result = await db.execute( text("SELECT location_name FROM booking_scrape_config WHERE is_active = TRUE LIMIT 1") ) location_row = location_result.fetchone() # Get last scrape info last_scrape_result = await db.execute( text(""" SELECT batch_id, scrape_type, started_at, completed_at, status, hotels_found, rates_scraped, error_message FROM booking_scrape_log ORDER BY started_at DESC LIMIT 1 """) ) last_scrape_row = last_scrape_result.fetchone() last_scrape = None if last_scrape_row: last_scrape = { 'batch_id': str(last_scrape_row.batch_id), 'scrape_type': last_scrape_row.scrape_type, 'started_at': last_scrape_row.started_at.isoformat() if last_scrape_row.started_at else None, 'completed_at': last_scrape_row.completed_at.isoformat() if last_scrape_row.completed_at else None, 'status': last_scrape_row.status, 'hotels_found': last_scrape_row.hotels_found, 'rates_scraped': last_scrape_row.rates_scraped, 'error_message': last_scrape_row.error_message, } return ScraperStatusResponse( enabled=config.get('booking_scraper_enabled', 'false') == 'true', backend=config.get('booking_scraper_backend', 'playwright_local'), location_configured=location_row is not None, location_name=location_row.location_name if location_row else None, last_scrape=last_scrape ) @router.get("/config/system") async def get_system_config( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Return all system_config rows as a flat dict (secrets masked).""" result = await db.execute(text("SELECT config_key, config_value FROM system_config")) out = {} for row in result.fetchall(): if row.config_value and ('password' in row.config_key or 'secret' in row.config_key): out[row.config_key] = '********' else: out[row.config_key] = row.config_value return out class SystemConfigUpdate(BaseModel): key: str value: str @router.post("/config/system") async def set_system_config( payload: SystemConfigUpdate, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Upsert a single system_config key.""" await db.execute( text(""" INSERT INTO system_config (config_key, config_value) VALUES (:key, :value) ON CONFLICT (config_key) DO UPDATE SET config_value = EXCLUDED.config_value """), {'key': payload.key, 'value': payload.value} ) await db.commit() return {"status": "success", "key": payload.key} # ── Booking.com scraper proxy config ───────────────────────────────────────── class ProxyConfig(BaseModel): enabled: bool = False host: str = '' port: str = '823' username: str = '' password: Optional[str] = None # None/'' => keep the stored password country: str = 'gb' @router.get("/config/proxy") async def get_proxy_config( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Return the scraper proxy config. Password is never returned — only a `password_set` flag indicating whether one is stored.""" result = await db.execute( text("SELECT config_key, config_value FROM system_config WHERE config_key LIKE 'booking_proxy_%'") ) cfg = {row.config_key: row.config_value for row in result.fetchall()} return { 'enabled': cfg.get('booking_proxy_enabled') == 'true', 'host': cfg.get('booking_proxy_host', ''), 'port': cfg.get('booking_proxy_port', '823'), 'username': cfg.get('booking_proxy_username', ''), 'country': cfg.get('booking_proxy_country', 'gb'), 'password_set': bool(cfg.get('booking_proxy_password')), } @router.post("/config/proxy") async def set_proxy_config( payload: ProxyConfig, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Upsert the scraper proxy config. A blank password leaves the stored one untouched, so the UI never has to round-trip the secret.""" values = { 'booking_proxy_enabled': 'true' if payload.enabled else 'false', 'booking_proxy_host': payload.host.strip(), 'booking_proxy_port': (payload.port or '823').strip(), 'booking_proxy_username': payload.username.strip(), 'booking_proxy_country': (payload.country or 'gb').strip(), } if payload.password: # only overwrite when a new value is provided values['booking_proxy_password'] = payload.password.strip() for key, value in values.items(): await db.execute( text(""" INSERT INTO system_config (config_key, config_value) VALUES (:key, :value) ON CONFLICT (config_key) DO UPDATE SET config_value = EXCLUDED.config_value """), {'key': key, 'value': value} ) await db.commit() return {"status": "success"} @router.post("/config/proxy/test") async def test_proxy_config( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Make a live request through the configured proxy and report the exit IP and country, so the user can confirm credentials + geo before scraping.""" result = await db.execute( text("SELECT config_key, config_value FROM system_config WHERE config_key LIKE 'booking_proxy_%'") ) raw = {row.config_key: row.config_value for row in result.fetchall()} # Test whatever creds are stored, regardless of the enabled toggle, so the # user can verify before switching the proxy on. cfg = { 'host': (raw.get('booking_proxy_host') or '').strip(), 'port': (raw.get('booking_proxy_port') or '823').strip(), 'username': (raw.get('booking_proxy_username') or '').strip(), 'password': (raw.get('booking_proxy_password') or '').strip(), 'country': (raw.get('booking_proxy_country') or 'gb').strip(), } if not (proxy_util.is_enabled(cfg) and cfg['password']): raise HTTPException(status_code=400, detail="Proxy host, username and password must be saved first.") proxy_url = proxy_util.httpx_proxy_url(cfg, proxy_util.new_session_id()) import httpx try: async with httpx.AsyncClient(proxy=proxy_url, timeout=40.0) as client: resp = await client.get("https://ipinfo.io/json") resp.raise_for_status() data = resp.json() except Exception as e: logger.warning(f"Proxy test failed: {e}") raise HTTPException(status_code=502, detail=f"Proxy test failed: {e}") return { 'ok': True, 'ip': data.get('ip'), 'country': data.get('country'), 'city': data.get('city'), 'org': data.get('org'), } @router.post("/config/location") async def set_location_config( config: LocationConfigRequest, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Set the location to scrape for competitor rates.""" from services.scraper_backends.playwright_local import location_params_from_url # If a full search URL was pasted, lift the destination params out of it — # dest_id for the column (visibility/fallback), and ss for the display name # when the user didn't type one. search_url = (config.location_search_url or '').strip() or None url_params = location_params_from_url(search_url) if search_url else {} dest_id = config.dest_id or url_params.get('dest_id') location_name = config.location_name.strip() or url_params.get('ss', '') if not location_name: raise HTTPException(status_code=400, detail="Provide a location name or a search URL.") # Deactivate existing configs await db.execute( text("UPDATE booking_scrape_config SET is_active = FALSE") ) # Insert new config await db.execute( text(""" INSERT INTO booking_scrape_config (location_name, pages_to_scrape, adults, dest_id, location_search_url, is_active) VALUES (:location, :pages, :adults, :dest_id, :search_url, TRUE) """), {'location': location_name, 'pages': config.pages_to_scrape, 'adults': config.adults, 'dest_id': dest_id, 'search_url': search_url} ) await db.commit() return {"status": "success", "location": location_name} @router.post("/config/enable") async def enable_scraper( enabled: bool = True, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Enable or disable the booking.com scraper.""" await db.execute( text("UPDATE system_config SET config_value = :val WHERE config_key = 'booking_scraper_enabled'"), {'val': 'true' if enabled else 'false'} ) await db.commit() return {"status": "success", "enabled": enabled} # ============================================ # MANUAL SCRAPE TRIGGER # ============================================ def run_scrape_sync(from_date: date, to_date: date): """Run scrape in sync context for background task.""" import asyncio from services.booking_scraper import run_manual_scrape, cleanup_stale_batches db = SyncSessionLocal() try: # Clean up any stale batches before starting cleanup_stale_batches(db, max_age_minutes=60) loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) try: result = loop.run_until_complete(run_manual_scrape(db, from_date, to_date)) logger.info(f"Background scrape completed: {result}") finally: loop.close() except Exception as e: logger.error(f"Background scrape failed: {e}", exc_info=True) finally: db.close() @router.post("/scrape") async def trigger_manual_scrape( request: ScrapeRequest, background_tasks: BackgroundTasks, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ Trigger a manual scrape for the specified date range. Runs in background - check /status for progress. """ try: from_date = date.fromisoformat(request.from_date) to_date = date.fromisoformat(request.to_date) if request.to_date else from_date except ValueError as e: raise HTTPException(status_code=400, detail=f"Invalid date format: {e}") if to_date < from_date: raise HTTPException(status_code=400, detail="to_date must be after from_date") if (to_date - from_date).days > 30: raise HTTPException(status_code=400, detail="Date range cannot exceed 30 days for manual scrape") # Check if location is configured location_result = await db.execute( text("SELECT id FROM booking_scrape_config WHERE is_active = TRUE LIMIT 1") ) if not location_result.fetchone(): raise HTTPException(status_code=400, detail="No scrape location configured. Set location first.") # Only one scrape at a time — concurrent Chromium runs cause the page # timeouts that produce partial results # # Best-effort early 409: the background task re-acquires the lock and will # no-op (logging "another scrape is running") if it loses a millisecond- # window race, so no double-run can slip through here. from services.booking_scraper import SCRAPE_LOCK if SCRAPE_LOCK.locked(): raise HTTPException(status_code=409, detail="A scrape is already running. Try again when it finishes.") # Start background task background_tasks.add_task(run_scrape_sync, from_date, to_date) return { "status": "started", "from_date": from_date.isoformat(), "to_date": to_date.isoformat(), "message": "Scrape started in background. Check /status for progress." } # ============================================ # HOTELS MANAGEMENT # ============================================ @router.get("/hotels", response_model=List[HotelResponse]) async def list_hotels( tier: Optional[str] = None, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ List all discovered hotels. Filter by tier: 'own', 'competitor', 'market', or None for all. """ query = """ SELECT id, booking_com_id, name, booking_com_url, star_rating, review_score, review_count, tier, display_order, notes, first_seen_at, last_seen_at, direct_hotel_id FROM booking_com_hotels WHERE is_active = TRUE """ params = {} if tier: if tier not in ('own', 'competitor', 'market'): raise HTTPException(status_code=400, detail="Invalid tier. Must be 'own', 'competitor', or 'market'") query += " AND tier = :tier" params['tier'] = tier query += " ORDER BY display_order, name" result = await db.execute(text(query), params) return [ HotelResponse( id=row.id, booking_com_id=row.booking_com_id or '', name=row.name, booking_com_url=row.booking_com_url, star_rating=float(row.star_rating) if row.star_rating else None, review_score=float(row.review_score) if row.review_score else None, review_count=row.review_count, tier=row.tier, display_order=row.display_order, notes=row.notes, first_seen_at=row.first_seen_at, last_seen_at=row.last_seen_at, direct_hotel_id=row.direct_hotel_id ) for row in result.fetchall() ] @router.put("/hotels/{hotel_id}/tier") async def update_hotel_tier( hotel_id: int, update: HotelTierUpdate, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ Update a hotel's tier and display order. Tiers: - 'own': Your hotel (for parity checking) - 'competitor': Main competitors (full tracking) - 'market': Other hotels (context only) """ if update.tier not in ('own', 'competitor', 'market'): raise HTTPException(status_code=400, detail="Invalid tier") # If setting as 'own', clear any existing 'own' hotel if update.tier == 'own': await db.execute( text("UPDATE booking_com_hotels SET tier = 'market' WHERE tier = 'own'") ) # Update the hotel set_clause = "tier = :tier" params = {'hotel_id': hotel_id, 'tier': update.tier} if update.display_order is not None: set_clause += ", display_order = :order" params['order'] = update.display_order result = await db.execute( text(f"UPDATE booking_com_hotels SET {set_clause} WHERE id = :hotel_id RETURNING id"), params ) if not result.fetchone(): raise HTTPException(status_code=404, detail="Hotel not found") await db.commit() # If this is now the own hotel, update system config if update.tier == 'own': await db.execute( text("UPDATE system_config SET config_value = :val WHERE config_key = 'booking_scraper_own_hotel_id'"), {'val': str(hotel_id)} ) await db.commit() return {"status": "success", "hotel_id": hotel_id, "tier": update.tier} @router.put("/hotels/{hotel_id}/notes") async def update_hotel_notes( hotel_id: int, notes: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Update notes for a hotel.""" result = await db.execute( text("UPDATE booking_com_hotels SET notes = :notes WHERE id = :hotel_id RETURNING id"), {'hotel_id': hotel_id, 'notes': notes} ) if not result.fetchone(): raise HTTPException(status_code=404, detail="Hotel not found") await db.commit() return {"status": "success"} class DirectLinkUpdate(BaseModel): direct_hotel_id: Optional[int] = None @router.put("/hotels/{hotel_id}/direct-link") async def update_hotel_direct_link( hotel_id: int, payload: DirectLinkUpdate, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Link (or unlink with null) a Booking.com hotel to a direct booking engine competitor, enabling the direct-rates sub-row in the rate matrix.""" if payload.direct_hotel_id is not None: exists = await db.execute( text("SELECT 1 FROM direct_competitor_hotels WHERE id = :id"), {'id': payload.direct_hotel_id} ) if not exists.fetchone(): raise HTTPException(status_code=404, detail="Direct hotel not found") result = await db.execute( text("UPDATE booking_com_hotels SET direct_hotel_id = :did WHERE id = :hotel_id RETURNING id"), {'hotel_id': hotel_id, 'did': payload.direct_hotel_id} ) if not result.fetchone(): raise HTTPException(status_code=404, detail="Hotel not found") await db.commit() return {"status": "success"} # ============================================ # COMPETITOR RATES MATRIX # ============================================ @router.get("/matrix") async def get_competitor_matrix( from_date: Optional[str] = None, to_date: Optional[str] = None, include_market: bool = False, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ Get rate comparison matrix for competitors. Returns rates for own hotel and competitors, organized by date. Set include_market=true to also include market tier hotels. """ today = date.today() start = date.fromisoformat(from_date) if from_date else today end = date.fromisoformat(to_date) if to_date else today + timedelta(days=30) if end < start: raise HTTPException(status_code=400, detail="to_date must be after from_date") if (end - start).days > 90: raise HTTPException(status_code=400, detail="Date range cannot exceed 90 days") tier_filter = "h.tier IN ('own', 'competitor')" if include_market: tier_filter = "h.tier IN ('own', 'competitor', 'market')" # Get hotels hotels_result = await db.execute( text(f""" SELECT id, name, tier, display_order, star_rating, review_score, booking_com_url, direct_hotel_id FROM booking_com_hotels WHERE is_active = TRUE AND {tier_filter.replace('h.', '')} ORDER BY display_order, name """) ) hotels = [dict(row._mapping) for row in hotels_result.fetchall()] # Get latest rates using the view rates_result = await db.execute( text(f""" SELECT DISTINCT ON (r.hotel_id, r.rate_date) r.hotel_id, r.rate_date, r.availability_status, r.rate_gross, r.room_type, r.breakfast_included, r.free_cancellation, r.no_prepayment, r.rooms_left, r.scraped_at FROM booking_com_rates r JOIN booking_com_hotels h ON r.hotel_id = h.id WHERE {tier_filter} AND h.is_active = TRUE AND r.rate_date >= :from_date AND r.rate_date <= :to_date ORDER BY r.hotel_id, r.rate_date, r.scraped_at DESC """), {'from_date': start, 'to_date': end} ) # Build matrix: hotel_id -> date -> rate data rates_by_hotel: Dict[int, Dict[str, dict]] = {} for row in rates_result.fetchall(): hotel_id = row.hotel_id rate_date = row.rate_date.isoformat() if hotel_id not in rates_by_hotel: rates_by_hotel[hotel_id] = {} rates_by_hotel[hotel_id][rate_date] = { 'availability_status': row.availability_status, 'rate_gross': float(row.rate_gross) if row.rate_gross else None, 'room_type': row.room_type, 'breakfast_included': row.breakfast_included, 'free_cancellation': row.free_cancellation, 'no_prepayment': row.no_prepayment, 'rooms_left': row.rooms_left, 'scraped_at': row.scraped_at.isoformat() if row.scraped_at else None, } # Most recent scrape touching each date, across ALL hotels — a partial # scrape may refresh the date without touching the displayed hotels, so # per-cell scraped_at can lag behind this column-level timestamp last_scraped_result = await db.execute( text(""" SELECT rate_date, MAX(scraped_at) AS last_scraped FROM booking_com_rates WHERE rate_date >= :from_date AND rate_date <= :to_date GROUP BY rate_date """), {'from_date': start, 'to_date': end} ) last_scraped = { row.rate_date.isoformat(): row.last_scraped.isoformat() if row.last_scraped else None for row in last_scraped_result.fetchall() } # Build date list dates = [] current = start while current <= end: dates.append(current.isoformat()) current += timedelta(days=1) return { 'from_date': start.isoformat(), 'to_date': end.isoformat(), 'dates': dates, 'hotels': hotels, 'rates': rates_by_hotel, 'last_scraped': last_scraped } # ============================================ # RATE PARITY (OWN HOTEL VS NEWBOOK) # ============================================ @router.get("/parity") async def get_rate_parity( from_date: Optional[str] = None, to_date: Optional[str] = None, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ Like-for-like parity comparison between our Booking.com rate and the comparable Newbook tariff (matched on flex/prepaid + board via the scraped rate's condition flags). Same logic as the daily alert job. """ import asyncio from jobs.check_rate_parity import get_parity_config, gather_comparisons, expected_booking_rate from database import SyncSessionLocal today = date.today() start = date.fromisoformat(from_date) if from_date else today end = date.fromisoformat(to_date) if to_date else today + timedelta(days=30) def _run(): sdb = SyncSessionLocal() try: cfg = get_parity_config(sdb) return cfg, gather_comparisons(sdb, start, end, cfg) finally: sdb.close() loop = asyncio.get_event_loop() cfg, comparisons = await loop.run_in_executor(None, _run) parity_issues = [ { 'rate_date': c["date"].isoformat(), 'booking_rate': c["booking_rate"], 'booking_basis': c["booking_basis"], 'booking_room_type': c["booking_room"], 'newbook_rate': c["newbook_rate"], 'newbook_tariff': c["newbook_tariff"], 'match_quality': c["match_quality"], 'expected_rate': round(expected_booking_rate(c["newbook_rate"], cfg), 2), 'difference_pct': round(c["dev_pct"], 2), 'difference_gbp': round(c["dev_gbp"], 2), 'alert_type': 'higher' if c["dev_pct"] > 0 else 'lower', } for c in comparisons if c["breach"] ] return { 'from_date': start.isoformat(), 'to_date': end.isoformat(), 'markup_value': cfg["markup_value"], 'markup_unit': cfg["markup_unit"], 'tolerance_value': cfg["tolerance_value"], 'tolerance_unit': cfg["tolerance_unit"], 'dates_compared': len(comparisons), 'issues_count': len(parity_issues), 'issues': parity_issues } @router.post("/parity/check") async def trigger_parity_check( current_user: dict = Depends(get_current_user) ): """Run the parity check now (same logic as the daily 06:45 job).""" import asyncio from jobs.check_rate_parity import run_parity_check loop = asyncio.get_event_loop() return await loop.run_in_executor(None, run_parity_check) # ============================================ # PARITY ALERTS # ============================================ @router.get("/parity/alerts") async def get_parity_alerts( status: Optional[str] = None, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Get rate parity alerts.""" query = """ SELECT id, rate_date, room_category, newbook_rate, booking_com_rate, difference_pct, alert_type, alert_status, created_at, acknowledged_at, acknowledged_by, notes FROM rate_parity_alerts """ params = {} if status: query += " WHERE alert_status = :status" params['status'] = status query += " ORDER BY rate_date DESC, created_at DESC LIMIT 100" result = await db.execute(text(query), params) return [ { 'id': row.id, 'rate_date': row.rate_date.isoformat(), 'room_category': row.room_category, 'newbook_rate': float(row.newbook_rate) if row.newbook_rate else None, 'booking_com_rate': float(row.booking_com_rate) if row.booking_com_rate else None, 'difference_pct': float(row.difference_pct) if row.difference_pct else None, 'alert_type': row.alert_type, 'alert_status': row.alert_status, 'created_at': row.created_at.isoformat() if row.created_at else None, 'acknowledged_at': row.acknowledged_at.isoformat() if row.acknowledged_at else None, 'acknowledged_by': row.acknowledged_by, 'notes': row.notes, } for row in result.fetchall() ] @router.put("/parity/alerts/{alert_id}/acknowledge") async def acknowledge_parity_alert( alert_id: int, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Acknowledge a parity alert.""" result = await db.execute( text(""" UPDATE rate_parity_alerts SET alert_status = 'acknowledged', acknowledged_at = NOW(), acknowledged_by = :username WHERE id = :alert_id RETURNING id """), {'alert_id': alert_id, 'username': current_user.get('username', 'unknown')} ) if not result.fetchone(): raise HTTPException(status_code=404, detail="Alert not found") await db.commit() return {"status": "success"} # ============================================ # QUEUE STATUS # ============================================ @router.get("/queue-status") async def get_queue_status( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Get current scrape queue status.""" result = await db.execute( text(""" SELECT status, COUNT(*) as count, MIN(rate_date) as earliest_date, MAX(rate_date) as latest_date FROM booking_scrape_queue GROUP BY status """) ) status_counts = {row.status: { 'count': row.count, 'earliest': row.earliest_date.isoformat() if row.earliest_date else None, 'latest': row.latest_date.isoformat() if row.latest_date else None, } for row in result.fetchall()} # Get retry items (failed but under max_attempts) retry_result = await db.execute( text(""" SELECT COUNT(*) as count FROM booking_scrape_queue WHERE status = 'pending' AND attempts > 0 """) ) retry_count = retry_result.fetchone().count return { 'statuses': status_counts, 'retries_pending': retry_count, 'total_pending': status_counts.get('pending', {}).get('count', 0), 'total_completed': status_counts.get('completed', {}).get('count', 0), 'total_failed': status_counts.get('failed', {}).get('count', 0), } # ============================================ # SCHEDULE INFO # ============================================ @router.get("/schedule-info") async def get_schedule_info( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Get information about the scraping schedule.""" # Get configured time time_result = await db.execute( text("SELECT config_value FROM system_config WHERE config_key = 'booking_scraper_daily_time'") ) time_row = time_result.fetchone() daily_time = time_row.config_value if time_row and time_row.config_value else '05:30' # Calculate what today's schedule would look like from jobs.scrape_booking_rates import get_high_priority_dates, get_medium_priority_dates, get_low_priority_dates high = get_high_priority_dates() medium = get_medium_priority_dates() low = get_low_priority_dates() today = date.today() weekday_name = today.strftime('%A') return { 'daily_time': daily_time, 'today': today.isoformat(), 'weekday': weekday_name, 'tiers': { 'high': { 'description': 'Next 30 days (scraped first)', 'dates_today': len(high), 'range': f'{high[0].isoformat()} to {high[-1].isoformat()}' if high else None, }, 'medium': { 'description': 'Days 31-180 (scraped after high priority)', 'dates_today': len(medium), 'range': f'{medium[0].isoformat()} to {medium[-1].isoformat()}' if medium else None, }, 'low': { 'description': 'Days 181-365 (scraped last, or until rate limit)', 'dates_today': len(low), 'range': f'{low[0].isoformat()} to {low[-1].isoformat()}' if low else None, }, }, 'total_dates_today': len(set(high + medium + low)), } # ============================================ # SCRAPE COVERAGE (365-day view) # ============================================ @router.get("/scrape-coverage") async def get_scrape_coverage( db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ Get 365-day scrape coverage showing last scraped time and next expected scrape for every date. """ today = date.today() end = today + timedelta(days=365) # Get latest scraped_at per date (across all hotels) result = await db.execute( text(""" SELECT rate_date, MAX(scraped_at) as last_scraped FROM booking_com_rates WHERE rate_date >= :from_date AND rate_date <= :to_date GROUP BY rate_date """), {'from_date': today, 'to_date': end} ) scraped_map = {row.rate_date: row.last_scraped for row in result.fetchall()} # Compute tier and next scrape for each date from jobs.scrape_booking_rates import compute_next_scrape_for_date coverage = [] for offset in range(366): d = today + timedelta(days=offset) tier, next_scrape = compute_next_scrape_for_date(d) last_scraped = scraped_map.get(d) coverage.append({ 'date': d.isoformat(), 'tier': tier, 'last_scraped': last_scraped.isoformat() if last_scraped else None, 'next_expected': next_scrape.isoformat() if next_scrape else None, }) return { 'today': today.isoformat(), 'coverage': coverage, } # ============================================ # BOOKING.COM AVAILABILITY CHECK (for Bookability page) # ============================================ @router.get("/booking-availability") async def get_booking_availability( from_date: Optional[str] = None, to_date: Optional[str] = None, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """ Check own hotel's availability on booking.com. Returns a simple summary: for each date in the range, whether the own hotel appears available on booking.com based on the latest scrape data. """ today = date.today() start = date.fromisoformat(from_date) if from_date else today end = date.fromisoformat(to_date) if to_date else today + timedelta(days=30) # Get own hotel's latest scraped availability result = await db.execute( text(""" SELECT DISTINCT ON (r.rate_date) r.rate_date, r.availability_status, r.rate_gross, r.scraped_at FROM booking_com_rates r JOIN booking_com_hotels h ON r.hotel_id = h.id WHERE h.tier = 'own' AND r.rate_date >= :from_date AND r.rate_date <= :to_date ORDER BY r.rate_date, r.scraped_at DESC """), {'from_date': start, 'to_date': end} ) rows = result.fetchall() if not rows: return { 'has_own_hotel': False, 'dates_checked': 0, 'dates_available': 0, 'dates_sold_out': 0, 'dates_no_data': 0, 'latest_scrape': None, 'dates': {}, } dates_map = {} dates_available = 0 dates_sold_out = 0 dates_no_data = 0 latest_scrape = None for row in rows: status = row.availability_status dates_map[row.rate_date.isoformat()] = { 'status': status, 'rate': float(row.rate_gross) if row.rate_gross else None, } if status == 'available': dates_available += 1 elif status == 'sold_out': dates_sold_out += 1 else: dates_no_data += 1 if row.scraped_at and (not latest_scrape or row.scraped_at > latest_scrape): latest_scrape = row.scraped_at return { 'has_own_hotel': True, 'dates_checked': len(rows), 'dates_available': dates_available, 'dates_sold_out': dates_sold_out, 'dates_no_data': dates_no_data, 'latest_scrape': latest_scrape.isoformat() if latest_scrape else None, 'dates': dates_map, } # ============================================ # SCRAPE HISTORY # ============================================ @router.get("/scrape-history") async def get_scrape_history( limit: int = 20, db: AsyncSession = Depends(get_db), current_user: dict = Depends(get_current_user) ): """Get recent scrape batch history.""" result = await db.execute( text(""" SELECT batch_id, scrape_type, started_at, completed_at, status, dates_queued, dates_completed, dates_failed, hotels_found, rates_scraped, error_message, blocked_at, resume_after FROM booking_scrape_log ORDER BY started_at DESC LIMIT :limit """), {'limit': limit} ) return [ { 'batch_id': str(row.batch_id), 'scrape_type': row.scrape_type, 'started_at': row.started_at.isoformat() if row.started_at else None, 'completed_at': row.completed_at.isoformat() if row.completed_at else None, 'status': row.status, 'dates_queued': row.dates_queued, 'dates_completed': row.dates_completed, 'dates_failed': row.dates_failed, 'hotels_found': row.hotels_found, 'rates_scraped': row.rates_scraped, 'error_message': row.error_message, 'blocked_at': row.blocked_at.isoformat() if row.blocked_at else None, 'resume_after': row.resume_after.isoformat() if row.resume_after else None, } for row in result.fetchall() ]