"""The collector and catalogue review area: /collectors, and its merge session.

Sources of truth: app/services/collector_view.py (the reads), app/services/merge_session.py
(the queue and the decisions). Owner-only (the same `sources.manage` permission as the kill
switch), computed on every request, so nothing here can show a number the database no
longer agrees with. The reads write nothing; the three merge writes record who decided
from the session, never from the payload.
"""

import os
import signal
from datetime import UTC, date, datetime
from typing import Literal

from fastapi import APIRouter, Depends, HTTPException, Query, Request
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from sqlalchemy import select
from sqlalchemy.orm import Session

from app.db import get_db
from app.models import Source
from app.models.schemas import ControlOut, LiveOut
from app.services import audit_log, identity, procinfo
from app.services import collector_view as view
from app.services import listings_table, merge_desk, merge_session, overrides
from app.services.collectors import control
from app.services.collectors.registry import COLLECTORS
from app.services.ingest import is_stuck, mark_stuck

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


@router.get("")
def read_collectors(db: Session = Depends(get_db)) -> dict:
    return view.collectors(db)


@router.get("/catalogue")
def read_catalogue(db: Session = Depends(get_db)) -> dict:
    return view.catalogue(db)


@router.get("/runs")
def read_runs(
    collector: str | None = None,
    limit: int = Query(default=60, ge=1, le=400),
    db: Session = Depends(get_db),
) -> dict:
    return view.runs(db, collector=collector, limit=limit)


@router.get("/live", response_model=LiveOut)
def read_live(db: Session = Depends(get_db)) -> LiveOut:
    """What is collecting right now and how it is going; the page polls it. A read: the nine
    state words, the counters, percent and ETA, the pace with its floor, memory, opening hours."""
    return view.live(db)


@router.get("/problems")
def read_problems(db: Session = Depends(get_db)) -> dict:
    return view.problems(db)


@router.get("/products")
def read_product_variants(
    q: str | None = None,
    collector: str | None = None,
    airport: str | None = None,
    category: str | None = None,
    vertical: str | None = None,
    brand: str | None = None,
    quantity_ml: float | None = Query(default=None, ge=0),
    missing: str | None = Query(default=None, pattern="^(image|size|gtin|category|brand)$"),
    min_airports: int | None = Query(default=None, ge=1, le=25),
    group: str | None = Query(default=None, pattern="^(brand|category|vertical|size)$"),
    sort: str = Query(default="name", pattern="^(name|brand|category|size|airports|newest)$"),
    descending: bool = False,
    page: int = Query(default=1, ge=1),
    per_page: int = Query(default=50, ge=1, le=200),
    db: Session = Depends(get_db),
) -> dict:
    filters = dict(q=q, collector=collector, airport=airport, category=category,
                   vertical=vertical, missing=missing, min_airports=min_airports,
                   brand=brand, quantity_ml=quantity_ml)
    if group:
        return view.grouped(db, group=group, **filters)
    return view.product_variants(db, sort=sort, descending=descending, page=page, per_page=per_page, **filters)


@router.get("/listings", response_model=None)
def read_listings(
    q: str | None = None,
    collector: str | None = None,
    airport: str | None = None,
    brand: str | None = None,
    listed_brand: str | None = None,
    quantity_unit: str | None = Query(default=None, pattern="^(ml|g|pcs)$"),
    quantity_state: str | None = Query(default=None, pattern="^(stated|none|unparsed)$"),
    form: str | None = Query(default=None, pattern="^(single|pack|set|refill)$"),
    attribute_kind: str | None = Query(default=None, pattern="^(concentration|color|flavor|age|cask|edition)$"),
    has_fragment: bool | None = None,
    pinned: bool | None = None,
    differs: str | None = Query(default=None, pattern="^(any|brand|name|quantity|size)$"),
    ignored: str | None = Query(default=None, pattern="^(only|any)$"),
    sort: str | None = None,
    dir: str = Query(default="asc", pattern="^(asc|desc)$"),
    columns: str | None = None,
    format: str | None = Query(default=None, pattern="^(csv)$"),
    page: int = Query(default=1, ge=1),
    per_page: int = Query(default=50, ge=1, le=listings_table.PER_PAGE_MAX),
    db: Session = Depends(get_db),
):
    """The giant Listings table (Stream L): every column sortable and filterable in SQL
    (`services/listings_table.py`); `format=csv` streams every matching row with the chosen
    columns, on this same path so the session cookie rides on the link."""
    filters = dict(q=q, brand=brand, collector=collector, airport=airport, listed_brand=listed_brand,
                   quantity_unit=quantity_unit, quantity_state=quantity_state, form=form,
                   attribute_kind=attribute_kind, has_fragment=has_fragment, ignored=ignored, pinned=pinned,
                   differs=differs)
    if format == "csv":
        stamp = date.today().isoformat()
        return StreamingResponse(
            listings_table.csv_rows(db, sort=sort, direction=dir, columns=columns, **filters),
            media_type="text/csv; charset=utf-8",
            headers={"Content-Disposition": f'attachment; filename="listings-{stamp}.csv"'},
        )
    return listings_table.query(db, sort=sort, direction=dir, page=page, per_page=per_page, columns=columns, **filters)


# --------------------------------------------------------------------------- the decided layer (Stream L)

class OverrideIn(BaseModel):
    """A product decision: which field (name, product_line_id, attribute, quantity), the value, why."""

    field: str
    value: object = None
    reason: str | None = None


class PinIn(BaseModel):
    variant_id: int
    reason: str | None = None


class ReasonIn(BaseModel):
    reason: str | None = None


def _decision_refused(exc: overrides.Refused) -> HTTPException:
    status = 404 if exc.code.endswith("_NOT_FOUND") else 422 if exc.code == "VALUE_INVALID" else 409
    return HTTPException(status_code=status, detail={"error_code": exc.code, "summary": exc.summary})


@router.post("/products/{variant_id}/override")
def override_product(variant_id: int, payload: OverrideIn, request: Request, db: Session = Depends(get_db)) -> dict:
    """Decide a product's name, line, attribute or quantity; the key follows and the
    decision survives every rekey. Who decided comes from the session, never the payload."""
    who = identity.actor(request)
    try:
        result = overrides.decide_product(db, variant_id, payload.field, payload.value, set_by=who.id, reason=payload.reason)
    except overrides.Refused as exc:
        raise _decision_refused(exc) from exc
    audit_log.record("product.override", entity_type="product", entity_key=str(variant_id),
                     detail={"field": payload.field, "value": result["value"], "reason": payload.reason})
    return result


@router.post("/listings/{listing_id}/pin")
def pin_listing(listing_id: int, payload: PinIn, request: Request, db: Session = Depends(get_db)) -> dict:
    who = identity.actor(request)
    try:
        result = overrides.pin_listing(db, listing_id, payload.variant_id, set_by=who.id, reason=payload.reason)
    except overrides.Refused as exc:
        raise _decision_refused(exc) from exc
    audit_log.record("listing.pin", entity_type="listing", entity_key=str(listing_id),
                     detail={"variant_id": payload.variant_id, "reason": payload.reason})
    return result


@router.post("/listings/{listing_id}/unpin")
def unpin_listing(listing_id: int, request: Request, db: Session = Depends(get_db)) -> dict:
    who = identity.actor(request)
    try:
        result = overrides.unpin_listing(db, listing_id, set_by=who.id)
    except overrides.Refused as exc:
        raise _decision_refused(exc) from exc
    audit_log.record("listing.unpin", entity_type="listing", entity_key=str(listing_id))
    return result


@router.post("/listings/{listing_id}/ignore")
def ignore_listing(listing_id: int, payload: ReasonIn, request: Request, db: Session = Depends(get_db)) -> dict:
    """Keep collecting the listing's prices, keep it out of the site's price readers and the
    collectors page. Tonight ignore reaches `catalog_queries.py` and `/collectors` only: the
    trip comparison, the home counts, IndexNow, coverage and the audit still see it."""
    who = identity.actor(request)
    try:
        result = overrides.ignore_listing(db, listing_id, set_by=who.id, reason=payload.reason)
    except overrides.Refused as exc:
        raise _decision_refused(exc) from exc
    audit_log.record("listing.ignore", entity_type="listing", entity_key=str(listing_id), detail={"reason": payload.reason})
    return result


@router.post("/listings/{listing_id}/unignore")
def unignore_listing(listing_id: int, request: Request, db: Session = Depends(get_db)) -> dict:
    who = identity.actor(request)
    try:
        result = overrides.unignore_listing(db, listing_id, set_by=who.id)
    except overrides.Refused as exc:
        raise _decision_refused(exc) from exc
    audit_log.record("listing.unignore", entity_type="listing", entity_key=str(listing_id))
    return result


# --------------------------------------------------------------------------- the merge session

class MergeConfirmIn(BaseModel):
    """The preferred name (one of the pair's spellings, or any), which side to keep when
    the name does not settle it, and a note (the barcode override records why)."""

    preferred_name: str | None = None
    keep: str | None = None
    note: str | None = None


class MergeKeepSeparateIn(BaseModel):
    note: str | None = None


def _refused(exc: merge_session.Refused) -> HTTPException:
    status = 404 if exc.code == "PAIR_NOT_FOUND" else 422 if exc.code == "NOTE_REQUIRED" else 409
    return HTTPException(status_code=status, detail={"error_code": exc.code, "summary": exc.summary})


@router.get("/merge")
def read_merge_queue(
    level: str = Query(default="brand", pattern="^(brand|line|product)$"),
    page: int = Query(default=1, ge=1),
    per_page: int = Query(default=20, ge=1, le=100),
    db: Session = Depends(get_db),
) -> dict:
    return merge_session.queue(db, level=level, page=page, per_page=per_page)


@router.post("/merge/{suggestion_id}/confirm")
def confirm_merge(suggestion_id: int, payload: MergeConfirmIn, request: Request, db: Session = Depends(get_db)) -> dict:
    who = identity.actor(request)
    try:
        result = merge_session.confirm(
            db, suggestion_id, decided_by=who.id, preferred_name=payload.preferred_name,
            keep=payload.keep if payload.keep in ("left", "right") else None, note=payload.note,
        )
    except merge_session.Refused as exc:
        raise _refused(exc) from exc
    audit_log.record("merge.confirm", entity_type=result["level"], entity_key=str(suggestion_id),
                     detail={"applied": result["applied"], "preferred_name": payload.preferred_name})
    result["progress"] = merge_session.progress(db)
    return result


@router.post("/merge/{suggestion_id}/keep-separate")
def keep_separate(suggestion_id: int, payload: MergeKeepSeparateIn, request: Request, db: Session = Depends(get_db)) -> dict:
    who = identity.actor(request)
    try:
        result = merge_session.keep_separate(db, suggestion_id, decided_by=who.id, note=payload.note)
    except merge_session.Refused as exc:
        raise _refused(exc) from exc
    audit_log.record("merge.keep_separate", entity_type=result["level"], entity_key=str(suggestion_id),
                     detail={"note": payload.note})
    result["progress"] = merge_session.progress(db)
    return result


# --------------------------------------------------------------------------- the merge desk (Stream L)

class BatchConfirmIn(BaseModel):
    suggestion_id: int
    preferred_name: str | None = None
    keep: str | None = None
    note: str | None = None


class BatchKeepSeparateIn(BaseModel):
    suggestion_id: int
    note: str | None = None


class BatchIn(BaseModel):
    """Up to 200 decisions, applied in id order as the signed-in person's; a refused pair is
    reported and the rest proceed; suggestions regenerate once at the end."""

    confirm: list[BatchConfirmIn] = []
    keep_separate: list[BatchKeepSeparateIn] = []


class ProposalEntry(BaseModel):
    level: str
    left_id: int
    right_id: int
    score: float = 0.5
    why: str


class ProposeIn(BaseModel):
    proposed_by: str | None = None
    entries: list[ProposalEntry]


@router.get("/desk/brands")
def read_desk_brands(
    text: str | None = None,
    has_suggestion: bool | None = None,
    min_score: float | None = Query(default=None, ge=0, le=1),
    sort: str = Query(default="score", pattern="^(name|score|product_variants|listings|lines)$"),
    dir: str = Query(default="desc", pattern="^(asc|desc)$"),
    page: int = Query(default=1, ge=1),
    per_page: int = Query(default=100, ge=1, le=500),
    db: Session = Depends(get_db),
) -> dict:
    return merge_desk.brands(db, text=text, has_suggestion=has_suggestion, min_score=min_score, sort=sort, direction=dir,
                             page=page, per_page=per_page)


@router.get("/desk/lines")
def read_desk_lines(
    brand: str | None = None,
    text: str | None = None,
    vertical: str | None = None,
    kind: str | None = None,
    reason: str | None = None,
    min_score: float | None = Query(default=None, ge=0, le=1),
    sort: str = Query(default="score", pattern="^(name|score|product_variants|listings|lines|brand)$"),
    dir: str = Query(default="desc", pattern="^(asc|desc)$"),
    page: int = Query(default=1, ge=1),
    per_page: int = Query(default=100, ge=1, le=500),
    db: Session = Depends(get_db),
) -> dict:
    return merge_desk.lines(db, brand=brand, text=text, vertical=vertical, kind=kind, reason=reason, min_score=min_score,
                            sort=sort, direction=dir, page=page, per_page=per_page)


@router.post("/merge/batch")
def batch_merge(payload: BatchIn, request: Request, db: Session = Depends(get_db)) -> dict:
    who = identity.actor(request)
    try:
        return merge_desk.batch(
            db, decided_by=who.id,
            confirm=[c.model_dump() for c in payload.confirm], keep_separate=[r.model_dump() for r in payload.keep_separate],
        )
    except merge_session.Refused as exc:
        raise _refused(exc) from exc


@router.post("/merge/propose")
def propose_merges(payload: ProposeIn, request: Request, db: Session = Depends(get_db)) -> dict:
    """File proposals (reason `proposed`, never withdrawn by the rules); nothing is merged."""
    who = identity.actor(request)
    result = merge_desk.propose(db, [e.model_dump() for e in payload.entries], proposed_by=payload.proposed_by or who.username)
    db.commit()
    audit_log.record("merge.propose", entity_type="merge", entity_key=payload.proposed_by or who.username,
                     detail={k: v for k, v in result.items() if k != "notes"})
    return result


@router.post("/merge/suggest")
def refresh_suggestions(request: Request, db: Session = Depends(get_db)) -> dict:
    identity.actor(request)
    counts = merge_session.refresh(db)
    return {"suggested": counts, "progress": merge_session.progress(db)}


# --------------------------------------------------------------------------- the control plane (Stream AW4)
# Six writes under `sources.manage`, each: the actor first, the source row locked for the
# request (FOR NO KEY UPDATE: a running collector's FK holds only KEY SHARE), the refusals in
# the brief's order, one audit row, the collector as the live read now sees it back. A refusal is 409
# `{error_code, summary}`; the summary is the one sentence the page can show, because the SPA's
# `ApiError` carries only the code (an issue on the running list).

class StartIn(BaseModel):
    mode: Literal["discover", "recheck"] | None = None
    limit: int | None = Field(default=None, ge=1, le=100000)


class StopIn(BaseModel):
    force: bool = False


class PaceIn(BaseModel):
    delay_seconds: float = Field(gt=0)


class ModeIn(BaseModel):
    mode: Literal["discover", "recheck"]


def _control_refused(refusal: control.Refusal) -> HTTPException:
    return HTTPException(status_code=409, detail={"error_code": refusal.code, "summary": refusal.summary})


def _locked_source(db: Session, slug: str, *, create: bool = False) -> tuple[Source, object]:
    collector = COLLECTORS.get(slug)
    if collector is None:
        raise HTTPException(status_code=404, detail={"error_code": "SOURCE_NOT_FOUND", "summary": f"no collector {slug}"})
    source = db.execute(select(Source).where(Source.slug == slug).with_for_update(key_share=True)).scalar_one_or_none()
    if source is None:
        if not create:
            raise HTTPException(status_code=404, detail={"error_code": "SOURCE_NOT_FOUND", "summary": f"{slug} has never run"})
        source = Source(slug=slug, name=collector.retailer_name)
        db.add(source)
        db.flush()
    return source, collector


def _live_run(db: Session, source: Source, now: datetime):
    """The running row that is really live, or None (a dead row is not a run to pause)."""
    run = control.running_run(db, source.id)
    if run is None or is_stuck(run, now):
        return None
    return run


@router.post("/{slug}/start", status_code=202)
def start_collector(slug: str, payload: StartIn, request: Request, db: Session = Depends(get_db)) -> ControlOut:
    who = identity.actor(request)
    now = datetime.now(UTC)
    source, collector = _locked_source(db, slug, create=True)
    mode = payload.mode or source.mode
    refusal = control.start_refusal(db, source, collector, mode, payload.limit, procinfo.container_memory(), now)
    if refusal is not None:
        db.rollback()
        raise _control_refused(refusal)
    source.control, source.control_set_by, source.control_set_at, source.mode = "run", who.username, now, mode
    db.commit()
    proc = control.spawn_collector(slug, mode, payload.limit, who.username)
    audit_log.record("source.start", entity_type="source", entity_key=slug,
                     detail={"mode": mode, "limit": payload.limit, "pid": proc.pid})
    return view.control_out(db, slug, now, started=True, pid=proc.pid)


@router.post("/{slug}/pause")
def pause_collector(slug: str, request: Request, db: Session = Depends(get_db)) -> ControlOut:
    who = identity.actor(request)
    now = datetime.now(UTC)
    source, collector = _locked_source(db, slug)
    run = _live_run(db, source, now)
    if run is None:
        db.rollback()
        raise _control_refused(control.Refusal("NOT_RUNNING", f"{slug} is not collecting."))
    effective = control.effective_control(source.control, source.control_set_at, run.started_at)
    if not source.enabled or effective == "stop":
        db.rollback()
        raise _control_refused(control.Refusal("ALREADY_STOPPING", f"{slug} is already stopping."))
    if effective == "pause":
        db.rollback()
        return view.control_out(db, slug, now, changed=False)
    source.control, source.control_set_by, source.control_set_at = "pause", who.username, now
    db.commit()
    audit_log.record("source.pause", entity_type="source", entity_key=slug, detail={"run_id": run.id})
    return view.control_out(db, slug, now)


@router.post("/{slug}/resume")
def resume_collector(slug: str, request: Request, db: Session = Depends(get_db)) -> ControlOut:
    who = identity.actor(request)
    now = datetime.now(UTC)
    source, collector = _locked_source(db, slug)
    run = _live_run(db, source, now)
    effective = control.effective_control(source.control, source.control_set_at, run.started_at) if run else "run"
    if run is not None and (not source.enabled or effective == "stop"):
        db.rollback()
        raise _control_refused(control.Refusal("ALREADY_STOPPING", f"{slug} is already stopping."))
    if run is None or effective != "pause":
        db.rollback()
        raise _control_refused(control.Refusal("NOT_PAUSED", f"{slug} is not paused."))
    source.control, source.control_set_by, source.control_set_at = "run", who.username, now
    db.commit()
    audit_log.record("source.resume", entity_type="source", entity_key=slug, detail={"run_id": run.id})
    return view.control_out(db, slug, now)


@router.post("/{slug}/stop")
def stop_collector(slug: str, payload: StopIn, request: Request, db: Session = Depends(get_db)) -> ControlOut:
    """A cooperative stop: the loop ends at its next request boundary. `force` sends SIGTERM to a
    run that was already asked to stop and did not answer; never SIGKILL. A row the rules call dead
    is marked ended by the actor (Mark as ended on the page), because a dead process holds no lock."""
    who = identity.actor(request)
    now = datetime.now(UTC)
    source, collector = _locked_source(db, slug)
    run = control.running_run(db, source.id)
    if run is None:
        db.rollback()
        raise _control_refused(control.Refusal("NOT_RUNNING", f"{slug} is not collecting."))
    if is_stuck(run, now):
        mark_stuck(run, f"ended by {who.username} from the page")
        db.commit()
        audit_log.record("source.stop", entity_type="source", entity_key=slug,
                         detail={"run_id": run.id, "dead": True, "marked": "stuck"})
        return view.control_out(db, slug, now)
    if payload.force:
        if control.effective_control(source.control, source.control_set_at, run.started_at) != "stop":
            db.rollback()
            raise _control_refused(control.Refusal("STOP_NOT_REQUESTED", "ask it to stop first; Stop now is for a run that did not answer."))
        if run.pid is None:
            db.rollback()
            raise _control_refused(control.Refusal("PID_UNKNOWN", f"run {run.id} recorded no process id."))
        living = procinfo.alive(run.pid, slug)
        if living is None:
            db.rollback()
            raise _control_refused(control.Refusal("PID_UNKNOWN", "the process table is not readable here."))
        if living:
            os.kill(run.pid, signal.SIGTERM)
            db.commit()
            audit_log.record("source.stop", entity_type="source", entity_key=slug,
                             detail={"run_id": run.id, "force": True, "pid": run.pid, "signal": "SIGTERM"})
        else:
            mark_stuck(run, f"stalled, process gone, stopped by {who.username}")
            db.commit()
            audit_log.record("source.stop", entity_type="source", entity_key=slug,
                             detail={"run_id": run.id, "force": True, "pid": run.pid, "marked": "stuck"})
        return view.control_out(db, slug, now)
    source.control, source.control_set_by, source.control_set_at = "stop", who.username, now
    db.commit()
    audit_log.record("source.stop", entity_type="source", entity_key=slug, detail={"run_id": run.id})
    return view.control_out(db, slug, now)


@router.post("/{slug}/pace")
def set_pace(slug: str, payload: PaceIn, request: Request, db: Session = Depends(get_db)) -> ControlOut:
    """The pace per source: never under the host's robots crawl delay as the loop last read it,
    never under our own second, never under the render floor on a rendered source."""
    who = identity.actor(request)
    now = datetime.now(UTC)
    source, collector = _locked_source(db, slug, create=True)
    refusal = control.pace_refusal(source, collector, payload.delay_seconds)
    if refusal is not None:
        db.rollback()
        raise _control_refused(refusal)
    before = float(source.delay_seconds or 0.0)
    source.delay_seconds, source.delay_set_by, source.delay_set_at = payload.delay_seconds, who.username, now
    db.commit()
    audit_log.record("source.pace", entity_type="source", entity_key=slug,
                     detail={"from": before, "to": payload.delay_seconds, "floor": control.pace_of(source, collector).floor})
    return view.control_out(db, slug, now)


@router.post("/{slug}/mode")
def set_mode(slug: str, payload: ModeIn, request: Request, db: Session = Depends(get_db)) -> ControlOut:
    who = identity.actor(request)
    now = datetime.now(UTC)
    source, collector = _locked_source(db, slug, create=True)
    if payload.mode == "recheck":
        refusal = control.recheck_refusal(collector)
        if refusal is not None:
            db.rollback()
            raise _control_refused(refusal)
    before = source.mode
    source.mode = payload.mode
    db.commit()
    audit_log.record("source.mode", entity_type="source", entity_key=slug,
                     detail={"from": before, "to": payload.mode, "by": who.username})
    return view.control_out(db, slug, now)
