from fastapi import APIRouter, Body, UploadFile, File, Form, HTTPException, Query, Request, Response
from fastapi.responses import RedirectResponse
from typing import List, Optional
from app.services.infra.database import (
    bootstrap_tenant, insert_vector_data, insert_summary, get_latest_summary, get_tickets, get_ticket_status_counts, 
    delete_vector_data_by_source, get_llm_usage, get_knowledge_sources, add_ticket_comment, get_ticket_comments,
    insert_crawled_url, update_ticket_status, get_ticket, delete_vector_data_by_url_subpath, get_source_ids_by_name
)
from app.services.infra.storage import storage_service
from app.services.infra.embedding import chunk_text, get_embeddings, generate_summary, generate_title
from app.services.content.website import extract_text_from_url
from app.services.chat.chat import get_chat_response
from app.services.integrations.firestore import get_all_threads, get_chat_history
from app.services.chat.analytics import get_kpis, get_general_kpis
from app.services.content.documents import (
    DEFAULT_CATEGORY, blob_url, delete_document, download_link, find_document,
    ingest_documents, list_documents, set_category,
)
from app.services.chat.tools_settings import (
    EMAIL_COPY_DEFAULTS, FEATURE_DEFAULTS, FEATURE_TOOLS, PROVIDER_OPTIONS,
    get_tool_settings, normalise_tenant_id, save_tool_settings,
)
from app.services.infra.quotas import QuotaExceededError
from app.services.enrichment.job import enrich_tenant
from app.services.pairing.recommendations import VALID_EVENTS, decide
from app.services.catalog.sync import has_any_products, sync_source
from app.services.catalog.sources import KIND_CRAWL, get_sources
from app.services.infra import jobs
from app.api.catalog import _run, start_job
from app.core.config import settings
import anyio
import asyncio
import json
import logging
import uuid
import re
import httpx
from datetime import datetime, timedelta, timezone

logger = logging.getLogger(__name__)

router = APIRouter()

@router.post("/bootstrap")
async def bootstrap(tenant_id: str):
    logger.info(f"Received bootstrap request for tenant: {tenant_id}")
    try:
        bootstrap_tenant(tenant_id)
        return {"status": "success", "message": f"Tenant {tenant_id} setup complete."}
    except ValueError as ve:
        logger.error(f"Validation error in bootstrap for tenant {tenant_id}: {ve}", exc_info=True)
        raise HTTPException(status_code=400, detail={"status": "error", "message": f"Validation failed: {str(ve)}"})
    except Exception as e:
        logger.error(f"Error in bootstrap endpoint for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Internal server error during tenant setup."})

@router.post("/ingest/files")
async def ingest_files(
    tenant_id: str = Form(...),
    user_id: Optional[str] = Form(None),
    category: Optional[str] = Form(None),
    files: List[UploadFile] = File(...)
):
    logger.info(f"Received ingest/files request for tenant: {tenant_id}, count: {len(files)}")
    try:
        bootstrap_tenant(tenant_id) # Ensure schema exists
        # Shared with POST /documents so the two upload paths cannot drift.
        uploads = [(f.filename, await f.read()) for f in files]
        result = ingest_documents(tenant_id, uploads, category=category, user_id=user_id)
        summary = result["summary"]

        now = datetime.utcnow()
        return {
            "status": "success",
            "message": f"Ingested {len(files)} files into {tenant_id}.",
            "summary": summary,
            "documents": result["documents"],
            "ingested_date": now.date().isoformat(),
            "ingested_time": now.time().strftime("%H:%M:%S"),
            "created_at": now.isoformat() + "Z"
        }
    except QuotaExceededError:
        raise
    except ValueError as ve:
        logger.error(f"Validation/Extraction error for tenant {tenant_id}: {ve}", exc_info=True)
        raise HTTPException(status_code=400, detail={"status": "error", "message": f"File ingestion failed: {str(ve)}"})
    except Exception as e:
        logger.error(f"Unexpected error in ingest/files for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "An error occurred during file ingestion."})

@router.post("/ingest/text")
async def ingest_text(
    tenant_id: str = Form(...),
    text_content: str = Form(...),
    user_id: Optional[str] = Form(None)
):
    logger.info(f"Received ingest/text request for tenant: {tenant_id}")
    if not text_content.strip():
        raise HTTPException(status_code=400, detail="text_content cannot be empty.")
        
    try:
        bootstrap_tenant(tenant_id) # Ensure schema exists
        source_id = str(uuid.uuid4())
        
        # Check for explicit Title/Content format
        title_match = re.search(r'^Title:\s*(.+?)\s*Content:\s*(.*)', text_content, re.IGNORECASE | re.DOTALL)
        if title_match:
            display_title = title_match.group(1).strip()
            # Clean up title for URL/filename safety
            display_title = re.sub(r'[^\w\s-]', '', display_title).strip().replace(' ', '_')
            if not display_title:
                display_title = "Document"
        else:
            display_title = generate_title(text_content, tenant_id, user_id=user_id)
            
        unique_source_name = f"{display_title}.txt"
        
        # Replacement Logic: Clear old data for this text title
        old_source_ids = get_source_ids_by_name(tenant_id, unique_source_name, "text")
        for sid in old_source_ids:
            logger.info(f"Replacing existing text {unique_source_name} (sid: {sid}) in tenant {tenant_id}")
            delete_vector_data_by_source(tenant_id, sid)

        chunks = chunk_text(text_content)
        metadata = {"source": unique_source_name, "source_id": source_id, "ingestion_type": "text"}
        if user_id:
            metadata["user_id"] = user_id
        logger.info(f"Created {len(chunks)} chunks for text input '{unique_source_name}'")
        for chunk in chunks:
            embedding = get_embeddings(chunk, tenant_id=tenant_id, user_id=user_id)
            insert_vector_data(tenant_id, chunk, embedding, metadata)
            
        # Generate and store summary
        summary = generate_summary(text_content, tenant_id=tenant_id, user_id=user_id)
        insert_summary(tenant_id, "text", summary)
        
        now = datetime.utcnow()
        return {
            "status": "success", 
            "message": f"Ingested text content into {tenant_id}.",
            "summary": summary,
            "ingested_date": now.date().isoformat(),
            "ingested_time": now.time().strftime("%H:%M:%S"),
            "created_at": now.isoformat() + "Z"
        }
    except QuotaExceededError:
        raise
    except Exception as e:
        logger.error(f"Error in ingest/text for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "An error occurred during text ingestion."})

@router.post("/ingest/website")
async def ingest_website(
    tenant_id: str = Form(...),
    url: str = Form(...),
    user_id: Optional[str] = Form(None)
):
    logger.info(f"Received ingest/website request for tenant: {tenant_id}, URL: {url}")
    
    # Early Validation Check
    try:
        async with httpx.AsyncClient(timeout=5.0, verify=False, follow_redirects=True) as client:
            resp = await client.head(url)
            # Some sites reject HEAD, so we'll fallback to GET if needed, but raise if clearly unreachable/error
            if resp.status_code >= 400 and resp.status_code != 405:
                # 405 Method Not Allowed is common for HEAD requests, we ignore it.
                logger.warning(f"URL Validation failed for {url} with status code {resp.status_code}")
                raise HTTPException(status_code=400, detail=f"Website validation failed: Received status {resp.status_code} from {url}")
    except httpx.RequestError as e:
        logger.error(f"URL Validation error for {url}: {e}")
        raise HTTPException(status_code=400, detail=f"Website validation failed: Could not connect to {url}")
    except HTTPException:
        raise
        
    try:
        bootstrap_tenant(tenant_id) # Ensure schema exists
        
        # Replacement Logic: Clear old data for this URL and its sub-paths
        logger.info(f"Applying replacement logic for {url} in tenant {tenant_id}")
        delete_vector_data_by_url_subpath(tenant_id, url)
        
        scrape_result = await extract_text_from_url(url, max_depth=20, max_pages=settings.CRAWL_MAX_PAGES, tenant_id=tenant_id)
        
        extracted_text = scrape_result["text"]
        if not extracted_text.strip():
            logger.error(f"No content could be extracted from {url}")
            raise HTTPException(status_code=400, detail=f"Failed to extract any content from the provided URL: {url}")

        source_id = str(uuid.uuid4())
        metadata_base = {
            "source": url,
            "source_id": source_id,
            "ingestion_type": "website",
            "urls_visited": scrape_result["urls_visited"],
            "colors": scrape_result["colors"],
            "top_images": scrape_result["images"]
        }
        if user_id:
            metadata_base["user_id"] = user_id
        
        # Store all visited URLs in the structured table
        for visited_url in scrape_result["urls_visited"]:
            insert_crawled_url(tenant_id, source_id, visited_url)
        
        total_chunks = 0
        for page in scrape_result.get("pages", []):
            page_text = page["text"]
            page_url = page["url"]
            
            # Chunk the page text specifically
            page_chunks = chunk_text(page_text)
            
            for chunk in page_chunks:
                # Prefix every chunk with its specific source URL inside the content
                referenced_chunk = f"[Source: {page_url}]\n{chunk}"
                
                # Create specific metadata for this chunk
                chunk_metadata = metadata_base.copy()
                chunk_metadata["source_url"] = page_url # Store deep link
                chunk_metadata["source"] = page_url     # Update primary source to deep link
                
                # Use page-specific images if available
                if "images" in page:
                    chunk_metadata["top_images"] = page["images"]
                
                embedding = get_embeddings(referenced_chunk, tenant_id=tenant_id, user_id=user_id)
                insert_vector_data(tenant_id, referenced_chunk, embedding, chunk_metadata)
                total_chunks += 1
                
        logger.info(f"Created total of {total_chunks} chunks for website crawl starting at {url} (across {len(scrape_result['pages'])} pages) with source_id {source_id}")
            
        # Generate and store summary
        summary = generate_summary(extracted_text, tenant_id=tenant_id, user_id=user_id)
        insert_summary(tenant_id, "website", summary)

        # This crawl builds the knowledge base and nothing else. Products come
        # from the catalogue's own website source, which registers the site,
        # runs on the build, and writes rows through normalise() -- with a
        # price, a currency, a taxonomy and a source_ref, none of which the
        # extraction here ever produced.
        #
        # Leaving it in gave one table two writers with incompatible ideas of
        # what a crawled product is, and the delete below it was scoped to
        # source_kind alone, so each would reap the other's rows.

        now = datetime.utcnow()
        metadata_base["ingested_date"] = now.date().isoformat()
        metadata_base["ingested_time"] = now.time().strftime("%H:%M:%S")
        metadata_base["created_at"] = now.isoformat() + "Z"

        return {
            "status": "success", 
            "message": f"Ingested website content from {url} (visited {len(scrape_result['urls_visited'])} pages) into {tenant_id}.",
            "summary": summary,
            "metadata": metadata_base,
            "ingested_date": metadata_base["ingested_date"],
            "ingested_time": metadata_base["ingested_time"],
            "created_at": metadata_base["created_at"]
        }
    except QuotaExceededError:
        raise
    except ValueError as ve:
        logger.error(f"Website ingestion validation failed for tenant {tenant_id}: {ve}", exc_info=True)
        raise HTTPException(status_code=400, detail={"status": "error", "message": f"Website ingestion failed: {str(ve)}"})
    except Exception as e:
        logger.error(f"Unexpected website ingestion error for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "An unexpected error occurred while crawling the website."})

@router.delete("/ingest/source")
async def delete_ingested_source(
    tenant_id: str,
    source_id: str,
    source_name: str = Query(None, description="The original filename to delete from Azure Blob Storage")
):
    """
    Deletes all data associated with a specific source_id from the knowledge base.
    """
    logger.info(f"Received request to delete source_id '{source_id}' for tenant {tenant_id}")
    try:
        # If it's a file, delete from Azure first. Use the combined UUID_filename format.
        if source_name:
            azure_blob_name = f"{source_id}_{source_name}"
            storage_service.delete_file(tenant_id, azure_blob_name)
        
        deleted_count = delete_vector_data_by_source(tenant_id, source_id)
        return {
            "status": "success",
            "message": f"Deleted {deleted_count} chunks for source_id '{source_id}' from {tenant_id} and Azure storage.",
            "deleted_count": deleted_count
        }
    except Exception as e:
        logger.error(f"Error deleting source_id {source_id} for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "An error occurred during data deletion."})

@router.get("/summary/{type}")
async def get_summary(tenant_id: str, type: str):
    logger.info(f"Received summary fetch request for tenant: {tenant_id}, type: {type}")
    if type not in ["files", "text", "website"]:
        raise HTTPException(status_code=400, detail="Invalid summary type. Use 'files', 'text', or 'website'.")
    
    try:
        summary, dt = get_latest_summary(tenant_id, type)
        if not summary:
            logger.warning(f"No summary found for {tenant_id} type {type}")
            raise HTTPException(status_code=404, detail="No summary found for this type and tenant.")
            
        return {
            "tenant_id": tenant_id, 
            "type": type, 
            "summary": summary,
            "ingested_date": dt.date().isoformat() if dt else None,
            "ingested_time": dt.time().strftime("%H:%M:%S") if dt else None,
            "created_at": dt.isoformat() if dt else None
        }
    except HTTPException as he:
        raise he
    except Exception as e:
        logger.error(f"Error fetching summary for tenant {tenant_id}, type {type}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Internal error fetching summary."})

@router.get("/ingest/files")
async def list_ingested_files(tenant_id: str):
    logger.info(f"Listing ingested files for tenant: {tenant_id}")
    try:
        sources = get_knowledge_sources(tenant_id, "file")
        return {"status": "success", "tenant_id": tenant_id, "count": len(sources), "files": sources}
    except Exception as e:
        logger.error(f"Error listing files for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error listing ingested files."})

@router.get("/ingest/websites")
async def list_ingested_websites(tenant_id: str):
    logger.info(f"Listing ingested websites for tenant: {tenant_id}")
    try:
        sources = get_knowledge_sources(tenant_id, "website")
        return {"status": "success", "tenant_id": tenant_id, "count": len(sources), "websites": sources}
    except Exception as e:
        logger.error(f"Error listing websites for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error listing ingested websites."})

@router.get("/ingest/text")
async def list_ingested_text(tenant_id: str):
    logger.info(f"Listing ingested text for tenant: {tenant_id}")
    try:
        sources = get_knowledge_sources(tenant_id, "text")
        return {"status": "success", "tenant_id": tenant_id, "count": len(sources), "text_sources": sources}
    except Exception as e:
        logger.error(f"Error listing text for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error listing ingested text."})

@router.post("/documents")
async def upload_documents_endpoint(
    tenant_id: str = Form(...),
    category: Optional[str] = Form(None),
    user_id: Optional[str] = Form(None),
    files: List[UploadFile] = File(...),
):
    """Upload documents under a category.

    Same indexing as /ingest/files, but the response carries each stored
    document's source_id, so the settings page can offer delete and
    re-categorise straight away without a second listing request.
    """
    try:
        tenant_id = normalise_tenant_id(tenant_id)
        bootstrap_tenant(tenant_id)
        uploads = [(f.filename, await f.read()) for f in files]
        result = ingest_documents(tenant_id, uploads, category=category, user_id=user_id)
        return {"status": "success",
                "message": f"Stored {len(result['documents'])} document(s).",
                **result}
    except QuotaExceededError:
        raise
    except ValueError as ve:
        logger.error(f"Extraction error for tenant {tenant_id}: {ve}", exc_info=True)
        raise HTTPException(status_code=400, detail={"status": "error",
                                                     "message": f"Upload failed: {ve}"})
    except Exception as e:
        logger.error(f"Error uploading documents for {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not store the documents."})


@router.delete("/documents")
async def delete_document_endpoint(tenant_id: str, source_id: str):
    """Delete a document: its knowledge base rows and its stored file."""
    try:
        tenant_id = normalise_tenant_id(tenant_id)
        removed = delete_document(tenant_id, source_id)
        if not removed:
            raise HTTPException(status_code=404, detail={"status": "error",
                                                         "message": "Document not found."})
        return {"status": "success", **removed}
    except HTTPException:
        raise
    except Exception as e:
        logger.error(f"Error deleting document {source_id} for {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not delete the document."})


@router.get("/documents")
async def list_documents_endpoint(tenant_id: str, category: Optional[str] = None):
    """The tenant's uploaded documents, for the Attach Documents settings card."""
    try:
        tenant_id = normalise_tenant_id(tenant_id)
        documents = list_documents(tenant_id, category)
        return {"status": "success", "count": len(documents), "documents": documents}
    except Exception as e:
        logger.error(f"Error listing documents for {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not list documents."})


@router.post("/documents/category")
async def set_document_category_endpoint(
    tenant_id: str = Form(...),
    source_id: str = Form(...),
    category: str = Form(...),
):
    """Re-categorise an already uploaded document."""
    try:
        tenant_id = normalise_tenant_id(tenant_id)
        if not set_category(tenant_id, source_id, category):
            raise HTTPException(status_code=404, detail={"status": "error",
                                                         "message": "Document not found."})
        return {"status": "success", "source_id": source_id,
                "category": category.strip() or DEFAULT_CATEGORY}
    except HTTPException:
        raise
    except Exception as e:
        logger.error(f"Error setting category for {source_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not update the category."})


@router.get("/documents/download")
async def document_download_endpoint(tenant_id: str, source_id: str = None,
                                     file_name: str = None):
    """Download a document. This is the URL handed to visitors in chat.

    It redirects to a storage link generated right now, so the link in a chat
    transcript never goes stale and the storage account never appears in a chat
    payload. Identify the document by `source_id`; `file_name` also works.
    """
    if not (source_id or file_name):
        raise HTTPException(status_code=400, detail={"status": "error",
                                                     "message": "source_id or file_name is required."})
    try:
        tenant_id = normalise_tenant_id(tenant_id)
        document = (find_document(tenant_id, source_id) if source_id
                    else find_document(tenant_id, (download_link(tenant_id, file_name) or {})
                                       .get("source_id")))
        url = blob_url(tenant_id, document)
        if not url:
            raise HTTPException(status_code=404, detail={"status": "error",
                                                         "message": "Document not found."})
        # 302, not 301: the storage link differs every time, so nothing about
        # this response may be cached by a browser or proxy.
        return RedirectResponse(url, status_code=302)
    except HTTPException:
        raise
    except Exception as e:
        logger.error(f"Error serving document {source_id or file_name}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not serve the document."})


@router.get("/tools-settings")
async def get_tools_settings_endpoint(tenant_id: str):
    """Every capability switch for a tenant, plus the email copy.

    `features` carries all eight switches so the settings page can render them
    without knowing which ones the backend implements. `wired` names the subset
    that currently changes chat behaviour; the rest are stored preferences whose
    tools do not exist yet.
    """
    try:
        settings = get_tool_settings(tenant_id)
        return {"status": "success",
                "settings": settings,
                "wired": sorted(FEATURE_TOOLS)}
    except Exception as e:
        logger.error(f"Error fetching tool settings for {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not load tool settings."})


def _parse_settings_form(form) -> tuple:
    """Read the flat form shape the settings page posts.

    Fields are named after the switches themselves (voice_mode=false), plus the
    three email strings. `enabled` is the old name for `send_email`.
    """
    def as_bool(raw):
        return str(raw).strip().lower() in ("1", "true", "yes", "on")

    features = {name: as_bool(form[name]) for name in FEATURE_DEFAULTS if name in form}
    if "enabled" in form:
        features.setdefault("send_email", as_bool(form["enabled"]))
    email = {name: form[name] for name in EMAIL_COPY_DEFAULTS if name in form}
    # Posted as e.g. ticket_management_provider=galaxiq_ticketmanagement, kept
    # apart from the switches above so a form field name never collides with a
    # feature of the same name.
    providers = {name: form[f"{name}_provider"] for name in PROVIDER_OPTIONS
                if f"{name}_provider" in form}
    return features, email, providers


@router.post("/tools-settings")
async def save_tools_settings_endpoint(request: Request):
    """Update capability switches and/or email copy. Only what is sent changes.

    Accepts either JSON:
        {"tenant_id": "...", "features": {"voice_mode": false},
         "email": {"header": "..."}, "ticket_emails": [{"email": "...", "category": "..."}]}
    or form fields, which is what the settings page already posts:
        tenant_id=...&voice_mode=false&header=...

    Both are supported so moving off the old /email-settings URL is a one-line
    change on the caller rather than a rewrite of how it posts. `ticket_emails`
    is JSON-only: a list of objects has no natural flat-form encoding and
    nothing in this codebase's form-based caller needs it.
    """
    if request.headers.get("content-type", "").startswith("application/json"):
        payload = await request.json()
        if not isinstance(payload, dict):
            raise HTTPException(status_code=400, detail={"status": "error",
                                                         "message": "Body must be an object."})
        tenant_id = (payload.get("tenant_id") or "").strip()
        features, email = payload.get("features") or {}, payload.get("email") or {}
        providers = payload.get("providers") or {}
        ticket_emails = payload.get("ticket_emails")
    else:
        form = await request.form()
        tenant_id = (form.get("tenant_id") or "").strip()
        features, email, providers = _parse_settings_form(form)
        ticket_emails = None

    if not tenant_id:
        raise HTTPException(status_code=400, detail={"status": "error",
                                                     "message": "tenant_id is required."})

    # A switch sent as null means "leave it alone", not "turn it off". A
    # provider sent as null is different -- None is how a choice is cleared,
    # so it must reach save_tool_settings rather than being filtered here.
    features = {k: v for k, v in features.items() if v is not None}
    email = {k: v for k, v in email.items() if v is not None}
    if not features and not email and not providers and ticket_emails is None:
        raise HTTPException(status_code=400, detail={"status": "error",
                                                     "message": "No settings supplied."})
    try:
        # Normalise first: bootstrapping a bare uuid would create a second,
        # empty schema that the chat never reads from.
        bootstrap_tenant(normalise_tenant_id(tenant_id))
        # ticket_emails is passed through as-is, NOT `ticket_emails or None`:
        # an empty list must reach save_tool_settings to clear the stored
        # list, but `[] or None` would silently turn that into "leave it
        # alone". It is already None when the key was absent from the payload.
        settings = save_tool_settings(tenant_id, features=features or None,
                                      email=email or None,
                                      providers=providers or None,
                                      ticket_emails=ticket_emails)
        return {"status": "success", "settings": settings, "wired": sorted(FEATURE_TOOLS)}
    except ValueError as ve:
        raise HTTPException(status_code=400, detail={"status": "error", "message": str(ve)})
    except Exception as e:
        logger.error(f"Error saving tool settings for {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error",
                                                     "message": "Could not save tool settings."})


@router.post("/recommendations/event")
async def recommendations_event_endpoint(payload: dict = Body(...)):
    """The widget reports page activity; we decide whether to suggest anything.

    Answers HTTP 200 with `recommend: false` rather than an error whenever there
    is nothing to show, because a suggestion the visitor never asked for must
    never surface as a failure in their browser.
    """
    tenant_id = (payload.get("tenant_id") or "").strip()
    visitor_id = (payload.get("visitor_id") or "").strip()
    event = (payload.get("event") or "").strip()
    page = payload.get("page") or {}

    if not tenant_id or not visitor_id or not page.get("url"):
        raise HTTPException(status_code=400, detail={
            "status": "error",
            "message": "tenant_id, visitor_id and page.url are required."})
    if event not in VALID_EVENTS:
        raise HTTPException(status_code=400, detail={
            "status": "error",
            "message": f"Unknown event. Expected one of: {', '.join(sorted(VALID_EVENTS))}"})

    try:
        outcome = await anyio.to_thread.run_sync(
            lambda: decide(normalise_tenant_id(tenant_id), event, page))
    except Exception as ex:
        # decide() already never raises; this is defense in depth so that even
        # an unforeseen failure here answers 200 rather than 500 -- a
        # suggestion nobody asked for must never surface as an error.
        logger.error(f"Recommendation endpoint failed for {visitor_id} on "
                    f"{page.get('url')}: {ex}", exc_info=True)
        return {"status": "success", "recommend": False, "reason": "no_match"}

    if not outcome.get("recommend"):
        logger.info(f"No recommendation for {visitor_id} on {page.get('url')}: "
                    f"{outcome.get('reason')} (title={page.get('title')!r})")
        return {"status": "success", **outcome}

    recommendation_id = f"rec_{uuid.uuid4().hex[:10]}"
    logger.info(f"{recommendation_id}: {len(outcome['products'])} products for "
                f"{visitor_id} on {page.get('url')}")
    return {"status": "success", "recommendation_id": recommendation_id, **outcome}


@router.post("/chat")
async def chat(
    tenant_id: str = Form(...),
    query: str = Form(...),
    thread_id: str = Form("default"),
    user_id: str = Form(None),
    user_name: str = Form(None),
    role: str = Form("user"),
    location: Optional[str] = Form(None)
):
    logger.info(f"Received chat request for tenant {tenant_id}, user_id: {user_id}, role: {role}, location: {location}")
    try:
        full_response = await get_chat_response(tenant_id, query, thread_id, user_name, role, user_id=user_id, location=location)
        return {
            "status": "success",
            "response": full_response.get("text", ""),
            "resources": full_response.get("resources", []),
            "tools": full_response.get("tools", []),
            "products": full_response.get("products", []),
            "thread_id": thread_id,
            "user_id": user_id
        }
    except QuotaExceededError as qe:
        return {
            "status": "error",
            "response": qe.message,
            "resources": {"urls": [], "images": []},
            "thread_id": thread_id,
            "user_id": user_id
        }
    except Exception as e:
        logger.error(f"Chat generation failed for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Failed to generate chat response. Please try again later."})

@router.get("/chat/threads")
async def list_threads(tenant_id: str, user_id: str = Query(None)):
    logger.info(f"Listing threads for tenant: {tenant_id}, user_id: {user_id}")
    try:
        threads = get_all_threads(tenant_id, user_id=user_id)
        return {
            "status": "success",
            "count": len(threads),
            "threads": threads
        }
    except Exception as e:
        logger.error(f"Error in list_threads endpoint for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error retrieving threads."})

@router.get("/chat/history")
async def fetch_chat_history(tenant_id: str, thread_id: str = "default", user_id: str = Query(None)):
    logger.info(f"Fetching history for tenant: {tenant_id}, user: {user_id}, thread: {thread_id}")
    try:
        messages = get_chat_history(tenant_id, thread_id, user_id=user_id)
        return {
            "status": "success",
            "thread_id": thread_id,
            "user_id": user_id,
            "messages": messages
        }
    except Exception as e:
        logger.error(f"Error in fetch_chat_history endpoint for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error retrieving chat history."})

@router.get("/tickets")
async def fetch_tickets(tenant_id: str):
    logger.info(f"Received request to fetch tickets for tenant: {tenant_id}")
    try:
        tickets = get_tickets(tenant_id)
        status_counts = get_ticket_status_counts(tenant_id)
        return {
            "status": "success",
            "count": len(tickets),
            "status_counts": status_counts,
            "tickets": tickets
        }
    except Exception as e:
        logger.error(f"Error in fetch_tickets endpoint for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error retrieving tickets."})

@router.post("/tickets/{ticket_id}/comments")
async def post_ticket_comment(
    ticket_id: str,
    tenant_id: str = Form(...),
    admin_name: str = Form("Admin"),
    comment: str = Form(...)
):
    logger.info(f"Received request to post comment on ticket {ticket_id} for tenant: {tenant_id}")
    try:
        new_comment = add_ticket_comment(tenant_id, ticket_id, admin_name, comment)
        
        # Real-time synchronization: Notify dashboard of the new comment
        try:
            from app.services.infra.redis import publish_event
            payload = {
                "type": "system",
                "event": "ticket_comment_added",
                "message": f"New comment added to ticket {ticket_id} by {admin_name}",
                "data": new_comment
            }
            await publish_event(f"tenant_{tenant_id}:tickets", payload)
        except Exception as ws_err:
            logger.warning(f"Failed to publish ticket comment to redis: {ws_err}")

        return {
            "status": "success",
            "message": "Comment added successfully",
            "comment": new_comment
        }
    except Exception as e:
        logger.error(f"Error posting comment to ticket {ticket_id} for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Failed to add comment."})

@router.get("/tickets/{ticket_id}/comments")
async def fetch_ticket_comments(ticket_id: str, tenant_id: str = Query(...)):
    logger.info(f"Received request to fetch comments for ticket {ticket_id} for tenant: {tenant_id}")
    try:
        comments = get_ticket_comments(tenant_id, ticket_id)
        return {
            "status": "success",
            "count": len(comments),
            "comments": comments
        }
    except Exception as e:
        logger.error(f"Error fetching comments for ticket {ticket_id} for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error retrieving comments."})

@router.post("/tickets/{ticket_id}/status")
async def update_ticket_status_endpoint(
    ticket_id: str,
    tenant_id: str = Form(...),
    status: str = Form(...)
):
    logger.info(f"Received request to update status of ticket {ticket_id} to {status} for tenant: {tenant_id}")
    
    allowed_statuses = ["In Progress", "Resolved", "Ongoing", "Open"]
    if status not in allowed_statuses:
        raise HTTPException(status_code=400, detail={"status": "error", "message": f"Invalid status. Must be one of {allowed_statuses}"})
        
    try:
        success = update_ticket_status(tenant_id, ticket_id, status)
        if not success:
            raise HTTPException(status_code=404, detail={"status": "error", "message": "Ticket not found or update failed."})
            
        # Real-time synchronization: Notify dashboard of the status change
        try:
            from app.services.infra.redis import publish_event
            ticket_data = get_ticket(tenant_id, ticket_id)
            payload = {
                "type": "system",
                "event": "ticket_status_updated",
                "message": f"Ticket {ticket_id} status updated to {status}",
                "data": ticket_data or {"ticket_id": ticket_id, "status": status}
            }
            await publish_event(f"tenant_{tenant_id}:tickets", payload)
        except Exception as ws_err:
            logger.warning(f"Failed to publish ticket status update to redis: {ws_err}")

        return {
            "status": "success",
            "message": f"Ticket status updated to {status}",
            "ticket_id": ticket_id,
            "new_status": status
        }
    except Exception as e:
        logger.error(f"Error updating status for ticket {ticket_id} for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Failed to update ticket status."})

@router.get("/analytics/usage")
async def fetch_llm_usage(tenant_id: str = Query(...), limit: int = Query(100)):
    """
    Returns LLM token usage history for a tenant.
    """
    try:
        usage_history = get_llm_usage(tenant_id, limit)
        return {
            "status": "success",
            "tenant_id": tenant_id,
            "count": len(usage_history),
            "usage": usage_history
        }
    except Exception as e:
        logger.error(f"Error in fetch_llm_usage endpoint for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Error retrieving usage analytics."})

@router.get("/analytics/refresh")
async def refresh_analytics(tenant_id: str):
    """
    Triggers the KPI calculation and returns the latest metrics.
    """
    logger.info(f"Received analytics refresh request for tenant: {tenant_id}")
    try:
        kpis = await get_kpis(tenant_id)
        return {
            "status": "success",
            "data": kpis
        }
    except QuotaExceededError:
        raise
    except Exception as e:
        logger.error(f"Analytics refresh failed for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Failed to compute analytics KPIs."})

@router.get("/analytics/general")
async def fetch_general_analytics(
    tenant_id: str, 
    time_range: str = Query("all", description="Time range (e.g., 1h, 24h, 7d, 30d, all)"),
    start_date: str = Query(None, description="Custom start date (YYYY-MM-DD)"),
    end_date: str = Query(None, description="Custom end date (YYYY-MM-DD)")
):
    """
    Returns a subset of KPIs optimized for the UI dashboard with dynamic filtering.
    """
    logger.info(f"Received request for metrics: tenant {tenant_id}, range {time_range}, custom {start_date}-{end_date}")
    try:
        metrics = await get_general_kpis(tenant_id, time_range, start_date, end_date)
        if metrics is None:
            # If dbt hasn't run yet or no data exists
            return {
                "status": "success",
                "tenant_id": tenant_id,
                "data": {
                    "top_kpis": {},
                    "charts": {
                        "total_users": {"labels": [], "datasets": []},
                        "conversation_channels": {"labels": [], "datasets": []},
                        "response_time": {"labels": [], "datasets": []}
                    }
                },
                "message": "No metrics available for the selected time range."
            }
            
        return {
            "status": "success",
            "tenant_id": tenant_id,
            "data": metrics
        }
    except HTTPException as he:
        raise he
    except Exception as e:
        logger.error(f"Error fetching general metrics for tenant {tenant_id}: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail={"status": "error", "message": "Internal error retrieving metrics."})

@router.post("/chat/takeover")
async def chat_takeover(
    tenant_id: str = Form(...),
    thread_id: str = Form(...),
    status: bool = Form(...),
    user_id: str = Form(None)
):
    """
    Toggles human takeover for a specific conversation thread.
    """
    from app.services.integrations.firestore import set_takeover_status
    logger.info(f"Takeover toggle: {tenant_id}/{user_id}/{thread_id} -> {status}")
    try:
        success = set_takeover_status(tenant_id, thread_id, status, user_id=user_id)
        if not success:
            raise HTTPException(status_code=404, detail="Conversation thread not found.")
        
        # Real-time synchronization: Broadcast the state change to any active WebSocket listeners
        try:
            from app.services.infra.redis import publish_event
            event_type = "takeover_started" if status else "takeover_ended"
            msg = "Human agent has joined the chat" if status else "AI assistant has resumed control"
            
            # 1. Notify the specific thread participants
            thread_payload = {
                "type": "system",
                "event": "takeover_started" if status else "takeover_ended",
                "message": msg,
                "data": {"takeover_active": status, "user_id": user_id}
            }
            await publish_event(f"tenant_{tenant_id}:chat_{thread_id}", thread_payload)
            
            # 2. Notify the conversation_list observers
            list_payload = {
                "type": "system",
                "event": "takeover_started" if status else "takeover_ended",
                "message": f"Takeover {'started' if status else 'ended'} for {thread_id}",
                "data": {"thread_id": thread_id, "takeover_active": status}
            }
            await publish_event(f"tenant_{tenant_id}:conversation_list", list_payload)
        except Exception as ws_err:
            logger.warning(f"Could not publish takeover status to redis: {ws_err}")

        return {
            "status": "success",
            "message": f"Human takeover {'activated' if status else 'deactivated'} for {thread_id}",
            "takeover_active": status
        }
    except HTTPException:
        raise
    except Exception as e:
        logger.error(f"Takeover toggle failed: {e}", exc_info=True)
        raise HTTPException(status_code=500, detail="Failed to toggle human takeover status.")

CRAWL_MIN_INTERVAL_HOURS = 24


def _crawled_recently(source: dict, tenant_id: str) -> bool:
    """Whether this website source was read recently enough to skip.

    Only websites are throttled, and only when the config does not override the
    window -- a merchant who wants their site read on every press sets
    min_interval_hours to 0 without a deploy.

    Also never skips a source with zero product rows, regardless of
    last_synced_at -- that timestamp only means the source was crawled
    recently, not that the products from that crawl still exist. A cleanup
    that deletes product rows without also clearing this timestamp would
    otherwise leave a build that looks successful but serves an empty
    catalogue.
    """
    if source["kind"] != KIND_CRAWL:
        return False

    last_synced = source.get("last_synced_at")
    if not last_synced:
        return False

    hours = source["config"].get("min_interval_hours")
    hours = CRAWL_MIN_INTERVAL_HOURS if hours is None else float(hours)
    if hours <= 0:
        return False

    age = datetime.now(timezone.utc) - last_synced
    if age >= timedelta(hours=hours):
        return False

    return has_any_products(tenant_id, KIND_CRAWL, source["external_ref"])


def _sync_sources(sources: list, force: bool = False):
    """A blocking, job-shaped wrapper around sync_source for all of a tenant's
    connected sources. One source's failure (e.g. a revoked token) must not
    hide the results of the others, so each is caught and reported inline
    rather than aborting the run.
    """
    def run(tenant_id: str, progress=None) -> dict:
        results = []
        total = len(sources)
        for i, source in enumerate(sources, start=1):
            if not force and _crawled_recently(source, tenant_id):
                # Re-reading a whole website is dozens of requests to the
                # merchant's own server, unlike one API call. Every other source
                # re-syncs on every press; this one earns its keep by usually
                # being a no-op, and force=true still forces it.
                results.append({"source": source["external_ref"],
                                "skipped": "crawled recently"})
                if progress is not None:
                    progress(i * 100 // total, f"synced {i} of {total} sources")
                continue
            try:
                # sync_source does blocking psycopg2 work (migrate, persist, delete)
                # across the whole catalog. Same precedent as recommendations_event_endpoint
                # above: push it to a worker thread so a full sync does not stall the
                # event loop for every other request. sync_source is itself async, so
                # anyio.run bridges it back into its own loop inside that thread.
                outcome = anyio.run(sync_source, tenant_id, source)
                results.append({"source": source["external_ref"], **outcome})
            except Exception as ex:
                # Only the exception's class name goes to the caller --
                # catalog_http.py already interpolates response-derived data
                # into its exception messages, so a raw str(ex) here could
                # leak that back out. Full detail still goes to the logs.
                logger.error("Catalog sync failed for %s", source["external_ref"],
                             exc_info=True)
                results.append({"source": source["external_ref"], "error": type(ex).__name__})
            if progress is not None:
                progress(i * 100 // total, f"synced {i} of {total} sources")
        return {"results": results}

    return run


@router.post("/catalog/sync")
async def catalog_sync_endpoint(
    tenant_id: str = Form(...),
    source_ref: Optional[str] = Form(None),
    wait: bool = Form(False),
    response: Response = None,
):
    sources = get_sources(tenant_id)
    if source_ref:
        sources = [s for s in sources if s["external_ref"] == source_ref]
    if not sources:
        raise HTTPException(status_code=404, detail="No connected source found")

    run = _sync_sources(sources)

    if wait:
        # The old contract, kept so existing callers add one parameter rather
        # than adopting polling.
        return await anyio.to_thread.run_sync(lambda: run(tenant_id))

    try:
        job_id = await start_job(tenant_id, "sync", run)
    except jobs.JobAlreadyRunning as running:
        raise HTTPException(status_code=409, detail={
            "error": "already_running", "job_id": running.job_id})

    response.status_code = 202
    return {"job_id": job_id, "status": "queued"}

@router.post("/catalog/enrich")
async def catalog_enrich_endpoint(
    tenant_id: str = Form(...),
    force: bool = Form(False),
    wait: bool = Form(False),
    response: Response = None,
):
    """Run attribute extraction over a tenant's catalog.

    Separate from /catalog/sync because it costs money and must be re-runnable
    without re-fetching a catalog.
    """
    if wait:
        # The old contract, kept so existing callers add one parameter rather
        # than adopting polling.
        try:
            return await anyio.to_thread.run_sync(lambda: enrich_tenant(tenant_id, force))
        except Exception as ex:
            logger.error("Enrichment failed for %s", tenant_id, exc_info=True)
            raise HTTPException(status_code=502, detail=type(ex).__name__)

    try:
        job_id = await start_job(
            tenant_id, "enrich",
            lambda t, progress=None: enrich_tenant(t, force, progress=progress))
    except jobs.JobAlreadyRunning as running:
        raise HTTPException(status_code=409, detail={
            "error": "already_running", "job_id": running.job_id})

    response.status_code = 202
    return {"job_id": job_id, "status": "queued"}
