"""Bluesky 収集 (Jetstream = 公式の公開 WebSocket、 認証不要、 2026-09-21)。

    python scripts/bluesky.py --dry --duration 30          # 動作確認 (表示だけ)
    python scripts/bluesky.py --genre entertainment         # 常駐 (daemon からは platform = "bluesky")

全投稿 (app.bsky.feed.post の create) が流れてくるので、 `langs` に ja を含む投稿だけ取る。
本文は YouTube と同じ経路 (mask → NFKC → SqliteAppender)。 投稿者 DID / URI / 時刻の細部は保存しない。
切断されたら数秒後に再接続 (cursor は持たない = 取りこぼしは許容、 dedup は DB 側)。
"""

from __future__ import annotations

import asyncio
import json
import pathlib
import re
import sys
import threading
import time

import websockets

from import_youtube_live import SqliteAppender, mask, normalize

ROOT = pathlib.Path(__file__).resolve().parent.parent
# 公式の Jetstream インスタンス (複数あるので順に試す)
ENDPOINTS = [
    "wss://jetstream2.us-east.bsky.network/subscribe",
    "wss://jetstream1.us-east.bsky.network/subscribe",
    "wss://jetstream2.us-west.bsky.network/subscribe",
    "wss://jetstream1.us-west.bsky.network/subscribe",
]
QUERY = "?wantedCollections=app.bsky.feed.post"
MENTION_RE = re.compile(r"@[\w.\-]+")


def clean(text: str) -> str:
    text = MENTION_RE.sub("", text)  # メンションは投稿者情報なので落とす
    return normalize(mask(text))


async def _run(genre, out_dir, stop_event, dry, duration, log_tag):
    prefix = f"[{log_tag}] " if log_tag else "[bluesky] "
    writer = None if dry else SqliteAppender.shared(out_dir / "comments.sqlite")
    deadline = time.time() + duration if duration > 0 else None
    count = 0
    ep = 0
    print(f"{prefix}[start] dry={dry}", flush=True)
    try:
        while True:
            if stop_event is not None and stop_event.is_set():
                break
            if deadline is not None and time.time() >= deadline:
                break
            url = ENDPOINTS[ep % len(ENDPOINTS)] + QUERY
            try:
                async with websockets.connect(url, max_size=2**20, open_timeout=30) as ws:
                    while True:
                        if stop_event is not None and stop_event.is_set():
                            return
                        if deadline is not None and time.time() >= deadline:
                            return
                        try:
                            raw = await asyncio.wait_for(ws.recv(), timeout=30)
                        except asyncio.TimeoutError:
                            continue
                        try:
                            ev = json.loads(raw)
                        except ValueError:
                            continue
                        c = ev.get("commit") or {}
                        if c.get("operation") != "create" or c.get("collection") != "app.bsky.feed.post":
                            continue
                        rec = c.get("record") or {}
                        if "ja" not in (rec.get("langs") or []):
                            continue
                        if rec.get("reply") and dry:
                            pass
                        text = rec.get("text") or ""
                        for line in text.split("\n"):
                            body = clean(line)
                            if len(body) < 2:
                                continue
                            if dry:
                                if count < 30:
                                    print("  ", body[:80], flush=True)
                            else:
                                writer.write(int(time.time() * 1000), body, genre)
                            count += 1
                            if count % 500 == 0 and not dry:
                                print(f"{prefix}  seen={count} touched={writer.touched} rejected={writer.rejected}", flush=True)
            except Exception as e:
                ep += 1
                print(f"{prefix}[ws] {e!r} → reconnect in 5s ({ENDPOINTS[ep % len(ENDPOINTS)]})", flush=True)
                for _ in range(5):
                    if stop_event is not None and stop_event.is_set():
                        return
                    await asyncio.sleep(1)
    finally:
        if writer is not None:
            writer.close()
        print(f"{prefix}[done] seen={count}", flush=True)


def collect_bluesky(out_dir: pathlib.Path, genre: str | None = None, stop_event: threading.Event | None = None,
                    dry: bool = False, duration_sec: int = 0, log_tag: str = "") -> None:
    asyncio.run(_run(genre, out_dir, stop_event, dry, duration_sec, log_tag))


def main() -> int:
    import argparse

    ap = argparse.ArgumentParser()
    ap.add_argument("--dry", action="store_true")
    ap.add_argument("--duration", type=int, default=0)
    ap.add_argument("--genre", default=None)
    ap.add_argument("--out-dir", default=str(ROOT / "data"))
    a = ap.parse_args()
    collect_bluesky(pathlib.Path(a.out_dir), a.genre, None, a.dry, a.duration)
    return 0


if __name__ == "__main__":
    sys.exit(main())
