#!/usr/bin/env python3
"""Aria QUEUE-Processor (KAR-591) — async Drop-Task-Loop.

Kais legt .md-Tasks in /root/aria/queue/ ab. Dieser Prozessor (per Timer):
- liest jeden Task, erkennt Modus aus Dateiname (RESEARCH/DRAFT/ANALYZE/generisch),
- holt Brain-Kontext (relevante Notes via brain-search) + WEEKLY-FOCUS,
- ruft Claude (Sonnet) → schreibt Output nach /root/aria/generated/,
- verschiebt erledigte Tasks nach queue/processed/ (NIE loeschen),
- Tasks die Tool-Access/Code/Deploy brauchen → queue/needs-session/ (echte Aria-Session noetig),
- Telegram-Notify.

Aufruf: python3 aria-queue-processor.py [--no-telegram] [--model sonnet|opus]
"""
from __future__ import annotations
import argparse
import os
import re
import shutil
import subprocess
import sys
import urllib.parse
import urllib.request
from datetime import datetime, timezone
from pathlib import Path

QUEUE = Path("/root/aria/queue")
PROCESSED = QUEUE / "processed"
NEEDS = QUEUE / "needs-session"
GENERATED = Path("/root/aria/generated")
FOCUS = Path("/root/aria/brain/WEEKLY-FOCUS.md")
BRAIN_SEARCH = "/root/aria/scripts/aria-brain-search.py"
CHAT_ID = "1164395546"
NOW = datetime.now(timezone.utc)
MODEL = {"sonnet": "claude-sonnet-4-6", "opus": "claude-opus-4-7"}

# tasks needing real tool-access (code/deploy) — flag, don't half-generate
TOOL_KEYWORDS = re.compile(r"\b(deploy|implement|code|build|fix bug|fixe|migration|push|PR\b|pull request|refactor|baue ein|programmier)", re.I)

MODES = {
    "RESEARCH": "Erstelle einen strukturierten Recherche-Brief: Kernfrage, 3-5 wichtigste Befunde mit Begruendung, was es fuer Kais/Aria bedeutet, offene Fragen. Nutze den Brain-Kontext wo relevant.",
    "DRAFT": "Erstelle einen ersten Entwurf (Text/Copy/Doc) gemaess dem Task. Aria-Stil: direkt, konkret, kein Marketing-Sprech, deutsche Umlaute. Markiere Stellen die Kais-Input brauchen mit [KAIS?].",
    "ANALYZE": "Analysiere strukturiert: Kernaussage, Staerken, Schwaechen/Kritik (Pflicht), Aria-First-Bezug (was bedeutet es fuer uns), Empfehlung mit Begruendung.",
}


def tg_token():
    for p in ["/root/.claude/channels/telegram/.env", "/root/aria/.env"]:
        f = Path(p)
        if f.exists():
            for line in f.read_text().splitlines():
                if "=" in line and not line.strip().startswith("#"):
                    k, _, v = line.partition("=")
                    if k.strip() in ("TG_TOKEN", "TELEGRAM_BOT_TOKEN", "TELEGRAM_TOKEN") and v.strip():
                        return v.strip().strip('"').strip("'")
    return os.environ.get("TG_TOKEN")


def md2(s):
    return re.sub(r"([_*\[\]()~`>#+\-=|{}.!\\])", r"\\\1", str(s))


def send_tg(text):
    tok = tg_token()
    if not tok:
        return False
    data = urllib.parse.urlencode({"chat_id": CHAT_ID, "text": text,
        "parse_mode": "MarkdownV2", "disable_web_page_preview": "true"}).encode()
    try:
        with urllib.request.urlopen(urllib.request.Request(
                f"https://api.telegram.org/bot{tok}/sendMessage", data=data, method="POST"), timeout=15) as r:
            return r.status == 200
    except Exception as e:
        print(f"[queue] tg err: {e}", file=sys.stderr)
        return False


def weekly_focus():
    if FOCUS.exists():
        m = re.search(r"\*Top-Fokus:\*\s*(.+)", FOCUS.read_text())
        if m:
            return m.group(1).strip()
    return ""


def brain_context(query):
    try:
        out = subprocess.run(["python3", BRAIN_SEARCH, query, "--top", "3"],
                             capture_output=True, text=True, timeout=20)
        return out.stdout[:2500]
    except Exception:
        return ""


def detect_mode(name):
    up = name.upper()
    for m in MODES:
        if up.startswith(m):
            return m
    return "GENERIC"


def process(client, task_path, model):
    name = task_path.stem
    content = task_path.read_text(errors="replace")
    mode = detect_mode(name)
    sys_prompt = (
        "Du bist Aria, Kais' KI-Partnerin, und bearbeitest einen async-QUEUE-Task waehrend Kais nicht da ist. "
        + MODES.get(mode, "Bearbeite den Task sorgfaeltig und liefere ein nuetzliches Ergebnis.")
        + " Sei konkret, ehrlich, kein Floskel. Wenn du etwas nicht ohne Tool-Access/Live-Daten kannst, sag es klar."
    )
    focus = weekly_focus()
    ctx = brain_context(name.replace("-", " ") + " " + content[:200])
    usr = (f"TASK ({mode}): {name}\n\n{content}\n\n"
           f"--- Wochen-Fokus: {focus or 'keiner'} ---\n"
           f"--- Brain-Kontext (relevante Notes) ---\n{ctx}")
    try:
        resp = client.messages.create(model=model, max_tokens=2000,
                                      system=sys_prompt, messages=[{"role": "user", "content": usr}])
        text = "".join(b.text for b in resp.content if hasattr(b, "text"))
    except Exception as e:
        return None, f"API-Fehler: {e}"
    GENERATED.mkdir(parents=True, exist_ok=True)
    out_path = GENERATED / f"{NOW.strftime('%Y-%m-%d')}-{name}.md"
    header = (f"---\ntask: {name}\nmode: {mode}\ngenerated: {NOW.isoformat()}\n"
              f"model: {model}\nsource_queue: {task_path.name}\nstatus: needs-review\n---\n\n")
    out_path.write_text(header + text)
    return out_path, None


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--no-telegram", action="store_true")
    ap.add_argument("--model", choices=["sonnet", "opus"], default="sonnet")
    args = ap.parse_args()

    tasks = [p for p in sorted(QUEUE.glob("*.md")) if p.name != "README.md"]
    if not tasks:
        print("[queue] no tasks")
        return 0

    api_key = os.environ.get("ANTHROPIC_API_KEY")
    if not api_key:
        print("[queue] ANTHROPIC_API_KEY missing")
        return 1
    from anthropic import Anthropic
    client = Anthropic(api_key=api_key)
    PROCESSED.mkdir(parents=True, exist_ok=True)
    NEEDS.mkdir(parents=True, exist_ok=True)

    done, flagged, failed = [], [], []
    for t in tasks:
        content = t.read_text(errors="replace")
        if TOOL_KEYWORDS.search(t.stem) or TOOL_KEYWORDS.search(content[:300]):
            shutil.move(str(t), str(NEEDS / t.name))
            flagged.append(t.name)
            continue
        out_path, err = process(client, t, MODEL[args.model])
        if err:
            failed.append((t.name, err))
        else:
            shutil.move(str(t), str(PROCESSED / t.name))
            done.append(out_path.name)

    print(f"[queue] done={len(done)} flagged={len(flagged)} failed={len(failed)}")

    if not args.no_telegram and (done or flagged or failed):
        L = ["*QUEUE\\-Processor* " + md2(NOW.strftime("%Y-%m-%d %H:%M")), ""]
        if done:
            L.append(f"✅ *{len(done)} Tasks bearbeitet* → `generated/`")
            for d in done[:6]:
                L.append(f"• {md2(d)}")
        if flagged:
            L.append("")
            L.append(f"🔧 *{len(flagged)} brauchen echte Aria\\-Session* \\(Code/Deploy\\) → `queue/needs-session/`")
            for f in flagged[:5]:
                L.append(f"• {md2(f)}")
        if failed:
            L.append("")
            L.append(f"⚠️ {len(failed)} Fehler")
        send_tg("\n".join(L))
    return 0


if __name__ == "__main__":
    sys.exit(main())
