# -*- coding: utf-8 -*-
"""
worker.py - Arbeitet die Job-Warteschlange ab (Hintergrundverarbeitung).

Wird regelmaessig aufgerufen (cPanel-Cronjob, z.B. jede Minute):
    /pfad/zu/python  /pfad/zu/worker.py

Ablauf je Job:
  1) Job aus jobs/pending holen -> nach jobs/processing verschieben (sperrt ihn)
  2) Geburtsdaten validieren und zerlegen (Datum, Zeit)
  3) Geburtsort -> Koordinaten + Zeitzone (geo.resolve_birthplace)
  4) generate_band.generate(...) aufrufen  (Berechnung + Claude-API + PDF)
  5) PDF per E-Mail an Kunden (email_sender.send_band_pdf), BCC zur Pruefung
  6) Job -> jobs/done  (oder bei Fehler -> jobs/failed, mit Fehlertext)

Robust: Jeder Job ist eine eigene Datei. Faellt einer aus, laufen die anderen
weiter. Ein 'processing'-Job, der zu lange haengt, wird beim naechsten Lauf
wieder freigegeben (siehe STALE_SECONDS).

Umgebungsvariablen: ANTHROPIC_API_KEY, SMTP_* (siehe email_sender), JOBS_DIR.
"""
import os
import json
import time
import shutil
import traceback
import datetime as dt

HERE = os.path.dirname(os.path.abspath(__file__))
JOBS_DIR = os.environ.get("JOBS_DIR", os.path.join(HERE, "jobs"))
PENDING = os.path.join(JOBS_DIR, "pending")
PROCESSING = os.path.join(JOBS_DIR, "processing")
DONE = os.path.join(JOBS_DIR, "done")
FAILED = os.path.join(JOBS_DIR, "failed")
for d in (PENDING, PROCESSING, DONE, FAILED):
    os.makedirs(d, exist_ok=True)

MAX_ATTEMPTS = 3
STALE_SECONDS = 1800   # haengende 'processing'-Jobs nach 30 min neu freigeben
LOCK_FILE = os.path.join(JOBS_DIR, "worker.lock")

# Band-Titel fuer die E-Mail
BAND_TITEL = {
    "I": "Band I \u2013 Grundprofil",
    "II": "Band II \u2013 Energie Timeline",
    "III": "Band III \u2013 Business & Finanzen",
    "IV": "Band IV \u2013 Liebe & Partnerschaft",
    "V": "Band V \u2013 Familie & Beziehungen",
}


def _acquire_lock():
    """Verhindert, dass zwei Worker-Laeufe gleichzeitig arbeiten."""
    if os.path.exists(LOCK_FILE):
        age = time.time() - os.path.getmtime(LOCK_FILE)
        if age < STALE_SECONDS:
            return False
    with open(LOCK_FILE, "w") as f:
        f.write(str(time.time()))
    return True


def _release_lock():
    try:
        os.remove(LOCK_FILE)
    except OSError:
        pass


def _requeue_stale():
    """Gibt zu lange haengende processing-Jobs wieder frei."""
    now = time.time()
    for fn in os.listdir(PROCESSING):
        p = os.path.join(PROCESSING, fn)
        if now - os.path.getmtime(p) > STALE_SECONDS:
            shutil.move(p, os.path.join(PENDING, fn))


def _parse_datum_zeit(datum_str, zeit_str):
    """'31.03.1976' + '15:00' -> (1976,3,31,15,0). Wirft ValueError bei Mist."""
    d = dt.datetime.strptime(datum_str.strip(), "%d.%m.%Y")
    if zeit_str and zeit_str.strip():
        t = dt.datetime.strptime(zeit_str.strip(), "%H:%M")
        stunde, minute = t.hour, t.minute
    else:
        # Unbekannte Geburtszeit -> 12:00 als neutraler Default.
        # (Shop-Entscheidung: spaeter ggf. AC/Haeuser-Hinweis ins PDF.)
        stunde, minute = 12, 0
    return d.year, d.month, d.day, stunde, minute


def process_job(job):
    """Verarbeitet einen einzelnen Job. Wirft bei Fehler (Aufrufer faengt)."""
    import geo
    import generate_band

    band = job["band"]
    name = job.get("kunde_name") or "Kundin/Kunde"
    geschlecht = (job.get("geschlecht") or "").strip().lower()
    email = job.get("kunde_email")
    if not email:
        raise ValueError("Keine Kunden-E-Mail in der Bestellung.")

    # 1) Datum/Zeit zerlegen
    jahr, monat, tag, stunde, minute = _parse_datum_zeit(
        job["geburtsdatum"], job.get("geburtszeit", ""))

    # 2) Ort -> Koordinaten + Zeitzone
    #    Bevorzugt: Koordinaten kamen schon aus dem Checkout (Ortsauswahl).
    #    Sonst: Geocoding anhand Ort + Land + PLZ (Fallback).
    lat_roh = (job.get("geburt_lat") or "").strip()
    lon_roh = (job.get("geburt_lon") or "").strip()
    if lat_roh and lon_roh:
        geo_info = geo.tz_from_coords(lat_roh, lon_roh)
        geo_info["display_name"] = job.get("geburtsort", "")
        geo_info["quelle"] = "checkout-auswahl"
    else:
        geo_info = geo.resolve_birthplace(
            job["geburtsort"],
            land=job.get("geburtsland"),
            plz=job.get("geburts_plz"),
        )
        geo_info["quelle"] = "geocoding-fallback"
    lat, lon, tz = geo_info["lat"], geo_info["lon"], geo_info["tz"]

    geburt_text = f"{job['geburtsdatum']} \u00b7 {job.get('geburtszeit','')} Uhr \u00b7 {job['geburtsort']}".replace("  ", " ")

    # 3) Pipeline: Berechnung + Claude-API + PDF
    #    Haeusersystem Placidus (b'P') als Shop-Standard - anpassbar.
    pdf_path, json_path, usage = generate_band.generate(
        name=name,
        geburt_text=geburt_text,
        jahr=jahr, monat=monat, tag=tag,
        stunde_lokal=stunde, minute=minute,
        lat=lat, lon=lon, tz=tz, hsys=b'P', band=band,
        geschlecht=geschlecht,
    )

    # 4) PDF per Mail an den Kunden
    import email_sender
    titel = BAND_TITEL.get(band, f"Band {band}")
    email_sender.send_band_pdf(email, name, titel, pdf_path)

    return {"pdf": pdf_path, "usage": usage, "geo": geo_info}


def run_once():
    if not _acquire_lock():
        print("Anderer Worker laeuft - uebersprungen.")
        return
    try:
        _requeue_stale()
        pending = sorted(os.listdir(PENDING))
        if not pending:
            print("Keine Jobs.")
            return
        for fn in pending:
            if not fn.endswith(".json"):
                continue
            src = os.path.join(PENDING, fn)
            proc = os.path.join(PROCESSING, fn)
            try:
                shutil.move(src, proc)   # Job sperren
            except OSError:
                continue                 # anderer Lauf war schneller
            with open(proc, "r", encoding="utf-8") as f:
                job = json.load(f)
            job["attempts"] = job.get("attempts", 0) + 1

            try:
                print(f"[{fn}] Band {job['band']} fuer {job.get('kunde_email')} ...")
                result = process_job(job)
                job["result"] = {"pdf": os.path.basename(result["pdf"]),
                                 "usage": result["usage"]}
                job["finished_at"] = time.time()
                with open(os.path.join(DONE, fn), "w", encoding="utf-8") as f:
                    json.dump(job, f, ensure_ascii=False, indent=2)
                os.remove(proc)
                print(f"[{fn}] OK -> versendet.")
            except Exception as e:
                job["error"] = f"{type(e).__name__}: {e}"
                job["traceback"] = traceback.format_exc()
                if job["attempts"] < MAX_ATTEMPTS:
                    # Zurueck in die Schlange fuer einen weiteren Versuch
                    with open(src, "w", encoding="utf-8") as f:
                        json.dump(job, f, ensure_ascii=False, indent=2)
                    os.remove(proc)
                    print(f"[{fn}] FEHLER ({job['attempts']}/{MAX_ATTEMPTS}), erneut eingereiht: {job['error']}")
                else:
                    with open(os.path.join(FAILED, fn), "w", encoding="utf-8") as f:
                        json.dump(job, f, ensure_ascii=False, indent=2)
                    os.remove(proc)
                    print(f"[{fn}] ENDGUELTIG FEHLGESCHLAGEN: {job['error']}")
    finally:
        _release_lock()


if __name__ == "__main__":
    run_once()
