#!/usr/bin/env python3
"""First controlled production run: 1 group, real dgred writes.

Scrapes 1 active czesci_tir group via Apify, runs full pipeline,
writes leads to dgred for real. Reports client IDs and links.
"""

from __future__ import annotations

import json
import logging
import sys
import time
from collections import Counter
from datetime import datetime, timedelta
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "src"))

import httpx

from agent_samochodowy.config import Settings
from agent_samochodowy.dgred.client import DgredClient
from agent_samochodowy.dgred.writer import LeadWriter
from agent_samochodowy.extraction import extract
from agent_samochodowy.ingestor.groups import GroupManager
from agent_samochodowy.matcher.engine import match
from agent_samochodowy.matcher.index import CatalogIndex
from agent_samochodowy.models import Post
from agent_samochodowy.status import map_status

logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
logger = logging.getLogger(__name__)
logging.getLogger("httpx").setLevel(logging.WARNING)
logging.getLogger("openai").setLevel(logging.WARNING)
logging.getLogger("httpcore").setLevel(logging.WARNING)

APIFY_BASE = "https://api.apify.com/v2"
STATUS_ID_MAP = {"Dopasowano": 67, "Do weryfikacji": 68, "Brak dopasowania": 69}
DGRED_CLIENT_URL = "https://app.dgred.com/client"

RESULTS_LIMIT = 20


def run_apify(token: str, actor_id: str, group_url: str) -> list[dict]:
    """Run Apify actor, poll for completion, return items."""
    with httpx.Client(timeout=180) as client:
        logger.info("Apify: starting scrape of %s", group_url)
        resp = client.post(
            f"{APIFY_BASE}/acts/{actor_id}/runs",
            params={"token": token},
            json={
                "startUrls": [{"url": group_url}],
                "resultsLimit": RESULTS_LIMIT,
                "sortBy": "New posts",
            },
        )
        resp.raise_for_status()
        run_data = resp.json().get("data", {})
        run_id = run_data.get("id")
        dataset_id = run_data.get("defaultDatasetId")

        if not run_id:
            logger.error("No run ID")
            return []

        for _ in range(90):
            time.sleep(2)
            try:
                sr = client.get(f"{APIFY_BASE}/actor-runs/{run_id}", params={"token": token})
                status = sr.json().get("data", {}).get("status")
                if status in ("SUCCEEDED", "FAILED", "ABORTED", "TIMED-OUT"):
                    break
            except Exception:
                continue

        if status != "SUCCEEDED":
            logger.warning("Apify run ended: %s", status)
            return []

        items_resp = client.get(
            f"{APIFY_BASE}/datasets/{dataset_id}/items",
            params={"token": token, "format": "json"},
        )
        items_resp.raise_for_status()
        return items_resp.json()


def map_item(item: dict, group_name: str) -> Post | None:
    text = item.get("text") or item.get("message") or ""
    if not text.strip():
        return None
    post_id = str(item.get("postId") or item.get("id") or item.get("url", ""))
    if not post_id:
        return None

    timestamp = item.get("timestamp") or item.get("time") or item.get("date")
    if timestamp and isinstance(timestamp, str):
        for fmt in ("%Y-%m-%dT%H:%M:%S.%fZ", "%Y-%m-%dT%H:%M:%SZ",
                    "%Y-%m-%dT%H:%M:%S", "%Y-%m-%d %H:%M:%S"):
            try:
                timestamp = datetime.strptime(timestamp, fmt)
                break
            except ValueError:
                continue
        else:
            timestamp = datetime.now()
    elif not timestamp:
        timestamp = datetime.now()

    author = (item.get("authorName")
              or (item.get("user", {}).get("name") if isinstance(item.get("user"), dict) else None)
              or "Unknown")
    profile_url = item.get("authorProfileUrl") or item.get("profileUrl") or ""
    post_url = item.get("postUrl") or item.get("url") or ""

    try:
        return Post(
            post_id=post_id, group=group_name, author_name=author,
            author_profile_url=profile_url, text=text,
            timestamp=timestamp, post_url=post_url,
        )
    except Exception:
        return None


def main() -> None:
    settings = Settings()
    settings.dry_run = False  # PRODUCTION

    print("=" * 80)
    print("FIRST PRODUCTION RUN — real dgred writes")
    print("=" * 80)
    print(f"LLM:            {settings.llm_provider} / {settings.llm_model}")
    print(f"Post max age:   {settings.post_max_age_hours}h")
    print(f"dry_run:        {settings.dry_run}")
    print()

    # Pick 1 group
    gm = GroupManager(settings.db_path)
    groups = gm.pick_groups(1)
    if not groups:
        print("ERROR: No active czesci_tir groups")
        return
    group = groups[0]
    print(f"Group: {group['nazwa']}")
    print(f"URL:   {group['url']}")
    print()

    # Scrape
    items = run_apify(settings.apify_token, settings.apify_actor_id, group["url"])
    if not items:
        print("No items returned (group private or Apify error)")
        gm.close()
        return

    posts = [p for p in (map_item(i, group["nazwa"]) for i in items) if p]
    total_fetched = len(posts)

    # Freshness filter
    cutoff = datetime.now() - timedelta(hours=settings.post_max_age_hours)
    fresh = [p for p in posts if p.timestamp.replace(tzinfo=None) >= cutoff]
    skipped_old = total_fetched - len(fresh)
    posts = fresh

    print(f"Posts fetched:      {total_fetched}")
    print(f"Too old (>{settings.post_max_age_hours}h):    {skipped_old}")
    print(f"To analyze:         {len(posts)}")
    print()

    # Pipeline with REAL dgred
    index = CatalogIndex(settings.db_path)
    dgred = DgredClient(settings)
    writer = LeadWriter(dgred, settings)

    statuses: Counter[str] = Counter()
    methods: Counter[str] = Counter()
    created_clients: list[dict] = []
    duplicates_found = 0

    for post in posts:
        query = extract(post, settings)
        if query is None or not query.is_part_request:
            continue

        result = match(query, index)
        status_name = map_status(result, settings)

        # Track if this author was already seen (dedup)
        ext_id = post.author_profile_url
        existing = dgred.find_client_by_external_id(ext_id) if ext_id else None
        is_dup = existing is not None

        # Write lead (creates or appends note)
        writer.write_lead(post, query, result)

        statuses[status_name] += 1
        methods[result.method] += 1

        # Get the client ID from writer's cache
        client_id = writer._client_cache.get(ext_id, "?")

        if is_dup:
            duplicates_found += 1

        created_clients.append({
            "client_id": client_id,
            "name": post.author_name,
            "status": status_name,
            "status_id": STATUS_ID_MAP.get(status_name, 0),
            "method": result.method,
            "confidence": result.confidence,
            "n_offers": len(result.offers),
            "post_text": post.text[:80],
            "is_dup": is_dup,
            "brand": query.brand,
            "category": query.category,
            "oe": ", ".join(query.oe_numbers) if query.oe_numbers else "—",
        })

    hits = statuses.get("Dopasowano", 0) + statuses.get("Do weryfikacji", 0)
    gm.record_scan(group["group_id"], queries=len(created_clients), hits=hits)
    index.close()
    dgred.close()
    gm.close()

    # Report
    leads_total = len(created_clients)
    print("=" * 80)
    print("RAPORT PRODUKCYJNY")
    print("=" * 80)
    print(f"Postów pobranych:       {total_fetched}")
    print(f"Za stare:               {skipped_old}")
    print(f"Do analizy:             {len(posts)}")
    print(f"Zapytania o część:      {leads_total}")
    print(f"Duplikaty autorów:      {duplicates_found}")
    print()
    print("Rozkład statusów:")
    for s in ["Dopasowano", "Do weryfikacji", "Brak dopasowania"]:
        print(f"  {STATUS_ID_MAP.get(s, '?')} {s:25s} {statuses.get(s, 0)}")
    print()
    print("Metody:")
    for m in ["oe", "engine_code", "atrybuty", "fuzzy", "brak"]:
        if methods.get(m, 0):
            print(f"  {m:20s} {methods[m]}")

    if created_clients:
        print()
        print("=" * 80)
        print(f"LEADY ZAPISANE DO DGRED ({leads_total}):")
        print("=" * 80)
        for cl in created_clients:
            dup_tag = " [DUP→nota]" if cl["is_dup"] else ""
            print(f"\n  Client ID: {cl['client_id']}{dup_tag}")
            print(f"  Link:      {DGRED_CLIENT_URL}/{cl['client_id']}")
            print(f"  Autor:     {cl['name']}")
            print(f"  Status:    {cl['status_id']} = {cl['status']}")
            print(f"  Match:     {cl['method']} ({cl['confidence']:.0%}), {cl['n_offers']} ofert")
            print(f"  Marka/Kat: {cl['brand']} / {cl['category']}")
            print(f"  OE:        {cl['oe']}")
            print(f"  Post:      {cl['post_text']}...")
    else:
        print()
        print("Brak leadów do zapisania (wszystkie posty odrzucone)")

    print()
    print("=" * 80)


if __name__ == "__main__":
    main()
