"""PeerDrop Core — Engine orchestrator.

The PeerEngine is the single source of truth. It wraps py-ipfs-lite Peer
and exposes high-level methods that all clients (CLI, REST, GUI, MCP) use.
"""

from __future__ import annotations

import logging
import os
from dataclasses import dataclass, field
from typing import Any, Callable

import multiaddr

from peerdrop.core.discovery import DiscoveryManager
from peerdrop.core.events import EventBus
from peerdrop.core.messaging import MessagingManager
from peerdrop.core.models import Peer, Transfer
from peerdrop.core.transfer import TRANSFER_PROTOCOL, TransferManager

logger = logging.getLogger(__name__)


@dataclass
class EngineConfig:
    """Configuration for the PeerEngine."""

    listen_port: int = 4001
    blockstore_type: str = "memory"
    blockstore_path: str = "~/.peerdrop/blocks"
    device_name: str = ""
    listen_addrs: list[str] = field(default_factory=lambda: ["/ip4/0.0.0.0/tcp/0"])
    identity_seed: str | None = None
    download_dir: str = "~/.peerdrop/downloads"


class PeerEngine:
    """The PeerDrop engine — single source of truth.

    All clients talk to this. It wraps py-ipfs-lite Peer and adds:
    - Peer discovery (mDNS)
    - Transfer protocol (stream handshake + bitswap)
    - Event bus for client notifications
    """

    def __init__(self, config: EngineConfig | None = None) -> None:
        self._config = config or EngineConfig()
        self._ipfs_peer: Any = None
        self._discovery: DiscoveryManager | None = None
        self._transfers: TransferManager | None = None
        self._messaging: MessagingManager | None = None
        self._events = EventBus()
        self._started = False

    async def start(self) -> None:
        """Start the engine.

        Creates py-ipfs-lite Peer with mDNS enabled via Config.enable_mdns,
        registers the transfer protocol handler, and starts discovery.
        """
        if self._started:
            return

        from py_ipfs_lite.config import Config
        from py_ipfs_lite.peer import Peer as IPFSPeer

        ipfs_config = Config(
            blockstore_type=self._config.blockstore_type,
            blockstore_path=os.path.expanduser(self._config.blockstore_path) if self._config.blockstore_type == "filesystem" else None,
            offline=False,
            enable_mdns=True,
        )

        # Build deterministic host key from seed if provided
        host_key = None
        if self._config.identity_seed:
            import hashlib
            from libp2p.crypto.ed25519 import create_new_key_pair
            seed_bytes = hashlib.sha256(self._config.identity_seed.encode()).digest()
            host_key = create_new_key_pair(seed=seed_bytes)

        # Get actual network interface addresses for listening
        from libp2p.utils.address_validation import get_available_interfaces
        listen_addrs = get_available_interfaces(self._config.listen_port)
        logger.info(f"Listening on interfaces: {listen_addrs}")

        self._ipfs_peer = IPFSPeer(
            ipfs_config,
            host_key=host_key,
            listen_addrs=listen_addrs,
        )

        # Register discovery handler BEFORE starting the peer so we don't
        # miss any mDNS discoveries that happen during host startup.
        self._discovery = DiscoveryManager(self._ipfs_peer, self._events)
        await self._discovery.start()

        await self._ipfs_peer.start()

        # Register transfer protocol handler
        self._ipfs_peer.host.set_stream_handler(
            TRANSFER_PROTOCOL,
            self._transfers_handler,
        )

        # Initialize transfer manager after peer is started
        self._transfers = TransferManager(self._ipfs_peer, self._events, self._config.download_dir)

        # Initialize messaging manager (started later with nursery)
        self._messaging = MessagingManager(self._ipfs_peer, self._events)

        self._started = True
        peer_id = str(self._ipfs_peer.host.id())
        addrs = [str(a) for a in self._ipfs_peer.host.addrs()]
        logger.info(f"PeerEngine started: {peer_id}")
        logger.info(f"Listening on: {addrs}")

    async def stop(self) -> None:
        """Stop the engine and clean up resources."""
        if not self._started:
            return

        if self._messaging:
            await self._messaging.stop()

        # Unregister from the global peerDiscovery singleton so restarts in the
        # same process don't accumulate duplicate handlers.
        if self._discovery:
            self._discovery.stop()

        self._events.unsubscribe_all()

        if self._ipfs_peer:
            await self._ipfs_peer.close()
        self._started = False
        logger.info("PeerEngine stopped")

    async def _transfers_handler(self, stream: Any) -> None:
        """Internal stream handler for incoming transfers."""
        if self._transfers:
            await self._transfers.handle_incoming_transfer(stream)

    # --- High-level API ---

    def get_peer_id(self) -> str:
        """Return this device's peer ID."""
        if not self._ipfs_peer:
            raise RuntimeError("Engine not started")
        return str(self._ipfs_peer.host.id())

    def get_addrs(self) -> list[str]:
        """Return this device's listening addresses."""
        if not self._ipfs_peer:
            raise RuntimeError("Engine not started")
        return [str(a) for a in self._ipfs_peer.host.addrs()]

    def get_download_dir(self) -> str:
        """Return the current download directory."""
        if not self._transfers:
            return self._config.download_dir
        return self._transfers.download_dir

    def set_download_dir(self, new_dir: str) -> str:
        """Change the download directory at runtime. Returns resolved path."""
        if not self._transfers:
            self._config.download_dir = new_dir
            return os.path.expanduser(new_dir)
        return self._transfers.set_download_dir(new_dir)

    def discover_peers(self) -> list[Peer]:
        """Return discovered peers on the network."""
        if not self._discovery:
            return []
        return self._discovery.get_peers()

    def get_peer(self, peer_id: str) -> Peer | None:
        """Lookup a peer by ID."""
        if not self._discovery:
            return None
        return self._discovery.get_peer(peer_id)

    async def connect_peer(self, addr: str) -> Peer:
        """Manually connect to a peer by multiaddress."""
        if not self._ipfs_peer:
            raise RuntimeError("Engine not started")

        maddr = multiaddr.Multiaddr(addr)
        host = self._ipfs_peer.host

        # Extract peer_id from multiaddr if present
        peer_id = None
        for protocol in maddr.protocols():
            if protocol.name == "p2p":
                peer_id = maddr.value_for_protocol(protocol)
                break

        if not peer_id:
            raise ValueError("Address must include /p2p/<peer_id>")

        # Connect via host
        from libp2p.peer.id import ID
        from libp2p.peer.peerinfo import PeerInfo

        remote_id = ID.from_base58(peer_id)
        peer_info = PeerInfo(peer_id=remote_id, addrs=[maddr])

        # Add to peerstore first
        host.get_network().peerstore.add_addrs(remote_id, [maddr], 3600)

        # Connect
        await host.connect(peer_info)

        # Ensure the pubsub stream with this peer is registered. py-libp2p's
        # pubsub only tries once per 'connected' notifee, which can race the
        # muxer handshake (especially with mDNS auto-connect) and leave
        # messaging dead for this peer until reconnect. Re-registering here
        # makes chat reliable after an explicit connect.
        if self._messaging:
            try:
                await self._messaging.ensure_peer_stream(peer_id)
            except Exception as e:
                logger.warning(f"Failed to ensure pubsub stream with {peer_id}: {e}")

        # Register in discovery
        peer = Peer(
            peer_id=peer_id,
            addrs=[addr],
            name="",
        )
        if self._discovery:
            self._discovery._peers[peer_id] = peer

        logger.info(f"Connected to peer: {peer_id}")
        return peer

    async def ping(self, peer_id: str, count: int = 1) -> dict:
        """Ping a peer and return latency stats."""
        if not self._ipfs_peer:
            raise RuntimeError("Engine not started")

        from libp2p.host.ping import PingService
        from libp2p.peer.id import ID

        host = self._ipfs_peer.host
        remote_id = ID.from_base58(peer_id)
        ping_service = PingService(host)

        latencies = []
        for _ in range(count):
            import time
            import trio

            start = time.monotonic()
            try:
                with trio.fail_after(5.0):
                    await ping_service.ping(remote_id, ping_amt=1)
                elapsed = (time.monotonic() - start) * 1000
                latencies.append(round(elapsed, 2))
            except Exception as e:
                logger.warning(f"Ping to {peer_id} failed: {e}")
                latencies.append(-1)

        avg = sum(l for l in latencies if l > 0) / max(1, len([l for l in latencies if l > 0]))
        return {
            "peer_id": peer_id,
            "count": count,
            "latencies_ms": latencies,
            "avg_ms": round(avg, 2),
            "loss_pct": round(100 * sum(1 for l in latencies if l < 0) / count, 1),
        }

    async def send_file(self, file_path: str, target_peer_id: str) -> Transfer:
        """Send a file to a target peer."""
        if not self._transfers:
            raise RuntimeError("Engine not started")
        return await self._transfers.send_file(file_path, target_peer_id)

    def list_transfers(self) -> list[Transfer]:
        """Return all active and completed transfers."""
        if not self._transfers:
            return []
        return self._transfers.list_transfers()

    def get_transfer(self, transfer_id: str) -> Transfer | None:
        """Lookup a transfer by ID."""
        if not self._transfers:
            return None
        return self._transfers.get_transfer(transfer_id)

    def cancel_transfer(self, transfer_id: str) -> bool:
        """Cancel a transfer and abort in-flight work. Returns True if cancelled."""
        if not self._transfers:
            return False
        return self._transfers.cancel_transfer(transfer_id)

    def on_event(self, event_type: type, handler: Callable) -> Any:
        """Subscribe to an event type. Returns a receive channel."""
        return self._events.subscribe(event_type)

    @property
    def event_bus(self) -> EventBus:
        """Access the event bus directly."""
        return self._events

    # --- Messaging API ---

    @property
    def messaging(self) -> MessagingManager | None:
        """Access the messaging manager."""
        return self._messaging

    async def subscribe_topic(self, topic: str) -> None:
        """Subscribe to a gossipsub topic."""
        if not self._messaging:
            raise RuntimeError("Engine not started")
        await self._messaging.subscribe(topic)

    async def unsubscribe_topic(self, topic: str) -> None:
        """Unsubscribe from a topic."""
        if not self._messaging:
            raise RuntimeError("Engine not started")
        await self._messaging.unsubscribe(topic)

    async def publish_message(self, topic: str, message: str) -> None:
        """Publish a message to a topic."""
        if not self._messaging:
            raise RuntimeError("Engine not started")
        await self._messaging.publish(topic, message)

    def get_topics(self) -> list[str]:
        """Return list of subscribed topics."""
        if not self._messaging:
            return []
        return self._messaging.get_topics()

    def get_messages(self, topic: str | None = None, limit: int = 50) -> list[dict]:
        """Return recent messages."""
        if not self._messaging:
            return []
        return self._messaging.get_messages(topic, limit)
