448 lines
15 KiB
Python
448 lines
15 KiB
Python
"""
|
|
AI Daily Insights Generation Job
|
|
|
|
Gathers Pickup-V2 forecast data, booking pace, competitor rates, and rate parity
|
|
information, then sends a compact prompt to Anthropic's Haiku model to generate
|
|
a daily briefing for hotel revenue staff.
|
|
|
|
Schedule: Daily at 7:15 AM (after all forecasts and accuracy calc complete)
|
|
Cost: ~$0.05/month at 1 run/day with Haiku
|
|
"""
|
|
import json
|
|
import logging
|
|
from datetime import date, timedelta, datetime, timezone
|
|
from typing import Dict, List, Any, Optional
|
|
|
|
from sqlalchemy import text
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
|
|
from database import AsyncSessionLocal
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
DEFAULT_MODEL = "claude-haiku-4-5-20251001"
|
|
DEFAULT_DAILY_TOKEN_BUDGET = 5000
|
|
MAX_OUTPUT_TOKENS = 400
|
|
|
|
|
|
async def get_config(db: AsyncSession) -> Dict[str, str]:
|
|
"""Get AI insights config from system_config."""
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT config_key, config_value, COALESCE(is_encrypted, false) as is_encrypted
|
|
FROM system_config
|
|
WHERE config_key LIKE 'ai_insights_%'
|
|
""")
|
|
)
|
|
config = {}
|
|
for row in result.fetchall():
|
|
value = row.config_value
|
|
if row.is_encrypted and value:
|
|
import base64
|
|
try:
|
|
value = base64.b64decode(value.encode()).decode()
|
|
except Exception:
|
|
pass
|
|
config[row.config_key] = value
|
|
return config
|
|
|
|
|
|
async def check_daily_budget(db: AsyncSession, budget: int) -> tuple[bool, int]:
|
|
"""Check if we're within the daily token budget. Returns (within_budget, tokens_used_today)."""
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT COALESCE(SUM(input_tokens + output_tokens), 0) as total
|
|
FROM ai_insights
|
|
WHERE generated_at >= CURRENT_DATE
|
|
""")
|
|
)
|
|
row = result.fetchone()
|
|
used = int(row.total) if row else 0
|
|
return used < budget, used
|
|
|
|
|
|
async def gather_occupancy_data(db: AsyncSession, days: int = 14) -> List[Dict]:
|
|
"""Gather Pickup-V2 occupancy forecast data for the next N days."""
|
|
from services.forecasting.pickup_v2_model import run_pickup_v2_forecast
|
|
|
|
today = date.today()
|
|
end = today + timedelta(days=days - 1)
|
|
|
|
try:
|
|
forecasts = await run_pickup_v2_forecast(
|
|
db, 'hotel_occupancy_pct', today, end, include_details=False
|
|
)
|
|
return forecasts
|
|
except Exception as e:
|
|
logger.warning(f"Failed to gather occupancy data: {e}")
|
|
return []
|
|
|
|
|
|
async def gather_revenue_data(db: AsyncSession, days: int = 14) -> List[Dict]:
|
|
"""Gather Pickup-V2 revenue forecast data for the next N days."""
|
|
from services.forecasting.pickup_v2_model import run_pickup_v2_forecast
|
|
|
|
today = date.today()
|
|
end = today + timedelta(days=days - 1)
|
|
|
|
try:
|
|
forecasts = await run_pickup_v2_forecast(
|
|
db, 'net_accom', today, end, include_details=False
|
|
)
|
|
return forecasts
|
|
except Exception as e:
|
|
logger.warning(f"Failed to gather revenue data: {e}")
|
|
return []
|
|
|
|
|
|
async def gather_budget_data(db: AsyncSession, days: int = 14) -> Dict[str, float]:
|
|
"""Gather budget values for forecast comparison."""
|
|
today = date.today()
|
|
end = today + timedelta(days=days - 1)
|
|
|
|
result = await db.execute(
|
|
text("""
|
|
SELECT date, budget_type, budget_value
|
|
FROM daily_budgets
|
|
WHERE date BETWEEN :start AND :end
|
|
AND budget_type IN ('net_accom', 'occupancy')
|
|
"""),
|
|
{"start": today, "end": end}
|
|
)
|
|
|
|
budgets = {}
|
|
for row in result.fetchall():
|
|
key = f"{row.date}_{row.budget_type}"
|
|
budgets[key] = float(row.budget_value) if row.budget_value else None
|
|
return budgets
|
|
|
|
|
|
async def gather_competitor_data(db: AsyncSession, days: int = 14) -> Dict[str, Any]:
|
|
"""Gather competitor rate data from Booking.com scraper."""
|
|
today = date.today()
|
|
end = today + timedelta(days=days - 1)
|
|
|
|
# Own hotel rate on Booking.com
|
|
own_result = await db.execute(
|
|
text("""
|
|
SELECT DISTINCT ON (r.rate_date)
|
|
r.rate_date,
|
|
r.rate_gross as booking_rate,
|
|
r.availability_status
|
|
FROM booking_com_rates r
|
|
JOIN booking_com_hotels h ON r.hotel_id = h.id
|
|
WHERE h.tier = 'own'
|
|
AND r.rate_date BETWEEN :start AND :end
|
|
ORDER BY r.rate_date, r.scraped_at DESC
|
|
"""),
|
|
{"start": today, "end": end}
|
|
)
|
|
own_rates = {str(row.rate_date): {
|
|
'rate': float(row.booking_rate) if row.booking_rate else None,
|
|
'status': row.availability_status
|
|
} for row in own_result.fetchall()}
|
|
|
|
# Cheapest competitor rate per date
|
|
comp_result = await db.execute(
|
|
text("""
|
|
SELECT DISTINCT ON (r.rate_date)
|
|
r.rate_date,
|
|
r.rate_gross as comp_rate,
|
|
h.name as comp_name
|
|
FROM booking_com_rates r
|
|
JOIN booking_com_hotels h ON r.hotel_id = h.id
|
|
WHERE h.tier IN ('primary', 'secondary')
|
|
AND r.rate_date BETWEEN :start AND :end
|
|
AND r.availability_status = 'available'
|
|
AND r.rate_gross IS NOT NULL
|
|
ORDER BY r.rate_date, r.rate_gross ASC
|
|
"""),
|
|
{"start": today, "end": end}
|
|
)
|
|
comp_rates = {str(row.rate_date): {
|
|
'rate': float(row.comp_rate),
|
|
'name': row.comp_name
|
|
} for row in comp_result.fetchall()}
|
|
|
|
# Own rack rate from Newbook
|
|
rack_result = await db.execute(
|
|
text("""
|
|
SELECT DISTINCT ON (rate_date)
|
|
rate_date,
|
|
rate_gross
|
|
FROM newbook_current_rates
|
|
WHERE rate_date BETWEEN :start AND :end
|
|
ORDER BY rate_date, valid_from DESC
|
|
"""),
|
|
{"start": today, "end": end}
|
|
)
|
|
rack_rates = {str(row.rate_date): float(row.rate_gross) if row.rate_gross else None
|
|
for row in rack_result.fetchall()}
|
|
|
|
return {
|
|
'own_booking': own_rates,
|
|
'competitors': comp_rates,
|
|
'rack': rack_rates
|
|
}
|
|
|
|
|
|
def _fmt_date(d: str) -> str:
|
|
"""Convert YYYY-MM-DD to DD/MM/YYYY for UK display."""
|
|
try:
|
|
parts = str(d).split('-')
|
|
if len(parts) == 3:
|
|
return f"{parts[2]}/{parts[1]}/{parts[0]}"
|
|
except Exception:
|
|
pass
|
|
return str(d)
|
|
|
|
|
|
def build_prompt(
|
|
occupancy: List[Dict],
|
|
revenue: List[Dict],
|
|
budgets: Dict[str, float],
|
|
competitor: Dict[str, Any]
|
|
) -> tuple[str, str]:
|
|
"""Build system and user prompts from gathered data. Returns (system_msg, user_msg)."""
|
|
|
|
system_msg = (
|
|
"You are an AI assistant for a hotel revenue manager in the UK. Analyze the data below and provide "
|
|
"a concise daily briefing (3-5 bullet points). Focus on: occupancy trends, pace vs prior "
|
|
"year, pricing opportunities, and anything unusual requiring attention. "
|
|
"Use UK date format (DD/MM/YYYY) and GBP (£) for all monetary values. "
|
|
"Be specific with numbers and dates. Keep it actionable — no fluff or generic advice."
|
|
)
|
|
|
|
lines = []
|
|
|
|
# Occupancy section
|
|
if occupancy:
|
|
lines.append("## Occupancy Forecast - Pickup-V2 (next 14 days)")
|
|
lines.append("Date | DoW | OTB | Forecast | PY Final | Pace vs LY | Budget")
|
|
for fc in occupancy:
|
|
d = fc.get('date', '')
|
|
dow = fc.get('day_of_week', '')
|
|
otb = fc.get('current_otb')
|
|
forecast = fc.get('forecast')
|
|
py_final = fc.get('prior_year_final')
|
|
pace = fc.get('pace_vs_prior_pct')
|
|
budget_key = f"{d}_occupancy"
|
|
budget_val = budgets.get(budget_key)
|
|
|
|
otb_str = f"{otb:.0f}%" if otb is not None else "-"
|
|
fc_str = f"{forecast:.0f}%" if forecast is not None else "-"
|
|
py_str = f"{py_final:.0f}%" if py_final is not None else "-"
|
|
pace_str = f"{pace:+.0f}%" if pace is not None else "-"
|
|
bud_str = f"{budget_val:.0f}%" if budget_val is not None else "-"
|
|
|
|
lines.append(f"{_fmt_date(d)} | {dow} | {otb_str} | {fc_str} | {py_str} | {pace_str} | {bud_str}")
|
|
lines.append("")
|
|
|
|
# Revenue summary
|
|
if revenue:
|
|
total_forecast = sum(fc.get('forecast', 0) or 0 for fc in revenue)
|
|
total_otb = sum(fc.get('current_otb_rev', 0) or 0 for fc in revenue)
|
|
total_py = sum(fc.get('prior_year_final_rev', 0) or 0 for fc in revenue)
|
|
opportunity_days = sum(1 for fc in revenue if fc.get('has_pricing_opportunity'))
|
|
total_lost = sum(fc.get('lost_potential', 0) or 0 for fc in revenue)
|
|
|
|
lines.append("## Revenue Signals")
|
|
lines.append(f"14-day forecast: £{total_forecast:,.0f} | OTB: £{total_otb:,.0f} | PY: £{total_py:,.0f}")
|
|
if opportunity_days > 0:
|
|
lines.append(f"Pricing opportunity days: {opportunity_days} | Total lost potential: £{total_lost:,.0f}")
|
|
lines.append("")
|
|
|
|
# Competitor rates section
|
|
own_booking = competitor.get('own_booking', {})
|
|
comp_rates = competitor.get('competitors', {})
|
|
rack_rates = competitor.get('rack', {})
|
|
|
|
if own_booking or comp_rates:
|
|
lines.append("## Competitor Rates (next 14 days)")
|
|
lines.append("Date | Own Rack | Own B.com | Cheapest Competitor | Competitor Name")
|
|
|
|
all_dates = sorted(set(list(own_booking.keys()) + list(comp_rates.keys()) + list(rack_rates.keys())))
|
|
for d in all_dates:
|
|
rack = rack_rates.get(d)
|
|
own = own_booking.get(d, {})
|
|
comp = comp_rates.get(d, {})
|
|
|
|
rack_str = f"£{rack:.0f}" if rack else "-"
|
|
own_str = f"£{own['rate']:.0f}" if own.get('rate') else "-"
|
|
comp_str = f"£{comp['rate']:.0f}" if comp.get('rate') else "-"
|
|
comp_name = comp.get('name', '-')
|
|
|
|
note = ""
|
|
if own.get('rate') and comp.get('rate') and own['rate'] < comp['rate']:
|
|
note = " <- cheapest on B.com"
|
|
|
|
lines.append(f"{_fmt_date(d)} | {rack_str} | {own_str} | {comp_str} | {comp_name}{note}")
|
|
lines.append("")
|
|
|
|
user_msg = "\n".join(lines)
|
|
return system_msg, user_msg
|
|
|
|
|
|
async def call_llm(api_key: str, system_msg: str, user_msg: str, model: str) -> Dict[str, Any]:
|
|
"""Call Anthropic API and return response with token usage."""
|
|
import anthropic
|
|
|
|
client = anthropic.AsyncAnthropic(api_key=api_key)
|
|
|
|
try:
|
|
response = await client.messages.create(
|
|
model=model,
|
|
max_tokens=MAX_OUTPUT_TOKENS,
|
|
temperature=0.2,
|
|
system=system_msg,
|
|
messages=[{"role": "user", "content": user_msg}]
|
|
)
|
|
|
|
content = response.content[0].text if response.content else ""
|
|
return {
|
|
"content": content,
|
|
"input_tokens": response.usage.input_tokens,
|
|
"output_tokens": response.usage.output_tokens,
|
|
"model": model,
|
|
}
|
|
finally:
|
|
await client.close()
|
|
|
|
|
|
async def save_insight(
|
|
db: AsyncSession,
|
|
content: str,
|
|
model: str,
|
|
input_tokens: int,
|
|
output_tokens: int,
|
|
data_snapshot: Dict,
|
|
triggered_by: str = "scheduler"
|
|
):
|
|
"""Save generated insight to database."""
|
|
await db.execute(
|
|
text("""
|
|
INSERT INTO ai_insights
|
|
(content, model, input_tokens, output_tokens, data_snapshot, triggered_by)
|
|
VALUES (:content, :model, :input_tokens, :output_tokens,
|
|
CAST(:data_snapshot AS jsonb), :triggered_by)
|
|
"""),
|
|
{
|
|
"content": content,
|
|
"model": model,
|
|
"input_tokens": input_tokens,
|
|
"output_tokens": output_tokens,
|
|
"data_snapshot": json.dumps(data_snapshot),
|
|
"triggered_by": triggered_by,
|
|
}
|
|
)
|
|
await db.commit()
|
|
|
|
|
|
async def cleanup_old_insights(db: AsyncSession, keep_days: int = 90):
|
|
"""Remove insights older than keep_days."""
|
|
cutoff = datetime.now(timezone.utc) - timedelta(days=keep_days)
|
|
await db.execute(
|
|
text("DELETE FROM ai_insights WHERE generated_at < :cutoff"),
|
|
{"cutoff": cutoff}
|
|
)
|
|
await db.commit()
|
|
|
|
|
|
async def generate_insight(db: AsyncSession, triggered_by: str = "scheduler") -> Dict[str, Any]:
|
|
"""
|
|
Core insight generation logic. Used by both scheduler and manual trigger.
|
|
Returns result dict with success/error status.
|
|
"""
|
|
config = await get_config(db)
|
|
|
|
# Check enabled
|
|
if config.get('ai_insights_enabled', 'false').lower() not in ('true', '1', 'yes'):
|
|
return {"success": False, "error": "AI insights disabled"}
|
|
|
|
# Check API key
|
|
api_key = config.get('ai_insights_api_key')
|
|
if not api_key:
|
|
return {"success": False, "error": "No API key configured"}
|
|
|
|
model = config.get('ai_insights_model', DEFAULT_MODEL)
|
|
budget = int(config.get('ai_insights_daily_token_budget', str(DEFAULT_DAILY_TOKEN_BUDGET)))
|
|
|
|
# Check daily budget
|
|
within_budget, used_today = await check_daily_budget(db, budget)
|
|
if not within_budget:
|
|
return {
|
|
"success": False,
|
|
"error": f"Daily token budget exceeded ({used_today}/{budget} tokens used today)"
|
|
}
|
|
|
|
# Gather data
|
|
logger.info("Gathering data for AI insight...")
|
|
occupancy = await gather_occupancy_data(db)
|
|
revenue = await gather_revenue_data(db)
|
|
budgets = await gather_budget_data(db)
|
|
competitor = await gather_competitor_data(db)
|
|
|
|
if not occupancy and not revenue:
|
|
return {"success": False, "error": "No forecast data available"}
|
|
|
|
# Build prompt
|
|
system_msg, user_msg = build_prompt(occupancy, revenue, budgets, competitor)
|
|
|
|
# Store data snapshot for debugging
|
|
data_snapshot = {
|
|
"occupancy_days": len(occupancy),
|
|
"revenue_days": len(revenue),
|
|
"competitor_dates": len(competitor.get('own_booking', {})),
|
|
"prompt_preview": user_msg[:500],
|
|
}
|
|
|
|
# Call LLM
|
|
logger.info(f"Calling {model} for AI insight...")
|
|
try:
|
|
result = await call_llm(api_key, system_msg, user_msg, model)
|
|
except Exception as e:
|
|
logger.error(f"LLM call failed: {e}")
|
|
return {"success": False, "error": f"LLM call failed: {str(e)}"}
|
|
|
|
# Save
|
|
await save_insight(
|
|
db,
|
|
content=result["content"],
|
|
model=result["model"],
|
|
input_tokens=result["input_tokens"],
|
|
output_tokens=result["output_tokens"],
|
|
data_snapshot=data_snapshot,
|
|
triggered_by=triggered_by,
|
|
)
|
|
|
|
logger.info(
|
|
f"AI insight generated: {result['input_tokens']} input, "
|
|
f"{result['output_tokens']} output tokens ({triggered_by})"
|
|
)
|
|
|
|
# Cleanup old insights
|
|
try:
|
|
await cleanup_old_insights(db)
|
|
except Exception as e:
|
|
logger.warning(f"Old insight cleanup failed: {e}")
|
|
|
|
return {
|
|
"success": True,
|
|
"content": result["content"],
|
|
"input_tokens": result["input_tokens"],
|
|
"output_tokens": result["output_tokens"],
|
|
"model": result["model"],
|
|
}
|
|
|
|
|
|
async def run_ai_insights_generation():
|
|
"""Scheduled job entry point."""
|
|
async with AsyncSessionLocal() as db:
|
|
try:
|
|
result = await generate_insight(db, triggered_by="scheduler")
|
|
if result.get("success"):
|
|
logger.info("Scheduled AI insight generation completed")
|
|
else:
|
|
logger.info(f"Scheduled AI insight skipped: {result.get('error')}")
|
|
except Exception as e:
|
|
logger.error(f"Scheduled AI insight generation failed: {e}", exc_info=True)
|