"""channels.toml を polling して YouTube LIVE を検知 → 自動 collect する daemon (並列版)。

Usage:
    python scripts/daemon.py                    # default: ./channels.toml
    python scripts/daemon.py path/to/channels.toml

各 channel を 1 thread に割り当てて並列実行。 thread は poll → LIVE 検知 → collect →
配信終了 → 次 poll、 のサイクルで自走。 SQLite は WAL モードで concurrent write 可。

LIVE 検知は `https://www.youtube.com/channel/<id>/live` または `/@<handle>/live` を
fetch して HTML から判定 (= API key 不要、 ただし HTML 構造変更に弱い)。

twitch は `scripts/twitch.py` (IRC 匿名 JOIN)。 .env に TWITCH_CLIENT_ID/SECRET があれば Helix で LIVE 検知して
LIVE 中だけ JOIN、 無ければ常時 JOIN (オフライン中のチャットもそのまま集める)。 niconico は未対応。
"""

from __future__ import annotations

import pathlib
import random
import re
import signal
import sys
import threading
import time
import tomllib
import urllib.error
import urllib.request
from typing import Any

from import_youtube_live import GENRES, collect, init_db
from twitch import check_twitch_live, collect_twitch, helix_configured
from import_5ch import collect as collect_5ch
from bluesky import collect_bluesky
from fediverse import collect_fediverse

# inotify は NTFS-FUSE (= ntfs-3g) 上では kernel event が発生しないため使えない。
# SMB 共有 dir の実体が FUSE mount に乗っているので、 一律 mtime polling で済ます
# (1s 間隔の stat() は micro-overhead で十分軽い)。
HAS_INOTIFY = False

UA = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) Gecko/20100101 Firefox/120.0"

# YouTube HTML の player config に出る LIVE 中シグナル。
RE_IS_LIVE = re.compile(r'"isLiveNow"\s*:\s*true')
RE_CANONICAL = re.compile(
    r'<link rel="canonical" href="https://www\.youtube\.com/watch\?v=([A-Za-z0-9_-]{11})"'
)


def live_url_for(id_or_handle: str) -> str:
    if id_or_handle.startswith("@"):
        return f"https://www.youtube.com/{id_or_handle}/live"
    if id_or_handle.startswith("UC") and len(id_or_handle) == 24:
        return f"https://www.youtube.com/channel/{id_or_handle}/live"
    return f"https://www.youtube.com/@{id_or_handle}/live"


def check_youtube_live(channel_id: str) -> str | None:
    url = live_url_for(channel_id)
    try:
        req = urllib.request.Request(url, headers={"User-Agent": UA})
        with urllib.request.urlopen(req, timeout=10) as r:
            html = r.read().decode("utf-8", errors="replace")
    except urllib.error.URLError as e:
        print(f"  [check err] {channel_id}: {e}")
        return None
    if not RE_IS_LIVE.search(html):
        return None
    m = RE_CANONICAL.search(html)
    return m.group(1) if m else None


class ChannelWatcher(threading.Thread):
    """1 channel を専属で監視する thread。 poll → LIVE → collect → 次 poll を自走。"""

    def __init__(
        self,
        channel: dict[str, Any],
        out_dir: pathlib.Path,
        interval_sec: int,
        stop_event: threading.Event,
        start_delay: float = 0.0,
    ) -> None:
        name = channel.get("name") or channel.get("id") or "?"
        super().__init__(daemon=True, name=name)
        self.channel = channel
        self.out_dir = out_dir
        self.interval_sec = interval_sec
        self.stop_event = stop_event
        # 起動 stagger: pytchat の同時 init で Windows 10035 socket error が出るため、
        # thread ごとに少しずつずらす。
        self.start_delay = start_delay

    def _wait(self, seconds: int) -> None:
        """interruptible sleep。 ±20% jitter で複数 thread の burst を分散。"""
        jittered = max(1, int(seconds * random.uniform(0.8, 1.2)))
        for _ in range(jittered):
            if self.stop_event.is_set():
                return
            time.sleep(1)

    def run(self) -> None:
        platform = self.channel.get("platform")
        ch_id = self.channel.get("id")
        if not ch_id:
            return
        if platform == "twitch":
            self._run_twitch(ch_id)
            return
        if platform == "5ch":
            self._run_5ch(ch_id)
            return
        if platform == "bluesky":
            self._run_bluesky()
            return
        if platform == "fediverse":
            self._run_fediverse(ch_id)
            return
        if platform != "youtube":
            print(f"[{self.name}] skip: platform {platform!r} not supported yet")
            return
        if self.start_delay > 0:
            self.stop_event.wait(timeout=self.start_delay)
        while not self.stop_event.is_set():
            # while 全体を try で包んで、 想定外例外でも thread が永久死亡しないように。
            try:
                vid = check_youtube_live(ch_id)
                if vid:
                    print(f"[{self.name}] LIVE → video {vid}")
                    collect(
                        vid,
                        max_messages=0,
                        duration_sec=0,
                        out_dir=self.out_dir,
                        stop_event=self.stop_event,
                        log_tag=self.name,
                        genre=self.channel.get("genre"),
                    )
            except Exception as e:
                print(f"[{self.name}] tick err: {e!r}")
            if self.stop_event.is_set():
                return
            self._wait(self.interval_sec)


    def _run_fediverse(self, base: str) -> None:
        """公開ローカル TL を polling (interval_sec、 既定 60 秒)。 id = インスタンスの URL。"""
        if self.start_delay > 0:
            self.stop_event.wait(timeout=self.start_delay)
        while not self.stop_event.is_set():
            try:
                collect_fediverse(base.rstrip("/"), self.out_dir, self.channel.get("genre"), self.stop_event,
                                  log_tag=self.name, interval=int(self.channel.get("interval_sec", 60)))
            except Exception as e:
                print(f"[{self.name}] tick err: {e!r}")
            if self.stop_event.is_set():
                return
            self._wait(30)

    def _run_bluesky(self) -> None:
        """Jetstream に常駐 (切断時は collect 側が再接続)。 例外で抜けたら 30 秒後に再開。"""
        if self.start_delay > 0:
            self.stop_event.wait(timeout=self.start_delay)
        while not self.stop_event.is_set():
            try:
                collect_bluesky(self.out_dir, self.channel.get("genre"), self.stop_event, log_tag=self.name)
            except Exception as e:
                print(f"[{self.name}] tick err: {e!r}")
            if self.stop_event.is_set():
                return
            self._wait(30)

    def _run_5ch(self, match: str) -> None:
        """定期 (interval_hours、 既定 6 時間) に板を発見し直して勢い上位スレを取る。 id = bbsmenu の板名/カテゴリに含まれる語。"""
        if self.start_delay > 0:
            self.stop_event.wait(timeout=self.start_delay)
        hours = float(self.channel.get("interval_hours", 6))
        while not self.stop_event.is_set():
            try:
                collect_5ch(
                    match,
                    self.channel.get("genre"),
                    self.out_dir,
                    limit_boards=int(self.channel.get("boards", 2)),
                    threads_per_board=int(self.channel.get("threads", 10)),
                    log_tag=self.name,
                    stop_event=self.stop_event,
                )
            except Exception as e:
                print(f"[{self.name}] tick err: {e!r}")
            if self.stop_event.is_set():
                return
            self._wait(int(hours * 3600))

    def _run_twitch(self, login: str) -> None:
        """Helix あり: poll → LIVE なら JOIN (配信終了で離脱)。 Helix なし: 常時 JOIN (切断時は collect 側が再接続)。"""
        if self.start_delay > 0:
            self.stop_event.wait(timeout=self.start_delay)
        use_helix = helix_configured()
        while not self.stop_event.is_set():
            try:
                if use_helix:
                    if check_twitch_live(login):
                        print(f"[{self.name}] LIVE (twitch #{login})")
                        collect_twitch(
                            login,
                            self.out_dir,
                            stop_event=self.stop_event,
                            log_tag=self.name,
                            genre=self.channel.get("genre"),
                            still_live=lambda: check_twitch_live(login),
                        )
                else:
                    collect_twitch(
                        login,
                        self.out_dir,
                        stop_event=self.stop_event,
                        log_tag=self.name,
                        genre=self.channel.get("genre"),
                    )
            except Exception as e:
                print(f"[{self.name}] tick err: {e!r}")
            if self.stop_event.is_set():
                return
            self._wait(self.interval_sec)


class Daemon:
    def __init__(self, config_path: pathlib.Path) -> None:
        self.config_path = config_path.resolve()
        with self.config_path.open("rb") as f:
            config: dict[str, Any] = tomllib.load(f)
        self.interval = int(config.get("poll", {}).get("interval_sec", 60))
        self.channels: list[dict[str, Any]] = config.get("channel", [])
        # genre は粗い固定語彙だけ許す (= 個別チャンネルを推定できる細かい値を DB に入れない)。
        # 語彙外 / 未指定の channel は genre 無しで集める。
        for ch in self.channels:
            g = ch.get("genre")
            if g is not None and g not in GENRES:
                print(f"[config] {ch.get('name')}: unknown genre {g!r} (ignored; allowed: {', '.join(GENRES)})")
                ch["genre"] = None
            elif g is None:
                print(f"[config] {ch.get('name')}: no genre")
        self.out_dir = pathlib.Path(__file__).resolve().parent.parent / "data"
        self.stop_event = threading.Event()
        # hot reload 用: 起動時の mtime を記憶、 変わったら exit → wrapper が再起動
        self.config_mtime = self.config_path.stat().st_mtime

    def request_stop(self, *args: Any) -> None:
        if not self.stop_event.is_set():
            print("\n[shutdown] requesting all watchers to stop...")
            self.stop_event.set()

    def _trigger_reload(self) -> None:
        print(f"[reload] {self.config_path.name} changed, exiting for restart")
        self.stop_event.set()

    def _wait_polling(self) -> None:
        """mtime polling (= Linux 以外 / inotify_simple 不在時の fallback)。 1s 間隔。"""
        while not self.stop_event.is_set():
            try:
                cur = self.config_path.stat().st_mtime
            except OSError:
                cur = self.config_mtime
            if cur != self.config_mtime:
                self._trigger_reload()
                return
            self.stop_event.wait(timeout=1.0)

    def _wait_inotify(self) -> None:
        """Linux inotify で push event を受ける。 polling 不要、 反映 ~10ms。

        watch は config_path の親 dir に張る (= file 上書き / atomic save に対応)。
        対象 event: MODIFY (= 内容書き換え) / CLOSE_WRITE (= save) / MOVED_TO (= editor の atomic 置換)。
        """
        ino = INotify()
        watch_mask = (
            inotify_flags.MODIFY
            | inotify_flags.CLOSE_WRITE
            | inotify_flags.MOVED_TO
            | inotify_flags.CREATE
        )
        # symlink を辿った先の親 dir を watch (= /var/www/.../stream-comments/)
        real_path = self.config_path.resolve()
        ino.add_watch(str(real_path.parent), watch_mask)
        target_name = real_path.name
        while not self.stop_event.is_set():
            for event in ino.read(timeout=1000):
                if event.name == target_name:
                    self._trigger_reload()
                    return

    def run(self) -> None:
        signal.signal(signal.SIGINT, self.request_stop)
        if hasattr(signal, "SIGTERM"):
            signal.signal(signal.SIGTERM, self.request_stop)
        # SQLite 初期化 (WAL 化 + schema)。 各 watcher thread の SqliteAppender は
        # connection を作るだけで、 ここでの初期化に依存する。
        init_db(self.out_dir / "comments.sqlite")
        print(
            f"[start] channels={len(self.channels)} "
            f"interval={self.interval}s out={self.out_dir}/comments.sqlite"
        )
        watchers = [
            ChannelWatcher(
                ch,
                self.out_dir,
                self.interval,
                self.stop_event,
                start_delay=i * 1.5,
            )
            for i, ch in enumerate(self.channels)
        ]
        for w in watchers:
            w.start()
        # main thread は config 変更を待つ。 Linux なら inotify (push)、 他は mtime polling。
        try:
            if HAS_INOTIFY:
                self._wait_inotify()
            else:
                self._wait_polling()
        finally:
            self.stop_event.set()
            for w in watchers:
                w.join(timeout=30)
                if w.is_alive():
                    print(f"[shutdown] {w.name} did not stop within 30s (forcing exit)")
        print("[shutdown] done")


def main() -> int:
    if len(sys.argv) > 1:
        path = pathlib.Path(sys.argv[1])
    else:
        path = pathlib.Path(__file__).resolve().parent.parent / "channels.toml"
    if not path.exists():
        print(f"[FAIL] channels.toml not found: {path}", file=sys.stderr)
        print("  hint: copy channels.toml.example to channels.toml and edit", file=sys.stderr)
        return 1
    Daemon(path).run()
    return 0


if __name__ == "__main__":
    raise SystemExit(main())
