#!/usr/bin/env python3
"""
Live interop test: py-ipfs-lite ↔ local Kubo

Starts a fresh local Kubo daemon, connects py-ipfs-lite to it, then:
  1. Measures how long identify takes
  2. Sends an initial ping and measures RTT
  3. Registers an inbound-ping handler and watches for Kubo re-pinging us
  4. Runs for 3 minutes, printing a timestamped event log

Run:  uv run python test_live_interop.py
"""

import logging
import os
import secrets
import subprocess
import sys
import tempfile
import time

import trio

# ── Logging ─────────────────────────────────────────────────────────────────
logging.basicConfig(
    level=logging.WARNING,
    format="%(asctime)s %(levelname)s [%(name)s] %(message)s",
    stream=sys.stderr,
)
# Only show our own event prints cleanly
LOG = logging.getLogger("interop")
LOG.setLevel(logging.DEBUG)
# Route our events to stdout so they don't mix with library warnings on stderr
_h = logging.StreamHandler(sys.stdout)
_h.setFormatter(logging.Formatter("%(message)s"))
LOG.addHandler(_h)
LOG.propagate = False

from libp2p import new_host
from libp2p.crypto.ed25519 import create_new_key_pair
from libp2p.crypto.x25519 import create_new_key_pair as x25519_kp
from libp2p.host.ping import PingService
from libp2p.host.ping import ID as PING_PROTO
from libp2p.peer.id import ID
from libp2p.peer.peerinfo import info_from_p2p_addr
from libp2p.security.noise.transport import Transport as NoiseTransport
from libp2p.security.tls.transport import TLSTransport
from multiaddr import Multiaddr

IDENTIFY_PROTO = "/ipfs/id/1.0.0"
OBSERVE_MINUTES = 3          # how long to watch for re-pings
T0 = time.monotonic()        # global wall-clock reference


def ts() -> str:
    """Return elapsed time as HH:MM:SS.mmm string."""
    elapsed = time.monotonic() - T0
    m, s = divmod(elapsed, 60)
    h, m = divmod(m, 60)
    ms = int((elapsed % 1) * 1000)
    return f"+{int(h):02d}:{int(m):02d}:{int(s):02d}.{ms:03d}"


def event(icon: str, msg: str) -> None:
    LOG.info(f"  {ts()}  {icon}  {msg}")


# ── Kubo helpers ─────────────────────────────────────────────────────────────

def start_kubo(ipfs_path: str):
    env = {**os.environ, "IPFS_PATH": ipfs_path}
    for cmd in [
        ["ipfs", "init", "--profile=test"],
        ["ipfs", "config", "--json", "Addresses.Swarm", '["/ip4/127.0.0.1/tcp/0"]'],
        ["ipfs", "bootstrap", "rm", "--all"],
        ["ipfs", "config", "Addresses.API",     "/ip4/127.0.0.1/tcp/0"],
        ["ipfs", "config", "Addresses.Gateway", "/ip4/127.0.0.1/tcp/0"],
    ]:
        subprocess.run(cmd, env=env, check=True, capture_output=True)

    log_path = os.path.join(ipfs_path, "daemon.log")
    proc = subprocess.Popen(
        ["ipfs", "daemon"], env=env,
        stdout=open(log_path, "w"), stderr=subprocess.STDOUT,
    )
    deadline = time.time() + 30
    while time.time() < deadline:
        time.sleep(0.5)
        if proc.poll() is not None:
            raise RuntimeError("Kubo exited:\n" + open(log_path).read()[-1000:])
        try:
            if "Daemon is ready" in open(log_path).read():
                break
        except FileNotFoundError:
            pass
    else:
        proc.terminate()
        raise RuntimeError("Kubo did not start in 30 s")

    peer_id_str = subprocess.check_output(["ipfs", "id", "-f=<id>"], env=env).decode().strip()
    addrs = subprocess.check_output(["ipfs", "id", "-f=<addrs>"], env=env).decode().strip().splitlines()
    addr = next((a.strip() for a in addrs if "127.0.0.1" in a and "/tcp/" in a), None)
    if not addr:
        proc.terminate()
        raise RuntimeError("No TCP loopback addr")
    return proc, ID.from_base58(peer_id_str), addr


# ── Inbound-ping handler that logs every incoming ping ────────────────────────

class LoggingPingService(PingService):
    """Extends PingService to log every inbound ping from Kubo."""

    def __init__(self, host, ping_log: list):
        super().__init__(host)
        self._ping_log = ping_log

    async def handle_ping(self, stream):
        peer_id = stream.muxed_conn.peer_id
        now = time.monotonic()
        self._ping_log.append(now)
        count = len(self._ping_log)
        if count == 1:
            event("📥", f"Kubo → us  INBOUND PING #{count}  (first re-ping from Kubo)")
        else:
            gap = now - self._ping_log[-2]
            event("📥", f"Kubo → us  INBOUND PING #{count}  (gap since last = {gap:.1f} s)")
        await super().handle_ping(stream)


# ── Identify helper ──────────────────────────────────────────────────────────

async def do_identify(host, kubo_id: ID) -> tuple:
    """Open /ipfs/id/1.0.0 stream, read response, return (elapsed_s, bytes_read)."""
    t_start = time.monotonic()
    stream = await host.new_stream(kubo_id, [IDENTIFY_PROTO])
    data = b""
    try:
        # Identify sends one varint-prefixed protobuf message then half-closes.
        # Read in chunks until EOF or timeout.
        with trio.move_on_after(8):
            while True:
                try:
                    chunk = await stream.read(4096)
                    if not chunk:
                        break
                    data += chunk
                except Exception:
                    break
    finally:
        try:
            await stream.close()
        except Exception:
            pass
    elapsed = time.monotonic() - t_start
    return elapsed, len(data)


# ── Main ─────────────────────────────────────────────────────────────────────

async def main():
    print()
    print("=" * 62)
    print("  LIVE INTEROP: py-ipfs-lite ↔ Kubo  (identify + ping)")
    print("=" * 62)

    # ── 1. Kubo ──────────────────────────────────────────────────
    with tempfile.TemporaryDirectory(prefix="kubo_live_") as ipfs_path:
        event("🚀", "Starting local Kubo daemon …")
        kubo_proc, kubo_id, kubo_addr = start_kubo(ipfs_path)
        event("✅", f"Kubo ready   peer={kubo_id}")
        event("   ", f"             addr={kubo_addr}")

        try:
            # ── 2. py-libp2p host ────────────────────────────────
            event("🔧", "Building py-libp2p host …")
            host_key  = create_new_key_pair()
            noise_key = x25519_kp()
            sec_opt = {
                "/noise":     NoiseTransport(host_key, noise_privkey=noise_key.private_key),
                "/tls/1.0.0": TLSTransport(host_key),
            }
            host = new_host(
                key_pair=host_key,
                listen_addrs=[Multiaddr("/ip4/127.0.0.1/tcp/0")],
                sec_opt=sec_opt,
            )

            async with host.run([Multiaddr("/ip4/127.0.0.1/tcp/0")]):
                event("✅", f"Host ready   id={host.get_id()}")

                # ── 3. Register inbound-ping handler ─────────────
                ping_log: list[float] = []
                ping_svc = LoggingPingService(host, ping_log)
                host.set_stream_handler(PING_PROTO, ping_svc.handle_ping)
                event("📋", f"Inbound ping handler registered on {PING_PROTO}")

                # ── 4. Connect ───────────────────────────────────
                event("🔗", "Connecting to Kubo …")
                t_conn = time.monotonic()
                peer_info = info_from_p2p_addr(Multiaddr(kubo_addr))
                with trio.fail_after(10):
                    await host.connect(peer_info)
                conn_ms = int((time.monotonic() - t_conn) * 1000)
                event("✅", f"Connected!   TCP + TLS/Noise + Yamux  ({conn_ms} ms)")

                # Introspect the negotiated security & muxer
                conns = host.get_network().connections.get(kubo_id, [])
                if conns:
                    conn = conns[0]
                    muxed = getattr(conn, "muxed_conn", None)
                    mux_name = type(muxed).__name__ if muxed else "unknown"
                    sec = (
                        getattr(conn, "security_protocol", None)
                        or getattr(getattr(conn, "secured_conn", None), "protocol_id", "unknown")
                    )
                    event("   ", f"             Security={sec}  Muxer={mux_name}")

                # ── 5. Identify ──────────────────────────────────
                event("🔍", f"Opening {IDENTIFY_PROTO} …")
                t_id = time.monotonic()
                try:
                    with trio.fail_after(10):
                        elapsed_id, id_bytes = await do_identify(host, kubo_id)
                    event("✅", f"Identify complete  ({elapsed_id*1000:.0f} ms, {id_bytes} bytes)")
                except Exception as e:
                    event("❌", f"Identify failed: {type(e).__name__}: {e}")

                # ── 6. Outbound ping (us → Kubo) ─────────────────
                event("📤", "Sending outbound ping (us → Kubo) …")
                try:
                    with trio.fail_after(10):
                        rtt_list = await ping_svc.ping(kubo_id, ping_amt=3)
                    for i, rtt in enumerate(rtt_list, 1):
                        event("✅", f"us → Kubo  ping {i}/3  RTT = {rtt} ms")
                except Exception as e:
                    event("❌", f"Outbound ping failed: {type(e).__name__}: {e}")

                # ── 7. Watch for re-pings for OBSERVE_MINUTES ────
                print()
                print(f"  Watching for Kubo inbound re-pings for {OBSERVE_MINUTES} min …")
                print(f"  (Kubo typically re-pings connected peers every ~10-15 min)")
                print()

                deadline = trio.current_time() + OBSERVE_MINUTES * 60
                ping_count_at_start = len(ping_log)

                while trio.current_time() < deadline:
                    remaining = int(deadline - trio.current_time())
                    new_pings = len(ping_log) - ping_count_at_start

                    # Check connection is still alive
                    conns = host.get_network().connections.get(kubo_id, [])
                    alive = any(not c.is_closed for c in conns)
                    conn_status = "connected" if alive else "DISCONNECTED"

                    print(
                        f"\r  [{ts()}]  Status: {conn_status} | "
                        f"Kubo re-pings received: {new_pings} | "
                        f"Time left: {remaining}s   ",
                        end="", flush=True,
                    )
                    await trio.sleep(2)

                print()  # newline after status line

                # ── Summary ──────────────────────────────────────
                print()
                print("=" * 62)
                print("  SUMMARY")
                print("=" * 62)
                print(f"  Connection time      : {conn_ms} ms")
                print(f"  Identify time        : {elapsed_id*1000:.0f} ms")
                print(f"  Outbound pings (us→Kubo): {len(rtt_list)} pings, "
                      f"RTTs={rtt_list} ms")
                new_pings = len(ping_log) - ping_count_at_start
                if new_pings:
                    gaps = [
                        f"{ping_log[i]-ping_log[i-1]:.1f}s"
                        for i in range(1, len(ping_log))
                    ]
                    print(f"  Inbound re-pings (Kubo→us): {new_pings} received")
                    if gaps:
                        print(f"  Re-ping intervals    : {', '.join(gaps)}")
                else:
                    print(f"  Inbound re-pings (Kubo→us): 0 in {OBSERVE_MINUTES} min")
                    print(f"  (Kubo re-ping interval > {OBSERVE_MINUTES} min in test profile)")

                conns = host.get_network().connections.get(kubo_id, [])
                alive = any(not c.is_closed for c in conns)
                print(f"  Connection at end    : {'✅ alive' if alive else '❌ dropped'}")
                print("=" * 62)

        finally:
            try:
                kubo_proc.terminate()
                kubo_proc.wait(timeout=5)
            except Exception:
                pass


if __name__ == "__main__":
    trio.run(main)
