"""PeerDrop Core — Event definitions and bus."""

from dataclasses import dataclass, field
from datetime import datetime
from typing import Any, Callable

import trio


@dataclass
class PeerDiscovered:
    """Emitted when a new peer is found via mDNS."""

    peer_id: str
    addrs: list[str]
    timestamp: datetime = field(default_factory=datetime.now)


@dataclass
class PeerLost:
    """Emitted when a peer goes offline."""

    peer_id: str
    timestamp: datetime = field(default_factory=datetime.now)


@dataclass
class TransferStarted:
    """Emitted when a transfer begins."""

    transfer_id: str
    file_name: str
    sender_peer_id: str
    receiver_peer_id: str
    timestamp: datetime = field(default_factory=datetime.now)


@dataclass
class TransferProgress:
    """Emitted during transfer to report progress."""

    transfer_id: str
    progress: float
    bytes_sent: int
    timestamp: datetime = field(default_factory=datetime.now)


@dataclass
class TransferCompleted:
    """Emitted when a transfer finishes successfully."""

    transfer_id: str
    file_name: str
    timestamp: datetime = field(default_factory=datetime.now)


@dataclass
class TransferFailed:
    """Emitted when a transfer fails."""

    transfer_id: str
    error: str
    timestamp: datetime = field(default_factory=datetime.now)


@dataclass
class MessageReceived:
    """Emitted when a pubsub message is received on a topic."""

    topic: str
    sender: str
    data: str
    timestamp: datetime = field(default_factory=datetime.now)


class EventBus:
    """Async event bus using trio memory channels.

    Usage:
        bus = EventBus()
        rx = bus.subscribe(TransferCompleted)
        # In a trio task:
        async with rx:
            async for event in rx:
                print(event)
        # Publishing:
        bus.publish(TransferCompleted(transfer_id="abc", file_name="test.txt"))
    """

    def __init__(self) -> None:
        self._subscribers: dict[type, list[trio.MemorySendChannel]] = {}

    def subscribe(self, event_type: type) -> trio.MemoryReceiveChannel:
        """Subscribe to an event type. Returns a receive channel."""
        send_channel, receive_channel = trio.open_memory_channel[Any](128)
        self._subscribers.setdefault(event_type, []).append(send_channel)
        return receive_channel

    def publish(self, event: Any) -> None:
        """Publish an event to all subscribers of its type."""
        event_type = type(event)
        for channel in self._subscribers.get(event_type, []):
            try:
                channel.send_nowait(event)
            except trio.WouldBlock:
                # Drop event if channel is full (best-effort delivery).
                # Progress events are high-volume; terminal events are also
                # recovered by clients via polling.
                pass
            except (trio.ClosedResourceError, trio.BrokenResourceError):
                # Subscriber closed its channel — drop it from the registry.
                self._subscribers[event_type].remove(channel)

    def unsubscribe_all(self) -> None:
        """Close all subscriber channels."""
        for channels in self._subscribers.values():
            for ch in channels:
                ch.close()
        self._subscribers.clear()
