"""
File watcher management and watcher route handlers.
"""
import asyncio
import os
from pathlib import Path
from typing import Any

from aiohttp import web
from mjr_am_backend.features.index.watcher_scope import (
    build_watch_paths,
    normalize_scope,
    persist_watcher_scope,
    resolve_service_watcher_scope,
)
from mjr_am_backend.features.watcher_settings import get_watcher_settings, update_watcher_settings
from mjr_am_backend.shared import Result, get_logger

from ..core import (
    _csrf_error,
    _json_response,
    _read_json,
    _require_services,
    _require_write_access,
    audit_log_write,
    safe_error_message,
)
from .scan_helpers import (
    _refresh_watcher_runtime_settings,
    _watcher_settings_from_body,
)

logger = get_logger(__name__)


def _current_request_user_id() -> str:
    try:
        from ..core import _current_user_id

        return str(_current_user_id() or "").strip()
    except Exception:
        return ""


def _persist_service_watcher_scope_preference(svc: dict[str, Any], scope: str, custom_root_id: str) -> str:
    payload = {"scope": normalize_scope(scope), "custom_root_id": str(custom_root_id or "").strip()}
    user_id = _current_request_user_id()
    if user_id:
        scope_map = svc.setdefault("watcher_scope_by_user", {})
        if isinstance(scope_map, dict):
            scope_map[user_id] = payload
    else:
        svc["watcher_scope"] = payload
    return user_id


def _set_active_watcher_scope(svc: dict[str, Any], scope: str, custom_root_id: str, user_id: str | None = None) -> None:
    svc["watcher_scope"] = {
        "scope": normalize_scope(scope),
        "custom_root_id": str(custom_root_id or "").strip(),
    }
    svc["watcher_scope_active_user_id"] = str(user_id or "").strip()


def _build_watch_paths_for_context(build_watch_paths_fn, scope: str, custom_root_id: str | None, user_id: str | None = None):
    effective_user_id = str(user_id or "").strip()
    if effective_user_id:
        try:
            return build_watch_paths_fn(scope, custom_root_id, user_id=effective_user_id)
        except TypeError:
            pass
    return build_watch_paths_fn(scope, custom_root_id)


def _watcher_is_running(watcher: Any) -> bool:
    if not watcher:
        return False
    try:
        raw = getattr(watcher, "is_running", False)
        value = raw() if callable(raw) else raw
        return bool(value)
    except Exception:
        return False


def _watcher_directories(watcher: Any) -> list[str]:
    if not watcher:
        return []
    try:
        raw = getattr(watcher, "watched_directories", [])
        value = raw() if callable(raw) else raw
        if isinstance(value, (list, tuple, set)):
            return [str(path) for path in value if path]
    except Exception:
        pass
    return []


# ---------------------------------------------------------------------------
# Watcher callback builders
# ---------------------------------------------------------------------------

def _watcher_scope_config(svc: dict[str, Any]) -> tuple[str, str | None]:
    scope_cfg = resolve_service_watcher_scope(svc if isinstance(svc, dict) else None, user_id=_current_request_user_id() or None)
    desired_scope = normalize_scope((scope_cfg or {}).get("scope"))
    desired_root_id = (scope_cfg or {}).get("custom_root_id")
    return desired_scope, desired_root_id


async def _delay_for_recent_generated_marker() -> None:
    # Skip the delay when there are no recently-generated entries to filter out.
    # The 200ms window exists so ComfyUI's own "executed" indexing path can mark
    # a file before the watcher callback processes it.  If the cache is empty no
    # ComfyUI generation is in-flight, so the delay is unnecessary.
    try:
        from mjr_am_backend.features.index.watcher import _RECENT_GENERATED
        if not _RECENT_GENERATED:
            return
    except Exception:
        pass
    try:
        await asyncio.sleep(0.2)
    except Exception:
        return


def _filter_recent_generated_files(filepaths: list[Any] | None) -> list[Any]:
    try:
        from mjr_am_backend.features.index.watcher import is_recent_generated

        return [f for f in (filepaths or []) if f and not is_recent_generated(f)]
    except Exception:
        return [f for f in (filepaths or []) if f]


def _build_watcher_callbacks(index_service: Any):
    async def index_callback(filepaths, base_dir, source=None, root_id=None):
        if not filepaths:
            return
        await _delay_for_recent_generated_marker()
        filepaths = _filter_recent_generated_files(filepaths)
        if not filepaths:
            return
        paths = [Path(f) for f in filepaths if f]
        if paths:
            await index_service.index_paths(
                paths=paths,
                base_dir=base_dir,
                incremental=True,
                source=source or "watcher",
                root_id=root_id,
            )

    async def remove_callback(filepaths, _base_dir, _source=None, _root_id=None):
        if not filepaths:
            return
        for fp in filepaths:
            try:
                await index_service.remove_file(str(fp))
            except Exception:
                continue

    async def move_callback(moves, _base_dir, source=None, root_id=None):
        if not moves:
            return
        for move in moves:
            try:
                old_fp, new_fp = move
            except Exception:
                continue
            try:
                res = await index_service.rename_file(str(old_fp), str(new_fp))
                if not res.ok:
                    await index_service.remove_file(str(old_fp))
                    await index_service.index_paths(
                        paths=[Path(str(new_fp))],
                        base_dir=str(_base_dir),
                        incremental=True,
                        source=source or "watcher",
                        root_id=root_id,
                    )
            except Exception:
                continue

    return index_callback, remove_callback, move_callback


# ---------------------------------------------------------------------------
# Watcher start / stop
# ---------------------------------------------------------------------------

async def _start_watcher_for_scope(svc: dict[str, Any], index_service: Any, *, _build_watch_paths=None) -> Result[dict[str, Any]]:
    from mjr_am_backend.features.index.watcher import OutputWatcher

    if _build_watch_paths is None:
        _build_watch_paths = build_watch_paths
    index_callback, remove_callback, move_callback = _build_watcher_callbacks(index_service)
    new_watcher = OutputWatcher(index_callback, remove_callback=remove_callback, move_callback=move_callback)
    desired_scope, desired_root_id = _watcher_scope_config(svc)
    request_user_id = _current_request_user_id()
    watch_paths = _build_watch_paths_for_context(_build_watch_paths, desired_scope, desired_root_id, request_user_id)
    if not watch_paths:
        return Result.Err("NO_DIRECTORIES", "No directories to watch")
    loop = asyncio.get_running_loop()
    await new_watcher.start(watch_paths, loop)
    svc["watcher"] = new_watcher
    _set_active_watcher_scope(svc, desired_scope, desired_root_id or "", request_user_id)
    return Result.Ok({"enabled": True, "directories": new_watcher.watched_directories})


async def _stop_watcher_if_running(svc: dict[str, Any], watcher: Any) -> Result[dict[str, Any]]:
    if watcher:
        await watcher.stop()
        svc["watcher"] = None
    return Result.Ok({"enabled": False, "directories": []})


# ---------------------------------------------------------------------------
# Watcher route handlers (registered via register_watcher_routes)
# ---------------------------------------------------------------------------

def register_watcher_routes(routes: web.RouteTableDef, *, deps: dict | None = None) -> None:
    """Register all /mjr/am/watcher/* route handlers.

    deps: if provided (pass globals() from the calling scan module), the handlers
          look up patchable functions from this dict AT CALL TIME so that test patches
          applied after route registration (e.g. _app() called before monkeypatch) work.
    """
    # _d captured by reference — scan.__dict__ is live, so patches applied
    # at any point (even after _app() is called) are picked up at handler call time.
    _d = deps  # May be None; fallback to local names when None.

    async def _audit_watcher_write(
        services: dict[str, Any] | None,
        request: web.Request,
        operation: str,
        target: str,
        result: Result,
        **details: Any,
    ) -> None:
        _audit_log_write_ = _d.get("audit_log_write", audit_log_write) if _d is not None else audit_log_write
        try:
            await _audit_log_write_(
                services if isinstance(services, dict) else {},
                request=request,
                operation=operation,
                target=target,
                result=result,
                details=details or None,
            )
        except Exception as exc:
            logger.debug("Watcher audit logging skipped for %s: %s", operation, exc)

    @routes.get("/mjr/am/watcher/status")
    async def watcher_status(request):
        """Get watcher status."""
        _require_services_ = _d.get("_require_services", _require_services) if _d is not None else _require_services
        svc, error_result = await _require_services_()
        if error_result:
            return _json_response(error_result)

        watcher = svc.get("watcher")
        return _json_response(Result.Ok({
            "enabled": _watcher_is_running(watcher),
            "directories": _watcher_directories(watcher),
        }))

    @routes.post("/mjr/am/watcher/flush")
    async def watcher_flush(request):
        """Flush watcher pending queue immediately (best-effort)."""
        _require_services_ = _d.get("_require_services", _require_services) if _d is not None else _require_services
        _csrf_error_ = _d.get("_csrf_error", _csrf_error) if _d is not None else _csrf_error
        _require_write_access_ = _d.get("_require_write_access", _require_write_access) if _d is not None else _require_write_access

        svc, error_result = await _require_services_()
        if error_result:
            return _json_response(error_result)

        csrf = _csrf_error_(request)
        if csrf:
            return _json_response(Result.Err("CSRF", csrf))
        auth = _require_write_access_(request)
        if not auth.ok:
            return _json_response(auth)

        watcher = svc.get("watcher")
        if not watcher:
            result = Result.Ok({"enabled": False, "flushed": False})
            await _audit_watcher_write(svc, request, "watcher_flush", "watcher:flush", result, enabled=False, flushed=False)
            return _json_response(result)

        flushed = False
        try:
            flush_fn = getattr(watcher, "flush_pending", None)
            if callable(flush_fn):
                flushed = bool(flush_fn())
        except Exception:
            flushed = False
        result = Result.Ok({"enabled": _watcher_is_running(watcher), "flushed": flushed})
        await _audit_watcher_write(
            svc,
            request,
            "watcher_flush",
            "watcher:flush",
            result,
            enabled=bool(result.data.get("enabled")) if isinstance(result.data, dict) else False,
            flushed=flushed,
        )
        return _json_response(result)

    @routes.post("/mjr/am/watcher/toggle")
    async def watcher_toggle(request):
        """Start or stop the file watcher."""
        _require_services_ = _d.get("_require_services", _require_services) if _d is not None else _require_services
        _csrf_error_ = _d.get("_csrf_error", _csrf_error) if _d is not None else _csrf_error
        _require_write_access_ = _d.get("_require_write_access", _require_write_access) if _d is not None else _require_write_access
        _read_json_ = _d.get("_read_json", _read_json) if _d is not None else _read_json
        _build_watch_paths_ = _d.get("build_watch_paths", build_watch_paths) if _d is not None else build_watch_paths

        svc, error_result = await _require_services_()
        if error_result:
            return _json_response(error_result)

        csrf = _csrf_error_(request)
        if csrf:
            return _json_response(Result.Err("CSRF", csrf))
        auth = _require_write_access_(request)
        if not auth.ok:
            return _json_response(auth)

        body_res = await _read_json_(request)
        if not body_res.ok:
            return _json_response(body_res)
        body = body_res.data or {}

        enabled = body.get("enabled", True)
        watcher = svc.get("watcher")

        if enabled:
            if _watcher_is_running(watcher):
                result = Result.Ok({"enabled": True, "directories": _watcher_directories(watcher)})
                await _audit_watcher_write(svc, request, "watcher_toggle", "watcher:toggle", result, enabled=True)
                return _json_response(result)

            index_service = svc.get("index")
            if not index_service:
                result = Result.Err("SERVICE_UNAVAILABLE", "Index service not available")
                await _audit_watcher_write(svc, request, "watcher_toggle", "watcher:toggle", result, enabled=True)
                return _json_response(result)
            result = await _start_watcher_for_scope(svc, index_service, _build_watch_paths=_build_watch_paths_)
            await _audit_watcher_write(svc, request, "watcher_toggle", "watcher:toggle", result, enabled=True)
            return _json_response(result)
        result = await _stop_watcher_if_running(svc, watcher)
        await _audit_watcher_write(svc, request, "watcher_toggle", "watcher:toggle", result, enabled=False)
        return _json_response(result)

    @routes.get("/mjr/am/watcher/settings")
    async def watcher_settings_get(request):
        _require_services_ = _d.get("_require_services", _require_services) if _d is not None else _require_services
        _get_watcher_settings_ = _d.get("get_watcher_settings", get_watcher_settings) if _d is not None else get_watcher_settings

        svc, error_result = await _require_services_()
        if error_result:
            return _json_response(error_result)
        try:
            settings = _get_watcher_settings_()
        except Exception as exc:
            logger.debug("Failed to read watcher settings: %s", exc)
            return _json_response(Result.Err("DEGRADED", "Watcher settings unavailable"))
        return _json_response(
            Result.Ok(
                {
                    "debounce_ms": settings.debounce_ms,
                    "dedupe_ttl_ms": settings.dedupe_ttl_ms,
                }
            )
        )

    @routes.post("/mjr/am/watcher/settings")
    async def watcher_settings_update(request):
        _require_services_ = _d.get("_require_services", _require_services) if _d is not None else _require_services
        _csrf_error_ = _d.get("_csrf_error", _csrf_error) if _d is not None else _csrf_error
        _require_write_access_ = _d.get("_require_write_access", _require_write_access) if _d is not None else _require_write_access
        _read_json_ = _d.get("_read_json", _read_json) if _d is not None else _read_json
        _update_watcher_settings_ = _d.get("update_watcher_settings", update_watcher_settings) if _d is not None else update_watcher_settings

        svc, error_result = await _require_services_()
        if error_result:
            return _json_response(error_result)

        csrf = _csrf_error_(request)
        if csrf:
            return _json_response(Result.Err("CSRF", csrf))
        auth = _require_write_access_(request)
        if not auth.ok:
            return _json_response(auth)

        body_res = await _read_json_(request)
        if not body_res.ok:
            return _json_response(body_res)
        body = body_res.data or {}

        try:
            debounce_ms, dedupe_ttl_ms = _watcher_settings_from_body(body)
        except ValueError as exc:
            return _json_response(Result.Err("INVALID_INPUT", str(exc)))

        if debounce_ms is None and dedupe_ttl_ms is None:
            return _json_response(Result.Err("INVALID_INPUT", "No watcher settings provided"))

        try:
            settings = _update_watcher_settings_(
                debounce_ms=debounce_ms,
                dedupe_ttl_ms=dedupe_ttl_ms,
            )
        except Exception as exc:
            result = Result.Err("DEGRADED", safe_error_message(exc, "Failed to update watcher settings"))
            await _audit_watcher_write(
                svc,
                request,
                "watcher_settings_update",
                "watcher:settings",
                result,
                debounce_ms=debounce_ms,
                dedupe_ttl_ms=dedupe_ttl_ms,
            )
            return _json_response(result)

        _refresh_watcher_runtime_settings(svc)

        result = Result.Ok(
            {
                "debounce_ms": settings.debounce_ms,
                "dedupe_ttl_ms": settings.dedupe_ttl_ms,
            }
        )
        await _audit_watcher_write(
            svc,
            request,
            "watcher_settings_update",
            "watcher:settings",
            result,
            debounce_ms=settings.debounce_ms,
            dedupe_ttl_ms=settings.dedupe_ttl_ms,
        )
        return _json_response(result)

    @routes.post("/mjr/am/watcher/scope")
    async def watcher_scope(request):
        """Configure watcher to follow a single scope (global mode)."""
        _require_services_ = _d.get("_require_services", _require_services) if _d is not None else _require_services
        _csrf_error_ = _d.get("_csrf_error", _csrf_error) if _d is not None else _csrf_error
        _require_write_access_ = _d.get("_require_write_access", _require_write_access) if _d is not None else _require_write_access
        _read_json_ = _d.get("_read_json", _read_json) if _d is not None else _read_json
        _build_watch_paths_ = _d.get("build_watch_paths", build_watch_paths) if _d is not None else build_watch_paths

        svc, error_result = await _require_services_()
        if error_result:
            return _json_response(error_result)

        csrf = _csrf_error_(request)
        if csrf:
            return _json_response(Result.Err("CSRF", csrf))
        auth = _require_write_access_(request)
        if not auth.ok:
            return _json_response(auth)

        body_res = await _read_json_(request)
        if not body_res.ok:
            return _json_response(body_res)
        body = body_res.data or {}

        scope = normalize_scope(body.get("scope"))
        custom_root_id = str(body.get("custom_root_id") or body.get("root_id") or "").strip()
        if scope != "custom":
            custom_root_id = ""

        request_user_id = _persist_service_watcher_scope_preference(svc, scope, custom_root_id)
        try:
            db = svc.get("db")
            if db:
                await persist_watcher_scope(db, scope, custom_root_id, user_id=request_user_id or None)
        except Exception:
            pass

        # If watcher is not enabled, just acknowledge.
        watcher = svc.get("watcher")
        if not _watcher_is_running(watcher):
            result = Result.Ok({"enabled": False, "directories": [], "scope": scope})
            await _audit_watcher_write(
                svc,
                request,
                "watcher_scope",
                f"watcher:scope:{scope}",
                result,
                scope=scope,
                custom_root_id=custom_root_id,
            )
            return _json_response(result)

        from mjr_am_backend.features.index.watcher import OutputWatcher

        index_service = svc.get("index")
        if not index_service:
            return _json_response(Result.Err("SERVICE_UNAVAILABLE", "Index service not available"))

        watch_paths = _build_watch_paths_for_context(_build_watch_paths_, scope, custom_root_id, request_user_id)
        if not watch_paths:
            if scope == "custom":
                result = Result.Err("INVALID_INPUT", "Invalid or missing custom_root_id for watcher scope")
                await _audit_watcher_write(
                    svc,
                    request,
                    "watcher_scope",
                    f"watcher:scope:{scope}",
                    result,
                    scope=scope,
                    custom_root_id=custom_root_id,
                )
                return _json_response(result)
            try:
                if _watcher_is_running(watcher):
                    await watcher.stop()
                    svc["watcher"] = None
            except Exception:
                pass
            result = Result.Ok({"enabled": False, "directories": [], "scope": scope})
            await _audit_watcher_write(
                svc,
                request,
                "watcher_scope",
                f"watcher:scope:{scope}",
                result,
                scope=scope,
                custom_root_id=custom_root_id,
            )
            return _json_response(result)

        # Idempotent scope update: avoid watcher stop/start churn when the
        # requested scope resolves to the exact same set of directories.
        try:
            current_dirs = []
            if _watcher_is_running(watcher):
                current_dirs = [os.path.normcase(os.path.normpath(str(p))) for p in (_watcher_directories(watcher) or []) if p]

            desired_dirs = []
            for entry in watch_paths:
                path_value = None
                if isinstance(entry, dict):
                    path_value = entry.get("path")
                else:
                    path_value = entry
                if path_value:
                    desired_dirs.append(os.path.normcase(os.path.normpath(str(path_value))))

            if _watcher_is_running(watcher) and set(current_dirs) == set(desired_dirs):
                _set_active_watcher_scope(svc, scope, custom_root_id, request_user_id)
                result = Result.Ok({"enabled": True, "directories": _watcher_directories(watcher), "scope": scope})
                await _audit_watcher_write(
                    svc,
                    request,
                    "watcher_scope",
                    f"watcher:scope:{scope}",
                    result,
                    scope=scope,
                    custom_root_id=custom_root_id,
                )
                return _json_response(result)
        except Exception:
            pass

        try:
            await watcher.stop()
        except Exception:
            pass

        # Use the shared callback builder so all watcher instances share the same
        # recent-generated filtering logic (BUG-02: was previously duplicated inline).
        index_cb, remove_cb, move_cb = _build_watcher_callbacks(index_service)
        new_watcher = OutputWatcher(index_cb, remove_callback=remove_cb, move_callback=move_cb)
        loop = asyncio.get_running_loop()
        await new_watcher.start(watch_paths, loop)
        svc["watcher"] = new_watcher
        _set_active_watcher_scope(svc, scope, custom_root_id, request_user_id)
        result = Result.Ok({"enabled": True, "directories": new_watcher.watched_directories, "scope": scope})
        await _audit_watcher_write(
            svc,
            request,
            "watcher_scope",
            f"watcher:scope:{scope}",
            result,
            scope=scope,
            custom_root_id=custom_root_id,
        )
        return _json_response(result)
