"""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);
"""


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 作成。
    journal_mode=WAL で LIVE 中の SELECT (= 別 process) も許可。
    """

    # 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

    def __init__(self, path: pathlib.Path) -> None:
        self.path = path
        # WAL 化と schema 作成は init_db() で済んでいる前提 (= daemon 起動時に 1 度)。
        # 各 connection は busy_timeout (= 30s) と synchronous=NORMAL だけ設定。
        self.conn = sqlite3.connect(str(path), timeout=30.0)
        self.conn.execute("PRAGMA busy_timeout = 30000")
        self.conn.execute("PRAGMA synchronous = NORMAL")
        self.touched = 0
        self.rejected = 0
        self._pending = 0
        self._last_commit = time.monotonic()

    def write(self, ts: int, body: str) -> bool:
        if len(set(body)) <= 1:
            self.rejected += 1
            return False
        assert self.conn is not None
        # ON CONFLICT で atomic upsert。 並列 thread から同 body が同時に来てもレース無し。
        self.conn.execute(
            "INSERT INTO comments(b, first_ts, count) VALUES (?, ?, 1) "
            "ON CONFLICT(b) DO UPDATE SET count = count + 1",
            (body, ts),
        )
        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:
        if self._pending > 0:
            self.conn.commit()
            self._pending = 0
        self.conn.close()


# ---------- 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 = "",
) -> None:
    """LIVE chat collect。 stop_event が set されたら次の chunk で離脱。

    並列利用時: 各 thread が独自の pytchat session と SqliteAppender (= 別 connection)
    を持つ。 SQLite は WAL で concurrent write 可、 ただし内部で serialize される。
    """
    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"
    writer = SqliteAppender(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)
                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())
