#!/usr/bin/env python3
"""DLRI Lead Engine API.

- POST /lead   : public lead capture (landing page). Appends to leads.jsonl,
                 fires a Telegram alert, best-effort creates a person in Twenty.
- GET  /health : liveness
- everything else : static files served from this directory (landing.html, privacy.html, terms.html)

Design notes:
- Public endpoint (browsers post it), so it is protected with a honeypot field,
  a per-IP rate limit, and a payload size cap -- NOT with the dashboard token.
- Secrets loaded from ~/.openclaw/credentials/*.env (0600), never hardcoded.
"""
import os, json, time, threading, datetime
import urllib.request, urllib.parse
from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer

HERE = os.path.dirname(os.path.abspath(__file__))
LEADS = os.path.join(HERE, "leads.jsonl")
CRED_DLRI = "/home/openclaw/.openclaw/credentials/dlri.env"
CRED_BOTS = "/home/openclaw/.openclaw/credentials/bot-tokens.env"
CRED_TWENTY = "/home/openclaw/.openclaw/credentials/twenty-dlri.env"
BIND = os.environ.get("DLRI_LEAD_BIND", "0.0.0.0")
PORT = int(os.environ.get("DLRI_LEAD_PORT", "9096"))

TWENTY_ENABLE = os.environ.get("TWENTY_ENABLE", "0") == "1"
MAX_BODY = 8 * 1024
RATE_WINDOW, RATE_MAX = 60, 8
_hits = {}
_lock = threading.Lock()


def _cred(path, key):
    try:
        for line in open(path):
            line = line.strip()
            if line and not line.startswith("#") and "=" in line:
                k, v = line.split("=", 1)
                if k.strip() == key:
                    return v.strip()
    except Exception:
        pass
    return None


TELEGRAM_TOKEN = _cred(CRED_BOTS, "MRPABLO_BOT_TOKEN")
# Lead alerts route to the DLRI Telegram group, "Leads" topic (see dlri.env).
TELEGRAM_CHAT = _cred(CRED_DLRI, "DLRI_LEADS_CHAT") or _cred(CRED_DLRI, "DLRI_TELEGRAM_CHAT")
TELEGRAM_THREAD = _cred(CRED_DLRI, "DLRI_LEADS_THREAD")


def send_telegram(text):
    if not (TELEGRAM_TOKEN and TELEGRAM_CHAT):
        return False
    payload = {"chat_id": TELEGRAM_CHAT, "text": text,
               "disable_web_page_preview": "true"}
    if TELEGRAM_THREAD:
        payload["message_thread_id"] = TELEGRAM_THREAD
    try:
        data = urllib.parse.urlencode(payload).encode()
        urllib.request.urlopen(urllib.request.Request(
            "https://api.telegram.org/bot%s/sendMessage" % TELEGRAM_TOKEN, data=data), timeout=10).read()
        return True
    except Exception:
        return False


def create_twenty_person(p):
    """Best-effort. Requires TWENTY_ENABLE=1 and a matching workspace schema."""
    if not TWENTY_ENABLE:
        return None
    base = _cred(CRED_TWENTY, "TWENTY_API_URL")
    key = _cred(CRED_TWENTY, "TWENTY_API_KEY")
    if not (base and key):
        return None
    parts = (p.get("name") or "").split()
    first, last = (parts[0] if parts else "Lead"), (" ".join(parts[1:]) or "-")
    body = json.dumps({
        "name": {"firstName": first, "lastName": last},
        "emails": {"primaryEmail": p.get("email", "")},
        "phones": {"primaryPhoneNumber": p.get("phone", ""), "primaryPhoneCountryCode": "US"},
    }).encode()
    req = urllib.request.Request(base.rstrip("/") + "/rest/people", data=body,
                                 headers={"Content-Type": "application/json",
                                          "Authorization": "Bearer " + key})
    try:
        return urllib.request.urlopen(req, timeout=10).read().decode()[:200]
    except Exception as e:
        return "err: %s" % e


def allowed(ip):
    now = time.time()
    with _lock:
        arr = [t for t in _hits.get(ip, []) if now - t < RATE_WINDOW]
        arr.append(now)
        _hits[ip] = arr
        return len(arr) <= RATE_MAX


class Handler(SimpleHTTPRequestHandler):
    def log_message(self, *a):  # quiet
        pass

    def _json(self, code, obj):
        b = json.dumps(obj).encode()
        self.send_response(code)
        self.send_header("Content-Type", "application/json")
        self.send_header("Access-Control-Allow-Origin", "*")
        self.send_header("Content-Length", str(len(b)))
        self.end_headers()
        self.wfile.write(b)

    def do_OPTIONS(self):
        self.send_response(204)
        self.send_header("Access-Control-Allow-Origin", "*")
        self.send_header("Access-Control-Allow-Headers", "Content-Type")
        self.send_header("Access-Control-Allow-Methods", "POST, OPTIONS")
        self.end_headers()

    def _lead_error(self):
        """Browser calls this when a /lead submission fails, so we can alert."""
        ip = self.client_address[0]
        if not allowed(ip):
            return self._json(429, {"error": "slow down"})
        n = int(self.headers.get("Content-Length", 0) or 0)
        if n > MAX_BODY:
            return self._json(413, {"error": "too big"})
        try:
            p = json.loads(self.rfile.read(n) or b"{}")
        except Exception:
            p = {}
        rec = {"_ip": ip, "_received_at": datetime.datetime.now(datetime.timezone.utc).isoformat(),
               "name": p.get("name"), "email": p.get("email"), "phone": p.get("phone"),
               "interest": p.get("interest"), "source": p.get("source")}
        with _lock:
            with open(os.path.join(HERE, "lead_errors.jsonl"), "a") as f:
                f.write(json.dumps(rec) + "\n")
        if self.headers.get("X-Test") != "1":
            send_telegram("\u26A0\uFE0F WEBSITE FORM SUBMISSION FAILED\n%s\n%s | %s | %s\nVisitor saw a 'contact me directly' fallback. Check the lead endpoint + follow up." % (
                rec["name"] or "?", rec["phone"] or "?", rec["email"] or "?", rec["interest"] or "-"))
        return self._json(200, {"ok": True})

    def do_GET(self):
        if self.path.split("?")[0] == "/health":
            return self._json(200, {"ok": True, "leads": sum(1 for _ in open(LEADS)) if os.path.exists(LEADS) else 0})
        return super().do_GET()

    def do_POST(self):
        path = self.path.split("?")[0]
        if path == "/lead-error":
            return self._lead_error()
        if path != "/lead":
            return self._json(404, {"error": "not found"})
        ip = self.client_address[0]
        if not allowed(ip):
            return self._json(429, {"error": "slow down"})
        n = int(self.headers.get("Content-Length", 0) or 0)
        if n > MAX_BODY:
            return self._json(413, {"error": "too big"})
        try:
            p = json.loads(self.rfile.read(n) or b"{}")
        except Exception:
            return self._json(400, {"error": "bad json"})
        if p.get("company"):  # honeypot -> pretend success, drop
            return self._json(200, {"ok": True})

        p["_ip"] = ip
        p["_received_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
        with _lock:
            with open(LEADS, "a") as f:
                f.write(json.dumps(p) + "\n")

        who = p.get("name", "?")
        note = "%s | %s | age %s | zip %s | %s" % (
            p.get("phone", "?"), p.get("email", "?"), p.get("age", "?"), p.get("zip", "?"),
            p.get("has_coverage", ""))
        if self.headers.get("X-Test") != "1":
            send_telegram("\U0001F4E5 NEW WEBSITE LEAD \u2014 %s\n%s\nSource: %s\n\u2192 call within 5 min" % (
                who, note, p.get("source", "landing")))
            create_twenty_person(p)
        return self._json(200, {"ok": True})


if __name__ == "__main__":
    os.chdir(HERE)
    srv = ThreadingHTTPServer((BIND, PORT), Handler)
    print("DLRI lead API on %s:%d (twenty=%s)" % (BIND, PORT, TWENTY_ENABLE), flush=True)
    srv.serve_forever()
