"""The merchant-facing surface for the product pairing graph."""
import asyncio
import logging
import pathlib

import anyio
from fastapi import APIRouter, File, Form, HTTPException, Response, UploadFile
from fastapi.responses import FileResponse, HTMLResponse
from pydantic import BaseModel
from psycopg2 import OperationalError as DbUnavailable

from app.services.catalog.build import build_catalog
from app.services.catalog.csv_import import import_csv
from app.services.catalog.csv_source import CsvFormatError
from app.services.infra import jobs
from app.services.infra.quotas import QuotaExceededError
from app.services.pairing.job import pair_tenant
from app.services.catalog.products import delete_all_products, delete_product
from app.services.catalog.sources import disconnect_all_sources
from app.services.pairing.decisions import record
from app.services.pairing.queries import (
    get_product, list_anchors, list_categories, list_pairings, list_products,
    pairings_for,
)

logger = logging.getLogger(__name__)

router = APIRouter(prefix="/catalog", tags=["catalog"])

# Read per request rather than at import: editing the page and refreshing is the
# whole development loop for a tool with no build step.
ADMIN_PAGE = pathlib.Path(__file__).resolve().parent.parent / "static" / "catalog_admin.html"
SAMPLE_CSV = pathlib.Path(__file__).resolve().parent.parent / "static" / "sample_products.csv"
DEMO_PAGE = pathlib.Path(__file__).resolve().parent.parent / "static" / "chat_ui.html"


class DecisionIn(BaseModel):
    anchor_key: str
    neighbor_key: str
    pair_type: str
    decision: str
    decided_by: str = None


class DecideIn(BaseModel):
    tenant_id: str
    decisions: list[DecisionIn]


async def _run(job_id: str, tenant_id: str, kind: str, fn) -> None:
    """Runs a blocking job off the event loop, reporting progress as it goes."""
    def progress(percent, step):
        # Both the job and the progress write are synchronous and share the
        # worker thread, so the thread records its own progress. This used to
        # hand a coroutine back to the event loop, which was only ever needed
        # because the store was an async Redis client.
        jobs.tick(job_id, percent, step)

    try:
        result = await anyio.to_thread.run_sync(
            lambda: fn(tenant_id, progress=progress))
        await jobs.finish(job_id, result=result)
    except QuotaExceededError as ex:
        # A plan limit, not a crash -- store the real message rather than the
        # exception's class name, so a merchant polling this job can see why
        # it stopped and what to do about it.
        logger.info("%s job stopped by quota for %s: %s", kind, tenant_id, ex.message)
        await jobs.finish(job_id, error=ex.message)
    except Exception as ex:
        logger.error("%s job failed for %s", kind, tenant_id, exc_info=True)
        await jobs.finish(job_id, error=type(ex).__name__)


async def start_job(tenant_id: str, kind: str, fn) -> str:
    """Takes the lock and schedules the work, or releases the lock trying.

    The window between create() and create_task() is small but real: an
    exception in it once left a tenant locked out of enrichment with no job
    running and nothing to release it until the TTL expired.
    """
    try:
        job_id = await jobs.create(tenant_id, kind)
    except jobs.JobAlreadyRunning:
        raise
    except DbUnavailable as ex:
        # A job whose progress could never be recorded is worse than no job, so
        # this refuses to start rather than running blind. 503 says "try
        # again", which is true; a 500 would say "this endpoint is broken",
        # which it is not, and sends people looking in the wrong place.
        logger.error("Job store unavailable starting %s job for %s", kind, tenant_id,
                     exc_info=True)
        raise HTTPException(status_code=503, detail={
            "error": "job_store_unavailable",
            "message": "Could not reach the job store. Try again in a moment.",
            "cause": type(ex).__name__})

    try:
        asyncio.create_task(_run(job_id, tenant_id, kind, fn))
    except Exception:
        logger.error("Could not schedule %s job for %s", kind, tenant_id,
                     exc_info=True)
        await jobs.finish(job_id, error="SchedulingFailed")
        raise
    return job_id


@router.post("/build")
async def build_endpoint(tenant_id: str = Form(...), force: bool = Form(False),
                         wait: bool = Form(False), response: Response = None):
    """Sync, enrich and pair as one job -- the whole pipeline behind one button.

    The lock covers the entire build rather than each stage, so a second press
    cannot start pairing while the first is still enriching.
    """
    def run(tenant, progress=None):
        return build_catalog(tenant, progress=progress, force=force)

    if wait:
        try:
            return await anyio.to_thread.run_sync(lambda: run(tenant_id))
        except QuotaExceededError as ex:
            raise HTTPException(status_code=402, detail=ex.message)
        except Exception as ex:
            logger.error("Build 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, "build", 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("/pair")
async def pair_endpoint(tenant_id: str = Form(...), wait: bool = Form(False),
                        response: Response = None):
    if wait:
        # The old contract, kept so existing callers add one parameter rather
        # than adopting polling. The admin page uses it: its jobs are small.
        try:
            return await anyio.to_thread.run_sync(lambda: pair_tenant(tenant_id))
        except Exception as ex:
            logger.error("Pairing 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, "pair", pair_tenant)
    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"}


# Registered before /products/{product_key} so "jobs" is never captured as a key.
@router.get("/jobs/{job_id}")
async def job_status_endpoint(job_id: str):
    try:
        state = await jobs.get(job_id)
    except DbUnavailable as ex:
        # A client polls this every couple of seconds, so a blip here is the
        # most likely error anyone sees. 503 tells them to keep polling; a 500
        # would look like the job itself had failed.
        logger.error("Job store unavailable reading job %s", job_id, exc_info=True)
        raise HTTPException(status_code=503, detail={
            "error": "job_store_unavailable",
            "message": "Could not reach the job store. Keep polling.",
            "cause": type(ex).__name__})

    if state is None:
        raise HTTPException(status_code=404, detail="Job not found")
    return state


# Registered before /products/{product_key} so "jobs" is never captured as a key.
@router.get("/jobs")
async def recent_jobs_endpoint(tenant_id: str, limit: int = 20):
    return await jobs.recent(tenant_id, limit)


@router.get("/categories")
async def categories_endpoint(tenant_id: str):
    return await anyio.to_thread.run_sync(lambda: list_categories(tenant_id))


@router.get("/products")
async def products_endpoint(tenant_id: str, category: str = None, limit: int = 50,
                             offset: int = 0):
    return await anyio.to_thread.run_sync(
        lambda: list_products(tenant_id, category=category, limit=limit, offset=offset))


@router.delete("/products")
async def delete_all_products_endpoint(tenant_id: str):
    """Wipes every product, pairing, and embedding for a tenant, and
    disconnects every source it had connected -- a source left connected
    would just get silently re-crawled on the next build, repopulating
    exactly what this just deleted. Reconnecting a source (POST
    /sources/website, etc.) is what makes the next build add anything again."""
    def _wipe():
        counts = delete_all_products(tenant_id)
        counts["sources_disconnected"] = disconnect_all_sources(tenant_id)
        return counts
    return await anyio.to_thread.run_sync(_wipe)


# Registered before /products/{product_key} so "pairings" is never captured.
@router.get("/pairings/anchors")
async def anchors_endpoint(tenant_id: str, status: str = "pending",
                           category: str = None, score: str = None,
                           limit: int = 50):
    """Products with matches, and how many of each state they hold.

    Backs a review screen that shows one product at a time -- a merchant thinks
    in products, so a flat list of every pair across the catalogue is not
    something they can work through.

    status   pending | approved | rejected | auto | all
    score    high | medium | low
    category a product category to narrow to

    Every anchor carries the full breakdown regardless of the filter, so a
    screen can show "4 pending, 2 approved" without asking again per product.
    """
    try:
        return await anyio.to_thread.run_sync(
            lambda: list_anchors(tenant_id, limit=limit, status=status,
                                 category=category, score=score))
    except ValueError as ex:
        raise HTTPException(status_code=422, detail=str(ex))


# Registered before /products/{product_key} so "pairings" is never captured.
@router.get("/pairings")
async def pairings_endpoint(tenant_id: str, anchor_key: str = None,
                            status: str = "pending", category: str = None,
                            score: str = None, limit: int = 100):
    """The matches themselves, filtered the same way as the anchor list.

    Pass anchor_key for one product's review screen; leave it off for the whole
    catalogue.
    """
    try:
        return await anyio.to_thread.run_sync(
            lambda: list_pairings(tenant_id, limit=limit, anchor_key=anchor_key,
                                  status=status, category=category, score=score))
    except ValueError as ex:
        raise HTTPException(status_code=422, detail=str(ex))


# Registered before /products/{product_key} so "admin" is never captured as a key.
@router.get("/admin", response_class=HTMLResponse)
async def admin_page():
    try:
        return HTMLResponse(ADMIN_PAGE.read_text())
    except OSError:
        logger.error("Could not read the admin page at %s", ADMIN_PAGE,
                     exc_info=True)
        raise HTTPException(status_code=500, detail="admin page unavailable")


# The chat and Connect-Shopify demo. Served from the API so it is same-origin:
# opened as a file:// page it would send Origin: null on every call, and this
# service sets allow_credentials, which makes a wildcard origin unreliable.
@router.get("/demo", response_class=HTMLResponse)
async def demo_page():
    try:
        return HTMLResponse(DEMO_PAGE.read_text())
    except OSError:
        logger.error("Could not read the demo page at %s", DEMO_PAGE,
                     exc_info=True)
        raise HTTPException(status_code=500, detail="demo page unavailable")


@router.post("/pairings/decide")
async def decide_endpoint(payload: DecideIn):
    decisions = [d.model_dump() for d in payload.decisions]
    try:
        recorded = await anyio.to_thread.run_sync(
            lambda: record(payload.tenant_id, decisions))
    except ValueError as ex:
        raise HTTPException(status_code=422, detail=str(ex))
    return {"recorded": recorded}


# Registered before /products/{product_key} so "sample.csv" is never captured
# as a key.
@router.get("/sample.csv")
async def sample_csv_endpoint():
    return FileResponse(SAMPLE_CSV, media_type="text/csv",
                        filename="sample_products.csv")


# Registered before /products/{product_key} so "import" is never captured as a key.
@router.post("/import/csv")
async def import_csv_endpoint(tenant_id: str = Form(...), currency: str = Form(None),
                              dry_run: bool = Form(False),
                              file: UploadFile = File(...)):
    body = await file.read()
    text = body.decode("utf-8", errors="replace")
    try:
        return await anyio.to_thread.run_sync(
            lambda: import_csv(tenant_id, text, file.filename, currency, dry_run=dry_run))
    except CsvFormatError as ex:
        raise HTTPException(status_code=422, detail=str(ex))
    except QuotaExceededError as ex:
        # A plan limit, not a server fault -- the merchant can act on this
        # (remove products, upgrade), so the real message must reach them
        # rather than being reduced to the exception's class name.
        raise HTTPException(status_code=402, detail=ex.message)
    except Exception as ex:
        logger.error("CSV import failed for %s", tenant_id, exc_info=True)
        raise HTTPException(status_code=502, detail=type(ex).__name__)


# Registered before /products/{product_key:path} -- a greedy path converter
# would otherwise match this URL too, treating "pairings" as part of the key.
@router.get("/products/{product_key:path}/pairings")
async def product_pairings_endpoint(product_key: str, tenant_id: str):
    return await anyio.to_thread.run_sync(lambda: pairings_for(tenant_id, product_key))


@router.get("/products/{product_key:path}")
async def product_endpoint(product_key: str, tenant_id: str):
    product = await anyio.to_thread.run_sync(lambda: get_product(tenant_id, product_key))
    if product is None:
        raise HTTPException(status_code=404, detail="Product not found")
    return product


@router.delete("/products/{product_key:path}")
async def delete_product_endpoint(product_key: str, tenant_id: str):
    removed = await anyio.to_thread.run_sync(lambda: delete_product(tenant_id, product_key))
    if not removed:
        raise HTTPException(status_code=404, detail="Product not found")
    return {"product_key": product_key, "status": "deleted"}
