#!/usr/bin/env python3
"""
XL 107.1 EAS local injector
============================
Runs on the same PC as your Icecast server. Receives EAS alert audio (MP3)
from the https://eas-cap.pages.dev EAS console and injects it into your
local Icecast as a source connection, paced in real time.

Why this exists: the EAS console is a website, and websites cannot push
audio through the zrok tunnel to Icecast (the tunnel blocks request
bodies). This little program runs on YOUR pc, next to Icecast, so it can
talk to Icecast directly at 127.0.0.1:8000 with no tunnel in the way.

Setup (2 minutes):
  1. Edit SOURCE_PASSWORD below (your Icecast <source-password>).
  2. Run:   python eas_injector.py
     (Windows: double-click it, or:  py eas_injector.py)
  3. Keep this window open.
  4. Open https://eas-cap.pages.dev in your browser.
     The console will show "Local injector: CONNECTED".
  5. Press "Test EAS" — you will hear it on your Icecast stream.

Requires: Python 3.8+ (standard library only, nothing to install).
"""

import base64
import json
import socket
import time
from http.server import BaseHTTPRequestHandler, HTTPServer

# ------------------------- config -------------------------
ICECAST_HOST = "127.0.0.1"      # Icecast runs on this PC
ICECAST_PORT = 8000             # Icecast port (plain HTTP, not the zrok URL)
ICECAST_USER = "source"         # Icecast source username
SOURCE_PASSWORD = "PUT-YOUR-SOURCE-PASSWORD-HERE"   # <-- EDIT THIS (your Icecast source password)
DEFAULT_MOUNT = "/XL107.1"
LISTEN_PORT = 8899              # this injector's own port (browser talks here)
MP3_BYTES_PER_SEC = 16000       # 128 kbps -> paced in real time
# ----------------------------------------------------------


def inject_to_icecast(mount, mp3_bytes):
    """Open a SOURCE connection to the local Icecast and stream the MP3
    paced at the mount bitrate so listeners hear it in real time."""
    creds = base64.b64encode(
        ("%s:%s" % (ICECAST_USER, SOURCE_PASSWORD)).encode()
    ).decode()
    s = socket.create_connection((ICECAST_HOST, ICECAST_PORT), timeout=10)
    try:
        headers = (
            "SOURCE %s HTTP/1.0\r\n"
            "Authorization: Basic %s\r\n"
            "Content-Type: audio/mpeg\r\n"
            "Ice-Public: 0\r\n"
            "Ice-Name: XL 107.1 EAS Alert\r\n"
            "Ice-Bitrate: 128\r\n"
            "Content-Length: %d\r\n"
            "\r\n"
        ) % (mount, creds, len(mp3_bytes))
        s.sendall(headers.encode("latin1"))
        resp = b""
        while b"\r\n\r\n" not in resp:
            chunk = s.recv(1024)
            if not chunk:
                break
            resp += chunk
        status_line = resp.split(b"\r\n", 1)[0].decode("latin1", "replace")
        print("  icecast says: %s" % status_line, flush=True)
        if "200" not in status_line:
            if "401" in status_line:
                return False, "Icecast rejected the password (401). Check SOURCE_PASSWORD."
            return False, "Icecast refused: %s" % status_line
        # stream paced in real time
        chunk_size = 4000
        pause = chunk_size / MP3_BYTES_PER_SEC  # 0.25 s
        for off in range(0, len(mp3_bytes), chunk_size):
            s.sendall(mp3_bytes[off:off + chunk_size])
            time.sleep(pause)
        secs = len(mp3_bytes) / MP3_BYTES_PER_SEC
        return True, "streamed %.1f s of EAS audio to %s" % (secs, mount)
    finally:
        try:
            s.close()
        except Exception:
            pass


def relay_stream_url(mount, url, seconds):
    """Fetch a remote audio stream and pipe it to the local Icecast
    in real time for up to `seconds`."""
    import urllib.request as urlreq
    try:
        resp = urlreq.urlopen(url, timeout=25)
    except Exception as e:
        return False, "could not open stream URL: %s" % e
    ctype = resp.headers.get("Content-Type", "audio/mpeg").split(";")[0].strip()
    s = socket.create_connection((ICECAST_HOST, ICECAST_PORT), timeout=10)
    try:
        creds = base64.b64encode(
            ("%s:%s" % (ICECAST_USER, SOURCE_PASSWORD)).encode()
        ).decode()
        head = (
            "SOURCE %s HTTP/1.0\r\n"
            "Authorization: Basic %s\r\n"
            "Content-Type: %s\r\n"
            "Ice-Public: 0\r\n"
            "Ice-Name: XL 107.1 EAS Relay\r\n"
            "\r\n"
        ) % (mount, creds, ctype)
        s.sendall(head.encode("latin1"))
        status = b""
        while b"\r\n" not in status:
            d = s.recv(1024)
            if not d:
                break
            status += d
        line = status.split(b"\r\n")[0].decode("latin1", "replace")
        print("  icecast says: %s" % line, flush=True)
        if "200" not in line:
            return False, "Icecast refused: %s" % line
        deadline = time.time() + seconds
        total = 0
        while time.time() < deadline:
            d = resp.read(8192)
            if not d:
                break
            s.sendall(d)
            total += len(d)
        return True, "relayed custom EAS stream to %s (%d bytes)" % (mount, total)
    finally:
        try:
            s.close()
        except Exception:
            pass


class Handler(BaseHTTPRequestHandler):
    def _cors(self):
        self.send_header("Access-Control-Allow-Origin", "*")
        self.send_header("Access-Control-Allow-Methods", "GET, POST, OPTIONS")
        self.send_header("Access-Control-Allow-Headers", "Content-Type")

    def do_OPTIONS(self):
        self.send_response(204)
        self._cors()
        self.end_headers()

    def _send_json(self, code, obj):
        body = json.dumps(obj).encode()
        self.send_response(code)
        self._cors()
        self.send_header("Content-Type", "application/json")
        self.send_header("Content-Length", str(len(body)))
        self.end_headers()
        self.wfile.write(body)

    def do_GET(self):
        if self.path == "/health":
            self._send_json(200, {"ok": True, "service": "eas-injector"})
        else:
            self._send_json(404, {"ok": False})

    def do_POST(self):
        try:
            length = int(self.headers.get("Content-Length", 0))
            data = json.loads(self.rfile.read(length) or b"{}")
        except Exception as e:
            self._send_json(400, {"ok": False, "error": "bad request: %s" % e})
            return
        if self.path == "/inject-stream":
            code, obj = self._handle_inject_stream(data)
            self._send_json(code, obj)
            return
        if self.path != "/inject":
            self._send_json(404, {"ok": False})
            return
        try:
            mp3 = base64.b64decode(data["audioMp3Base64"])
            mount = data.get("mount") or DEFAULT_MOUNT
            if not mount.startswith("/"):
                mount = "/" + mount
        except Exception as e:
            self._send_json(400, {"ok": False, "error": "bad request: %s" % e})
            return
        stamp = time.strftime("%H:%M:%S")
        print("[%s] EAS inject -> %s (%d bytes MP3)" % (stamp, mount, len(mp3)), flush=True)
        try:
            ok, note = inject_to_icecast(mount, mp3)
        except Exception as e:
            ok, note = False, "injector error: %s" % e
        print("  %s: %s" % ("OK" if ok else "FAILED", note), flush=True)
        self._send_json(200 if ok else 502, {"ok": ok, "note": note})

    def _handle_inject_stream(self, data):
        url = data.get("streamUrl")
        if not url:
            return 400, {"ok": False, "error": "missing streamUrl"}
        mount = data.get("mount") or DEFAULT_MOUNT
        if not mount.startswith("/"):
            mount = "/" + mount
        try:
            secs = int(data.get("seconds") or 120)
            if secs <= 0 or secs > 600:
                secs = 120
        except (TypeError, ValueError):
            secs = 120
        stamp = time.strftime("%H:%M:%S")
        print("[%s] EAS stream relay -> %s from %s (%ds)" % (stamp, mount, url, secs), flush=True)
        try:
            ok, note = relay_stream_url(mount, url, secs)
        except Exception as e:
            ok, note = False, "injector error: %s" % e
        print("  %s: %s" % ("OK" if ok else "FAILED", note), flush=True)
        return (200 if ok else 502), {"ok": ok, "note": note}

    def log_message(self, *args):
        pass  # keep the console clean; we print our own lines


if __name__ == "__main__":
    if SOURCE_PASSWORD == "PUT-YOUR-SOURCE-PASSWORD-HERE":
        print("ERROR: open eas_injector.py and set SOURCE_PASSWORD first")
        print("       (it's your Icecast <source-password>, near the top).")
        raise SystemExit(1)
    server = HTTPServer(("127.0.0.1", LISTEN_PORT), Handler)
    print("=" * 60)
    print("XL 107.1 EAS injector running.")
    print("Listening : http://127.0.0.1:%d  (the EAS console talks here)" % LISTEN_PORT)
    print("Target    : %s:%d%s as '%s'" % (ICECAST_HOST, ICECAST_PORT, DEFAULT_MOUNT, ICECAST_USER))
    print("Keep this window OPEN, then open https://eas-cap.pages.dev")
    print("=" * 60)
    try:
        server.serve_forever()
    except KeyboardInterrupt:
        print("\nstopped.")
