#!/usr/bin/env python3
"""Run Maestro flows through the Maestro MCP server so the Maestro Viewer is always up.

The founder watches every test run: the device and the live command list in the Maestro Viewer
(http://localhost:7777). `maestro test` never starts the Viewer; only `maestro mcp` does. This
runner starts `maestro mcp`, opens the Viewer in the browser, runs each flow file with the MCP
`run` tool, and reports pass/fail per flow.

    maestro-live --device <udid> [--tags smoke] [--junit out.xml] [--json out.json] <flow|dir>...

A directory is expanded to its top-level *.yaml flows (config.yaml excluded); --tags keeps only
flows whose `tags:` list contains one of the given tags. Port 7777 is used when free, otherwise
the next free port (another session's MCP may hold 7777); the URL actually used is printed and
opened. Exit code 0 only if every flow passed; 3 if another simulator's XCTest driver holds port
22087 (`maestro mcp` always uses that port, so the run would drive the wrong simulator).
"""
from __future__ import annotations

import argparse
import json
import os
import re
import socket
import subprocess
import sys
import time
from pathlib import Path
from xml.sax.saxutils import escape

MAESTRO = os.environ.get("MAESTRO_BIN") or str(Path.home() / ".maestro" / "bin" / "maestro")
JAVA_HOME = "/opt/homebrew/opt/openjdk@17/libexec/openjdk.jdk/Contents/Home"


def free_port(start: int = 7777) -> int:
    for port in range(start, start + 50):
        with socket.socket() as s:
            if s.connect_ex(("127.0.0.1", port)) != 0:
                return port
    raise SystemExit("no free port for the Maestro Viewer")


# `maestro mcp` has no driver-port option: it always talks to the XCTest driver on 22087. If a
# driver for ANOTHER simulator already listens there, every command of this run would read and tap
# that simulator instead (seen live: a run on a fresh simulator read another, already running demo
# simulator). So refuse to start in that case, and stop the driver this run started when it ends.
XCTEST_PORT = 22087


def driver_on_port(port: int = XCTEST_PORT) -> tuple[int, str] | None:
    """(pid, simulator udid or '') of the process listening on the XCTest driver port."""
    r = subprocess.run(["lsof", "-nP", f"-iTCP:{port}", "-sTCP:LISTEN", "-Fp"],
                       capture_output=True, text=True, check=False)
    pids = [int(line[1:]) for line in r.stdout.splitlines() if line.startswith("p")]
    if not pids:
        return None
    cmd = subprocess.run(["ps", "-o", "command=", "-p", str(pids[0])],
                         capture_output=True, text=True, check=False).stdout
    m = re.search(r"/Devices/([0-9A-F-]{36})/", cmd)
    return pids[0], (m.group(1) if m else "")


def flow_tags(path: Path) -> set[str]:
    tags, in_tags = set(), False
    for line in path.read_text(encoding="utf-8").splitlines():
        if line.strip() == "---":
            break
        if line.startswith("tags:"):
            in_tags = True
            rest = line[5:].strip()
            if rest.startswith("["):
                tags |= {t.strip().strip("'\"") for t in rest.strip("[]").split(",") if t.strip()}
                in_tags = False
            continue
        if in_tags:
            if line.lstrip().startswith("- "):
                tags.add(line.strip()[2:].strip().strip("'\""))
            elif line.strip():
                in_tags = False
    return tags


def expand(targets: list[str], tags: set[str]) -> list[Path]:
    files: list[Path] = []
    for t in targets:
        p = Path(t)
        if p.is_dir():
            files += sorted(f for f in p.glob("*.yaml") if f.name != "config.yaml")
            files += sorted(f for f in (p / "flows").glob("*.yaml")) if (p / "flows").is_dir() else []
        else:
            files.append(p)
    if tags:
        files = [f for f in files if flow_tags(f) & tags]
    return files


class Mcp:
    def __init__(self, port: int, workdir: str):
        env = dict(os.environ, MAESTRO_CLI_NO_ANALYTICS="true",
                   MAESTRO_CLI_ANALYSIS_NOTIFICATION_DISABLED="true", MAESTRO_DISABLE_UPDATE_CHECK="true")
        env.setdefault("JAVA_HOME", JAVA_HOME)
        env.setdefault("MAESTRO_DRIVER_STARTUP_TIMEOUT", "120000")
        self.proc = subprocess.Popen([MAESTRO, "mcp", "--viewer-port", str(port), "--working-dir", workdir],
                                     stdin=subprocess.PIPE, stdout=subprocess.PIPE, text=True, env=env)
        self.n = 0

    def call(self, method: str, params: dict | None = None, notify: bool = False):
        msg = {"jsonrpc": "2.0", "method": method, "params": params or {}}
        if not notify:
            self.n += 1
            msg["id"] = self.n
        self.proc.stdin.write(json.dumps(msg) + "\n")
        self.proc.stdin.flush()
        if notify:
            return None
        while True:
            line = self.proc.stdout.readline()
            if not line:
                raise SystemExit("maestro mcp exited unexpectedly")
            try:
                r = json.loads(line)
            except ValueError:
                continue
            if r.get("id") == self.n:
                return r

    def close(self):
        self.proc.terminate()
        try:
            self.proc.wait(10)
        except subprocess.TimeoutExpired:
            self.proc.kill()


def run_result(resp: dict) -> tuple[bool, str]:
    res = resp.get("result") or {}
    texts = [c.get("text", "") for c in res.get("content", []) if c.get("type") == "text"]
    for t in texts:
        if t.lstrip().startswith("{"):
            try:
                d = json.loads(t, strict=False)
            except ValueError:
                continue
            ok = bool(d.get("success"))
            detail = ""
            for r in d.get("results") or []:
                if not r.get("success"):
                    detail = json.dumps(r)[:600]
            return ok, detail
    return (not res.get("isError")) and bool(texts), " | ".join(texts)[:600]


def main() -> int:
    ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
    ap.add_argument("targets", nargs="+")
    ap.add_argument("--device", required=True, help="simulator udid")
    ap.add_argument("--tags", default="", help="comma-separated tags to include")
    ap.add_argument("--junit", help="write a JUnit report here")
    ap.add_argument("--json", help="write per-flow results as JSON here")
    ap.add_argument("--no-open", action="store_true", help="print the Viewer URL but do not open it")
    ap.add_argument("--head-start", type=float, default=4.0,
                    help="seconds to wait after opening the Viewer before the first flow")
    a = ap.parse_args()

    files = expand(a.targets, {t.strip() for t in a.tags.split(",") if t.strip()})
    if not files:
        print("no flows matched", file=sys.stderr)
        return 2
    held = driver_on_port()
    if held and held[1] != a.device:
        print(f"XCTest driver port {XCTEST_PORT} is held by pid {held[0]} "
              f"(simulator {held[1] or 'unknown'}), not {a.device}. `maestro mcp` cannot use another "
              "port, so this run would drive that simulator. Stop that Maestro session first.",
              file=sys.stderr)
        return 3
    port = free_port()
    mcp = Mcp(port, os.getcwd())
    results = []
    try:
        mcp.call("initialize", {"protocolVersion": "2025-06-18", "capabilities": {},
                                "clientInfo": {"name": "maestro-live", "version": "1"}})
        mcp.call("notifications/initialized", notify=True)
        mcp.call("tools/call", {"name": "open_maestro_viewer", "arguments": {}})
        url = f"http://localhost:{port}"
        print(f"== Maestro Viewer: {url}", flush=True)
        if not a.no_open:
            subprocess.run(["open", url], check=False)
        time.sleep(a.head_start)
        for f in files:
            t0 = time.time()
            for attempt in range(3):
                resp = mcp.call("tools/call", {"name": "run", "arguments": {
                    "files": [str(f.resolve())], "device_id": a.device}})
                ok, detail = run_result(resp)
                # The XCTest driver sometimes is not up yet right after an install; that is
                # infrastructure, not a failing flow, so retry instead of reporting it.
                if ok or "unreachable" not in detail.lower():
                    break
                print(f"..  {f}: device unreachable, retrying", flush=True)
                time.sleep(5)
            dt = time.time() - t0
            results.append({"name": f.stem, "file": str(f), "passed": ok, "seconds": round(dt, 1),
                            "tags": sorted(flow_tags(f)), "detail": "" if ok else detail})
            print(f"{'PASS' if ok else 'FAIL'}  {f}  ({dt:.0f}s){'' if ok else '  ' + detail}", flush=True)
    finally:
        mcp.close()
        # Free port 22087 for the next run on another simulator: stop the driver this run started.
        mine = driver_on_port()
        if not held and mine and mine[1] == a.device:
            subprocess.run(["kill", str(mine[0])], check=False)

    if a.json:
        Path(a.json).parent.mkdir(parents=True, exist_ok=True)
        Path(a.json).write_text(json.dumps(results, indent=2))
    if a.junit:
        fails = sum(not r["passed"] for r in results)
        cases = "".join(
            f'<testcase name="{escape(r["name"])}" classname="maestro" time="{r["seconds"]}">'
            + ("" if r["passed"] else f'<failure message="{escape(r["detail"][:200])}"/>')
            + "</testcase>" for r in results)
        Path(a.junit).parent.mkdir(parents=True, exist_ok=True)
        Path(a.junit).write_text(f'<?xml version="1.0"?><testsuites><testsuite name="maestro" '
                                 f'tests="{len(results)}" failures="{fails}">{cases}</testsuite></testsuites>')
    return 0 if results and all(r["passed"] for r in results) else 1


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