#!/usr/bin/env python3
"""Media library watcher: watches /mnt/DataDisk/待处理 for new files and
fires the Hermes media-library-watch webhook (HMAC-SHA256 V2 signature).

Durable daemon: run via systemd user service (see media-watcher.service),
survives reboots. Debounces bursts so a multi-file download = one webhook call.
"""
import hashlib
import hmac
import json
import logging
import os
import sys
import threading
import time
import urllib.request

from watchdog.events import FileSystemEventHandler
from watchdog.observers import Observer

WATCH_DIR = "/mnt/DataDisk/待处理"
SUBS_FILE = os.path.expanduser("~/.hermes/webhook_subscriptions.json")
ROUTE = "media-library-watch"
DEBOUNCE_SECONDS = 5
# Only media files that matter for the two organize rules
RELEVANT_SUFFIXES = (".mp4", ".funscript")
# Download-in-progress markers: ignore these events entirely
IGNORE_SUFFIXES = (".aria2", ".part", ".tmp", ".crdownload")

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(levelname)s %(message)s",
    handlers=[
        logging.FileHandler(os.path.expanduser("~/.hermes/logs/media-watcher.log")),
        logging.StreamHandler(),
    ],
)
log = logging.getLogger("media-watcher")


def load_route():
    """Return (url, secret) for the subscription route."""
    with open(SUBS_FILE) as f:
        subs = json.load(f)
    route = subs[ROUTE]
    secret = route["secret"]
    if not secret or secret == "INSECURE_NO_AUTH":
        raise RuntimeError(f"route {ROUTE} has no usable secret")
    return f"http://localhost:8644/webhooks/{ROUTE}", secret


def fire_webhook(url: str, secret: str, payload: dict) -> bool:
    """POST payload with V2 signature; returns True on 2xx."""
    body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
    ts = str(int(time.time()))
    signed = hmac.new(
        secret.encode(), ts.encode() + b"." + body, hashlib.sha256
    ).hexdigest()
    req = urllib.request.Request(
        url,
        data=body,
        method="POST",
        headers={
            "Content-Type": "application/json",
            "X-Webhook-Signature-V2": signed,
            "X-Webhook-Timestamp": ts,
        },
    )
    try:
        with urllib.request.urlopen(req, timeout=15) as resp:
            ok = 200 <= resp.status < 300
            log.info("webhook POST -> %s (%s)", resp.status, payload.get("file"))
            return ok
    except Exception as e:  # network/gateway down: log, don't crash
        log.error("webhook POST failed: %s", e)
        return False


class DebouncedFirer:
    """Coalesces rapid events into a single webhook call after quiet time."""

    def __init__(self, url: str, secret: str):
        self.url = url
        self.secret = secret
        self._lock = threading.Lock()
        self._timer = None
        self._pending = {}

    def _schedule(self):
        with self._lock:
            if self._timer:
                self._timer.cancel()
            self._timer = threading.Timer(DEBOUNCE_SECONDS, self._fire)
            self._timer.daemon = True
            self._timer.start()

    def notify(self, file_path: str, event: str):
        with self._lock:
            self._pending = {"file": file_path, "event": event, "ts": int(time.time())}
        self._schedule()
        log.info("queued %s (%s)", file_path, event)

    def _fire(self):
        with self._lock:
            payload, self._pending = self._pending, {}
        if not payload:
            return
        # Fire at most ~3 times; if the gateway is down, drop this event
        # (next file landing will re-trigger anyway).
        for _ in range(3):
            if fire_webhook(self.url, self.secret, payload):
                return
            time.sleep(2)
        log.warning("gave up after retries for %s", payload.get("file"))


class MediaHandler(FileSystemEventHandler):
    def __init__(self, firer: DebouncedFirer):
        self.firer = firer

    def _consider(self, path, event):
        name = os.path.basename(path)
        lower = name.lower()
        if lower.endswith(IGNORE_SUFFIXES):
            return
        if not lower.endswith(RELEVANT_SUFFIXES):
            return
        self.firer.notify(path, event)

    def on_created(self, ev):
        self._consider(ev.src_path, "created")

    def on_moved(self, ev):
        # aria2 renames xxx.mp4.aria2 -> xxx.mp4 on completion: that is the
        # "download finished" signal we actually want to fire on.
        self._consider(ev.dest_path, "moved_in")

    def on_modified(self, ev):
        self._consider(ev.src_path, "modified")


def main():
    url, secret = load_route()
    firer = DebouncedFirer(url, secret)
    handler = MediaHandler(firer)
    observer = Observer()
    observer.schedule(handler, WATCH_DIR, recursive=True)
    observer.daemon = True
    observer.start()
    log.info("watching %s -> %s (debounce %ss)", WATCH_DIR, url, DEBOUNCE_SECONDS)
    try:
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        observer.stop()
    observer.join()


if __name__ == "__main__":
    try:
        main()
    except Exception:
        log.exception("fatal")
        sys.exit(1)
