"""Import 1 YouTube Live into stream-comments unified SQLite (anonymized + dedup).

Usage:
    python scripts/import_youtube_live.py <video_url_or_id> [max_messages] [duration_sec]

defaults: max_messages=5000、 duration_sec=1800 (= 30 min)
0 を渡すとその cap を無効化、 両方 0 で完全無制限 (= 配信終了 or Ctrl+C 待ち)。

LIVE 専用 (pytchat 経由)。 mask は src/mask.ts と等価ロジック (URL/MAIL/TEL)。

source-anonymized + dedup + single file:
- 出力先は `data/comments.sqlite` の 1 ファイル固定 (= 全 collect を統合、 配信跨ぎで dedup)
- comments table: b (TEXT PRIMARY KEY) + first_ts + count
  → 同一本文は count++ で集計、 record は重複しない
- 「同一文字のみ」 (= 「あああ」 「www」 「!!!」) は INSERT 前に reject
- 配信元 (channel / stream / video / 投稿者) を辿る情報は一切保存しない
"""

from __future__ import annotations

import pathlib
import re
import sqlite3
import sys
import threading
import time
import unicodedata

import pytchat

# ---------- mask (= src/mask.ts 等価) ----------

URL_RE = re.compile(r'(?:https?://|www\.)[^\s<>"\'　]+')
EMAIL_RE = re.compile(r'[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}')
PHONE_INTL_RE = re.compile(r'\+81[\- ]?\d{1,4}[\- ]?\d{1,4}[\- ]?\d{4}')
PHONE_HYPHEN_RE = re.compile(r'0\d{1,4}-\d{1,4}-\d{4}')
PHONE_PLAIN_RE = re.compile(r'(?<!\d)0\d{9,10}(?!\d)')
URL_TRAIL_RE = re.compile(r'[.,;:!?。、！？)\]）】]+$')


def mask(body: str) -> str:
    def _url_sub(m: re.Match[str]) -> str:
        trail = URL_TRAIL_RE.search(m.group())
        return "[URL]" + (trail.group() if trail else "")
    s = URL_RE.sub(_url_sub, body)
    s = EMAIL_RE.sub("[MAIL]", s)
    s = PHONE_INTL_RE.sub("[TEL]", s)
    s = PHONE_HYPHEN_RE.sub("[TEL]", s)
    s = PHONE_PLAIN_RE.sub("[TEL]", s)
    return s


WHITESPACE_RE = re.compile(r"\s+")
# YouTube channel custom emoji の textual form。 代表例:
#   :_cha: / :face-blue-smiling: / :hololive-pekora-peko-peko: 等
# alpha or underscore 始まり + 英数 / underscore / ハイフン の連続 + コロン終わり。
# 「19:30」 「これは:こう」 のような時刻 / 漢字始まりは intentionally 拾わない。
YT_STICKER_RE = re.compile(r":[_A-Za-z][_A-Za-z0-9\-]*:")


def strip_sticker(text: str) -> str:
    return YT_STICKER_RE.sub(" ", text)


def normalize(text: str) -> str:
    """NFKC + 改行 → 空白 + 連続空白 collapse + 両端 trim。
    dedup を実効化するため collect 時に適用する (= 「？？？」 と「???」 を同 row に集約)。
    """
    text = unicodedata.normalize("NFKC", text)
    text = text.replace("\r", " ").replace("\n", " ")
    return WHITESPACE_RE.sub(" ", text).strip()


# ---------- video id ----------

VIDEO_ID_RE = re.compile(r"(?:v=|youtu\.be/|/live/)([A-Za-z0-9_-]{11})")


def extract_video_id(url_or_id: str) -> str:
    m = VIDEO_ID_RE.search(url_or_id)
    if m:
        return m.group(1)
    if re.fullmatch(r"[A-Za-z0-9_-]{11}", url_or_id):
        return url_or_id
    raise ValueError(f"could not extract video id from {url_or_id!r}")


# ---------- SQLite writer (= 単一 file appender) ----------

SCHEMA = """
CREATE TABLE IF NOT EXISTS comments (
    b TEXT PRIMARY KEY,
    first_ts INTEGER NOT NULL,
    count INTEGER NOT NULL DEFAULT 1
);
CREATE INDEX IF NOT EXISTS idx_first_ts ON comments(first_ts);

-- 粗いジャンル別の出現回数 (★2026-09-15)。 チャンネル名・配信・投稿者は持たず、
-- channels.toml の genre (vtuber / news / sports 等の十数語) だけを comments の rowid に紐付ける。
-- genre 導入前に集めた行には対応行が無い (= 出どころ不明)。
CREATE TABLE IF NOT EXISTS genres (
    id INTEGER PRIMARY KEY,
    name TEXT NOT NULL UNIQUE
);
CREATE TABLE IF NOT EXISTS comment_genres (
    cid INTEGER NOT NULL,
    gid INTEGER NOT NULL,
    count INTEGER NOT NULL DEFAULT 1,
    PRIMARY KEY (cid, gid)
) WITHOUT ROWID;
"""

# channels.toml の genre に書ける語。 粗さを保つため (= 個別チャンネルを推定できないよう) ここで固定する。
GENRES = (
    "vtuber", "game", "news", "sports", "music", "business",
    "education", "entertainment", "product", "hobby",
    "sns", "board",  # 2026-09-21: 大量ソース (Bluesky / Fediverse = sns、 5ch = board)。 prune_db.py が間引く対象
)


def init_db(path: pathlib.Path) -> None:
    """daemon 起動時に 1 度呼ぶ: journal_mode 設定 + schema 作成。

    journal_mode は DELETE (= classic rollback journal)。 WAL は NTFS-FUSE (ntfs-3g) 等で
    mmap が不安定 → 'disk I/O error' を引き起こすので避ける。 trade-off: reader 並列性
    (= LIVE 中の同時 SELECT) は失われる、 ただし read は集計時の short query のみ想定。
    """
    path.parent.mkdir(parents=True, exist_ok=True)
    conn = sqlite3.connect(str(path), timeout=30.0)
    try:
        conn.execute("PRAGMA journal_mode = DELETE")
        conn.execute("PRAGMA synchronous = NORMAL")
        conn.executescript(SCHEMA)
        conn.commit()
    finally:
        conn.close()


class SqliteAppender:
    """data/comments.sqlite 単一 file に append する writer。

    既存 DB に対しては INSERT OR UPDATE で merge、 fresh なら schema 作成。

    ## 並列書き込みの扱い (★2026-09-11)

    DB は ntfs-3g (fuseblk) 上にあり **WAL が使えない** (mmap 不安定 → disk I/O error)
    ため journal_mode = DELETE。 その状態で channel thread がそれぞれ connection を
    持つと、 各 connection が最長 `FLUSH_SECS` 秒 transaction を開けたままになり、
    thread 数が増えるほど `database is locked` が出る (43ch に増やした時点で発生)。

    そこで **process 内で connection を 1 本だけ共有** し、 全 thread が
    `_SHARED_LOCK` 越しに書く。 SQLite の file lock 競合そのものが起きなくなり、
    待ちは in-process の短い lock だけになる。 `shared()` で取得し、 `close()` は
    参照カウントが 0 になった時だけ実際に閉じる (= 他 thread の書き込みを壊さない)。
    """

    # commit を batch 化して fsync 頻度を下げる (= ntfs-3g/FUSE 上で per-write fsync が
    # mount.ntfs を飽和させ load を上げていたため)。 `BATCH` 件 か `FLUSH_SECS` 秒の
    # 早い方で commit。 transaction は最長 FLUSH_SECS 秒しか開かないので、 並列 thread の
    # lock 競合は busy_timeout (= 30s) で吸収される。 graceful stop (close) で必ず flush、
    # 停電も UPS→SIGTERM 経由で flush されるため最悪でも直近 batch 分のみ損失。
    BATCH = 50
    FLUSH_SECS = 3.0

    # process 内で共有する単一 connection (= 書き込みの唯一の入口)。
    _shared: "SqliteAppender | None" = None
    _shared_refs = 0
    _shared_guard = threading.Lock()

    @classmethod
    def shared(cls, path: pathlib.Path) -> "SqliteAppender":
        """共有 writer を取得する (無ければ作る)。 呼んだ数だけ `close()` すること。"""
        with cls._shared_guard:
            if cls._shared is None:
                cls._shared = cls(path)
            cls._shared_refs += 1
            return cls._shared

    def __init__(self, path: pathlib.Path) -> None:
        self.path = path
        # journal_mode / schema は init_db() で済んでいる前提 (= daemon 起動時に 1 度)。
        # check_same_thread=False: 複数 channel thread から同 connection を使うため
        # (排他は self._lock で保証する)。
        self.conn = sqlite3.connect(str(path), timeout=30.0, check_same_thread=False)
        self.conn.execute("PRAGMA busy_timeout = 30000")
        self.conn.execute("PRAGMA synchronous = NORMAL")
        self._lock = threading.Lock()
        self.touched = 0
        self.rejected = 0
        self._pending = 0
        self._last_commit = time.monotonic()
        self._gids: dict[str, int] = {}

    def _gid(self, genre: str) -> int:
        """genre 名 → genres.id (無ければ作る)。 self._lock の中で呼ぶこと。"""
        gid = self._gids.get(genre)
        if gid is None:
            self.conn.execute("INSERT OR IGNORE INTO genres(name) VALUES (?)", (genre,))
            gid = self.conn.execute("SELECT id FROM genres WHERE name = ?", (genre,)).fetchone()[0]
            self._gids[genre] = gid
        return gid

    def write(self, ts: int, body: str, genre: str | None = None) -> bool:
        if len(set(body)) <= 1:
            self.rejected += 1
            return False
        assert self.conn is not None
        # ON CONFLICT で atomic upsert。 並列 thread から同 body が同時に来てもレース無し。
        # connection は process 内で共有なので、 execute + commit を lock で括る。
        with self._lock:
            self.conn.execute(
                "INSERT INTO comments(b, first_ts, count) VALUES (?, ?, 1) "
                "ON CONFLICT(b) DO UPDATE SET count = count + 1",
                (body, ts),
            )
            if genre:
                cid = self.conn.execute("SELECT rowid FROM comments WHERE b = ?", (body,)).fetchone()[0]
                self.conn.execute(
                    "INSERT INTO comment_genres(cid, gid, count) VALUES (?, ?, 1) "
                    "ON CONFLICT(cid, gid) DO UPDATE SET count = count + 1",
                    (cid, self._gid(genre)),
                )
            self.touched += 1
            self._pending += 1
            if (
                self._pending >= self.BATCH
                or (time.monotonic() - self._last_commit) >= self.FLUSH_SECS
            ):
                self.conn.commit()
                self._pending = 0
                self._last_commit = time.monotonic()
        return True

    def close(self) -> None:
        """flush する。 共有 writer の場合、 最後の参照が閉じた時だけ connection を閉じる。"""
        cls = type(self)
        with cls._shared_guard:
            is_shared = cls._shared is self
            if is_shared:
                cls._shared_refs -= 1
            last = (not is_shared) or cls._shared_refs <= 0
            with self._lock:
                if self._pending > 0:
                    self.conn.commit()
                    self._pending = 0
                if last:
                    self.conn.close()
            if is_shared and last:
                cls._shared = None
                cls._shared_refs = 0


# ---------- main flow ----------

def collect(
    url_or_id: str,
    max_messages: int,
    duration_sec: int,
    out_dir: pathlib.Path,
    stop_event: threading.Event | None = None,
    log_tag: str = "",
    genre: str | None = None,
) -> None:
    """LIVE chat collect。 stop_event が set されたら次の chunk で離脱。

    並列利用時: 各 thread が独自の pytchat session を持つが、 **SQLite writer は
    process 内で 1 connection を共有** する (`SqliteAppender.shared`)。 ntfs-3g 上では
    WAL が使えず、 connection を分けると thread 数に比例して
    `database is locked` が出るため (★2026-09-11)。
    """
    video_id = extract_video_id(url_or_id)
    # interruptable=False: pytchat が SIGINT hook を立てるのを抑止 (= 非 main thread でも動かす)。
    # 停止は stop_event 経由で daemon 側から制御。
    chat = pytchat.create(video_id=video_id, interruptable=False)

    db_path = out_dir / "comments.sqlite"
    # 書き込みは process 内で 1 connection に集約する (= file lock 競合を作らない)。
    writer = SqliteAppender.shared(db_path)

    count = 0
    unlimited_msgs = max_messages <= 0
    unlimited_dur = duration_sec <= 0
    deadline = None if unlimited_dur else time.time() + duration_sec
    stop_reason = "live_ended"

    prefix = f"[{log_tag}] " if log_tag else ""
    print(
        f"{prefix}[start] db={db_path.name} "
        f"max={'∞' if unlimited_msgs else max_messages} "
        f"dur={'∞' if unlimited_dur else f'{duration_sec}s'} "
        f"is_alive={chat.is_alive()}"
    )
    try:
        while chat.is_alive() and (unlimited_msgs or count < max_messages):
            if stop_event is not None and stop_event.is_set():
                stop_reason = "stopped"
                break
            if deadline is not None and time.time() >= deadline:
                stop_reason = "duration_cap"
                break
            items = list(chat.get().sync_items())
            if not items:
                # stop 応答性のため細かく sleep
                for _ in range(2):
                    if stop_event is not None and stop_event.is_set():
                        break
                    time.sleep(1)
                continue
            for c in items:
                raw = (c.message or "").strip()
                if not raw:
                    continue
                # sticker (= :_xxx: 形式) は本文 noise として除去、 mixed comment の本文部分は残す。
                body = normalize(strip_sticker(mask(raw)))
                if not body:
                    continue
                writer.write(int(c.timestamp), body, genre)
                count += 1
                if count % 100 == 0:
                    print(
                        f"{prefix}  seen={count} touched={writer.touched} "
                        f"rejected={writer.rejected}"
                    )
                if not unlimited_msgs and count >= max_messages:
                    stop_reason = "message_cap"
                    break
                if stop_event is not None and stop_event.is_set():
                    stop_reason = "stopped"
                    break
            else:
                continue
            break
    except KeyboardInterrupt:
        stop_reason = "interrupted"
        print(f"\n{prefix}[interrupted] flushing...")
    finally:
        chat.terminate()
        writer.close()

    size_kib = db_path.stat().st_size / 1024 if db_path.exists() else 0
    print(
        f"{prefix}[done] reason={stop_reason} seen={count} "
        f"touched={writer.touched} rejected={writer.rejected} "
        f"db_size={size_kib:.1f}KiB"
    )


def main() -> int:
    if len(sys.argv) < 2:
        print(__doc__)
        return 2
    url = sys.argv[1]
    max_messages = int(sys.argv[2]) if len(sys.argv) > 2 else 5000
    duration_sec = int(sys.argv[3]) if len(sys.argv) > 3 else 1800

    here = pathlib.Path(__file__).resolve().parent.parent
    out_dir = here / "data"
    collect(url, max_messages, duration_sec, out_dir)
    return 0


if __name__ == "__main__":
    raise SystemExit(main())
