# -*- coding: utf-8 -*-
# FILE: _St/Cache8701/cache8701_main.py | ROLE: Chrome CDP 캐시·응답 이미지 추출 서비스
from __future__ import annotations

import argparse
import base64
import json
import os
import re
import socket
import struct
import threading
import time
import traceback
import urllib.parse
import urllib.request
import uuid
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from datetime import datetime
from pathlib import Path
from typing import Any, Dict, Optional

ROOT = Path(__file__).resolve().parents[1]
CONFIG_PATH = Path(__file__).resolve().with_name("cache8701_config.json")
STATE_LOCK = threading.RLock()
OBSERVED: Dict[str, dict] = {}  # key: tab_id|normalized_image_url
SAVED: Dict[str, float] = {}
CACHE_RUNS: Dict[str, dict] = {}
# [Cache Image 완료 기준][회귀 금지]
# 자동 MP4 저장 여부와 무관하게 관찰된 이미지 메타데이터를 처리한다.
# 파일은 {shortcode}.{실제확장자}, 상태 원본은 SQLite이며 해시·정적 보고서를 새로 만들지 않는다.
LOG_ONCE: Dict[str, str] = {}
LOG_LOCK = threading.RLock()
LOG_FILE = None
LOG_FILE_PATH = ""
STOP = threading.Event()
CONFIG: dict = {}
RUNTIME: dict = {
    "connection": {},
    "watcher": None,
    "watcher_stop": None,
    "watcher_signature": "",
    "sessions": {},
    "last_error": "",
}


def load_config() -> dict:
    base = {
        "schema": "cache8701_config_v2",
        "version": "2.3.8.324",
        "host": "127.0.0.1",
        "port": 8701,
        "auto8700_host": "127.0.0.1",
        "auto8700_port": 8707,
        "auto8700_connection_path": "/api/cdp/connection",
        "cdp_host": "127.0.0.1",
        "cdp_port": 8700,
        "save_root": r"D:\_StDown\instagram",
        "ttl_sec": 120,
        "poll_sec": 1.0,
        "run_timeout_sec": 90,
        "file_logging": False,
        "log_keep_files": 20,
        "enabled": True,
        "page_types": {"explore": True, "single_post": True, "profile": True, "profile_reels": True},
        "mime_extensions": {"image/jpeg": ".jpg", "image/jpg": ".jpg", "image/webp": ".webp", "image/avif": ".avif", "image/png": ".png"},
    }
    try:
        raw = json.loads(CONFIG_PATH.read_text(encoding="utf-8"))
        if isinstance(raw, dict):
            base.update(raw)
            if isinstance(raw.get("page_types"), dict):
                base["page_types"].update(raw["page_types"])
            if isinstance(raw.get("mime_extensions"), dict):
                base["mime_extensions"].update(raw["mime_extensions"])
    except FileNotFoundError:
        pass
    return base


def save_config() -> None:
    tmp = CONFIG_PATH.with_suffix(".json.tmp")
    tmp.write_text(json.dumps(CONFIG, ensure_ascii=False, indent=2), encoding="utf-8")
    os.replace(tmp, CONFIG_PATH)


def _cache_log_dir() -> Path:
    return Path(__file__).resolve().with_name("logs")


def _close_file_log() -> None:
    global LOG_FILE, LOG_FILE_PATH
    with LOG_LOCK:
        current = LOG_FILE
        LOG_FILE = None
        LOG_FILE_PATH = ""
        if current is not None:
            try:
                current.flush()
                current.close()
            except Exception:
                pass


def _prune_file_logs_before_create(root: Path, keep_files: int) -> None:
    # [8701 로그 파일 보관][회귀 금지]
    # 프로그램 시작 때 새 파일을 만들기 전에 기존 파일이 20개 이상이면 오래된 파일부터 지운다.
    keep = max(1, int(keep_files or 20))
    while True:
        files = []
        for path in root.glob("cache8701_log_*.log"):
            try:
                files.append((path.stat().st_mtime, path))
            except Exception:
                continue
        if len(files) < keep:
            return
        files.sort(key=lambda item: item[0])
        try:
            files[0][1].unlink()
        except Exception:
            return


def _open_file_log(reason: str = "startup") -> str:
    global LOG_FILE, LOG_FILE_PATH
    with LOG_LOCK:
        if LOG_FILE is not None:
            return LOG_FILE_PATH
        root = _cache_log_dir()
        root.mkdir(parents=True, exist_ok=True)
        keep = max(1, int(CONFIG.get("log_keep_files") or 20))
        _prune_file_logs_before_create(root, keep)
        stamp = datetime.now().strftime("%Y%m%d_%H%M%S")
        path = root / f"cache8701_log_{stamp}_pid{os.getpid()}.log"
        LOG_FILE = path.open("a", encoding="utf-8", buffering=1)
        LOG_FILE_PATH = str(path)
        LOG_FILE.write(
            f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] "
            f"PY_LOG_START app=cache8701 pid={os.getpid()} log={path} keep_files={keep} mode=feature_check reason={reason}\n"
        )
        LOG_FILE.flush()
        return LOG_FILE_PATH


def set_file_logging(enabled: bool, reason: str = "settings") -> str:
    CONFIG["file_logging"] = bool(enabled)
    if bool(enabled):
        return _open_file_log(reason)
    _close_file_log()
    return ""


def _write_file_log(line: str) -> None:
    if not bool(CONFIG.get("file_logging", False)):
        return
    try:
        _open_file_log("lazy_start")
        with LOG_LOCK:
            if LOG_FILE is not None:
                LOG_FILE.write(f"[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] {line}\n")
                LOG_FILE.flush()
    except Exception:
        pass


def check_log(check_id: str, status: str, text: str, key: str = "") -> None:
    fingerprint = f"{status}|{text}"
    once_key = f"{check_id}|{key}"
    with STATE_LOCK:
        if LOG_ONCE.get(once_key) == fingerprint:
            return
        LOG_ONCE[once_key] = fingerprint
    line = f"[캐시이미지][{check_id}][{status}] {text}"
    print(line, flush=True)
    _write_file_log(line)


def log_exception(check_id: str, text: str, exc: BaseException, key: str = "") -> None:
    if not bool(CONFIG.get("file_logging", False)):
        return
    fingerprint = f"{type(exc).__name__}|{exc}"
    once_key = f"{check_id}|EXC|{key}"
    with STATE_LOCK:
        if LOG_ONCE.get(once_key) == fingerprint:
            return
        LOG_ONCE[once_key] = fingerprint
    detail = "".join(traceback.format_exception(type(exc), exc, exc.__traceback__)).rstrip()
    _write_file_log(f"[캐시이미지][{check_id}][예외추적] {text} / error={exc}\n{detail}")


def normalize_url(value: Any) -> str:
    text = str(value or "").strip()
    if not text:
        return ""
    try:
        u = urllib.parse.urlsplit(text)
        return urllib.parse.urlunsplit((u.scheme.lower(), u.netloc.lower(), u.path, u.query, ""))
    except Exception:
        return text.split("#", 1)[0]


def safe_name(value: Any, fallback: str = "unknown") -> str:
    text = re.sub(r'[\\/:*?"<>|\s]+', "_", str(value or "").strip()).strip("_.")
    return text[:120] or fallback


def page_type_from_url(url: str) -> str:
    path = urllib.parse.urlsplit(str(url or "")).path.strip("/")
    parts = [part for part in path.split("/") if part]
    if not parts or parts[0] in ("explore", "search"):
        return "explore"
    if parts[0] in ("p", "reel", "reels"):
        return "single_post"
    if len(parts) >= 2 and parts[1] == "reels":
        return "profile_reels"
    return "profile"


def first_value(item: dict, *keys: str) -> Any:
    for key in keys:
        value = item.get(key)
        if value not in (None, ""):
            return value
    return None


def iter_observations(payload: dict):
    source_url = str(payload.get("source_url") or payload.get("url") or "")
    for raw in payload.get("info") or payload.get("items") or []:
        if not isinstance(raw, dict):
            continue
        media = raw.get("media") if isinstance(raw.get("media"), dict) else raw
        code = first_value(raw, "shortcode", "code") or first_value(media, "shortcode", "code")
        account = first_value(raw, "username", "owner_username", "account", "account_name") or first_value(media, "username", "owner_username", "account", "account_name")
        image_url = first_value(raw, "currentSrc", "current_src", "img_origin", "image_url", "thumbnail_url", "display_url", "img_small") or first_value(media, "currentSrc", "current_src", "img_origin", "image_url", "thumbnail_url", "display_url", "img_small")
        post_url = first_value(raw, "post_url", "media_url", "href") or first_value(media, "post_url", "media_url", "href") or source_url
        if not code or not account or not image_url:
            continue
        page_type = str(first_value(raw, "page_type") or page_type_from_url(source_url or str(post_url or "")))
        metrics = {
            "views": first_value(media, "play_count", "view_count", "views"),
            "likes": first_value(media, "like_count", "likes"),
            "comments": first_value(media, "comment_count", "comments"),
            "reposts": first_value(media, "media_repost_count", "repost_count", "reposts"),
            "shares": first_value(media, "share_count", "shares"),
            "published_at": first_value(media, "created_at", "taken_at", "published_at"),
            "engagement_rate": first_value(media, "engagement_rate"),
        }
        yield {
            "account": str(account),
            "shortcode": str(code),
            "image_url": str(image_url),
            "post_url": str(post_url or ""),
            "page_type": page_type,
            "metrics": metrics,
            "tab_id": int(payload.get("tab_id") or payload.get("tabId") or 0),
            "run_id": str(payload.get("run_id") or payload.get("cache_run_id") or "").strip(),
            "observed_at": time.time(),
        }


def cache_run_snapshot(run_id: str) -> dict:
    key = str(run_id or "").strip()
    now = time.time()
    timed_out_pending = 0
    with STATE_LOCK:
        source = CACHE_RUNS.get(key) if key else None
        if source:
            target = int(source.get("target") or 0)
            processed = int(source.get("processed") or 0)
            ttl = max(10.0, float(CONFIG.get("run_timeout_sec") or 90.0))
            created_at = float(source.get("created_at") or source.get("updated_at") or now)
            pending = max(0, target - processed)
            if pending and now - created_at >= ttl:
                timed_out_pending = pending
                source["failed"] = int(source.get("failed") or 0) + pending
                source["processed"] = target
                source["updated_at"] = now
            item = dict(source)
        else:
            item = {}
    if timed_out_pending:
        tab_id = int(item.get("tab_id") or 0) if item else 0
        check_log("CI-03", "체크실패", f"CDP 이미지 URL 매칭 시간초과 / tab_id={tab_id or '-'} run_id={key} 미처리={timed_out_pending}", key or "run_timeout")
    item.setdefault("run_id", key)
    item.setdefault("target", 0)
    item.setdefault("processed", 0)
    item.setdefault("new_saved", 0)
    item.setdefault("duplicate_skipped", 0)
    item.setdefault("failed", 0)
    item.pop("shortcodes", None)
    item.pop("processed_shortcodes", None)
    item["pending"] = max(0, int(item["target"]) - int(item["processed"]))
    item["complete"] = bool(int(item["target"]) == 0 or int(item["processed"]) >= int(item["target"]))
    return item


def cache_run_register_target(obs: dict) -> None:
    run_id = str(obs.get("run_id") or "").strip()
    if not run_id:
        return
    code = str(obs.get("shortcode") or "").strip()
    with STATE_LOCK:
        item = CACHE_RUNS.setdefault(run_id, {"run_id": run_id, "tab_id": int(obs.get("tab_id") or 0), "target": 0, "processed": 0, "new_saved": 0, "duplicate_skipped": 0, "failed": 0, "shortcodes": set(), "created_at": time.time(), "updated_at": time.time()})
        codes = item.setdefault("shortcodes", set())
        if code and code not in codes:
            codes.add(code)
            item["target"] = int(item.get("target") or 0) + 1
        item["updated_at"] = time.time()


def cache_run_record(obs: dict, result: dict) -> None:
    run_id = str(obs.get("run_id") or "").strip()
    if not run_id:
        return
    code = str(obs.get("shortcode") or "").strip()
    with STATE_LOCK:
        item = CACHE_RUNS.setdefault(run_id, {"run_id": run_id, "tab_id": int(obs.get("tab_id") or 0), "target": 0, "processed": 0, "new_saved": 0, "duplicate_skipped": 0, "failed": 0, "shortcodes": set(), "processed_shortcodes": set(), "created_at": time.time(), "updated_at": time.time()})
        processed_codes = item.setdefault("processed_shortcodes", set())
        if code and code in processed_codes:
            return
        if code:
            processed_codes.add(code)
        item["processed"] = int(item.get("processed") or 0) + 1
        if not result.get("ok"):
            item["failed"] = int(item.get("failed") or 0) + 1
        elif result.get("duplicate"):
            item["duplicate_skipped"] = int(item.get("duplicate_skipped") or 0) + 1
        else:
            item["new_saved"] = int(item.get("new_saved") or 0) + 1
        item["updated_at"] = time.time()


def accept_observation(obs: dict) -> bool:
    if not CONFIG.get("enabled", True):
        return False
    page_type = str(obs.get("page_type") or "")
    if not bool((CONFIG.get("page_types") or {}).get(page_type, False)):
        return False
    key = normalize_url(obs.get("image_url"))
    if not key:
        return False
    tab_id = int(obs.get("tab_id") or 0)
    observed_key = f"{tab_id}|{key}"
    with STATE_LOCK:
        OBSERVED[observed_key] = obs
    cache_run_register_target(obs)
    check_log("CI-01", "체크완료", f"메타데이터 수신 / tab_id={tab_id or '-'} 계정={obs['account']} shortcode={obs['shortcode']} page={page_type} run_id={obs.get('run_id') or '-'}", observed_key)
    return True


def purge_old() -> None:
    cutoff = time.time() - float(CONFIG.get("ttl_sec") or 120)
    with STATE_LOCK:
        for table in (OBSERVED, SAVED):
            for key, value in list(table.items()):
                ts = value if isinstance(value, (int, float)) else float(value.get("observed_at") or value.get("ts") or 0)
                if ts < cutoff:
                    table.pop(key, None)
        for key, value in list(CACHE_RUNS.items()):
            if float((value or {}).get("updated_at") or 0) < time.time() - 3600:
                CACHE_RUNS.pop(key, None)


class SimpleCDPWebSocket:
    def __init__(self, ws_url: str, timeout: float = 5.0):
        self.ws_url = ws_url
        self.timeout = timeout
        self.sock: Optional[socket.socket] = None
        self._id = 0
        self._send_lock = threading.Lock()
        self._backlog: list[dict] = []

    def connect(self) -> None:
        parsed = urllib.parse.urlparse(self.ws_url)
        host, port = parsed.hostname or "127.0.0.1", int(parsed.port or 80)
        path = parsed.path or "/"
        if parsed.query:
            path += "?" + parsed.query
        key = base64.b64encode(os.urandom(16)).decode("ascii")
        sock = socket.create_connection((host, port), timeout=self.timeout)
        request = (
            f"GET {path} HTTP/1.1\r\nHost: {host}:{port}\r\nUpgrade: websocket\r\n"
            f"Connection: Upgrade\r\nSec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\n\r\n"
        ).encode("ascii")
        sock.sendall(request)
        raw = b""
        while b"\r\n\r\n" not in raw:
            chunk = sock.recv(4096)
            if not chunk:
                break
            raw += chunk
        if " 101 " not in raw.decode("latin1", "ignore").split("\r\n", 1)[0]:
            sock.close()
            raise RuntimeError("WebSocket handshake 실패")
        self.sock = sock
        sock.settimeout(1.0)

    def close(self) -> None:
        try:
            if self.sock:
                self.sock.close()
        finally:
            self.sock = None

    def _read_exact(self, size: int) -> bytes:
        if not self.sock:
            raise RuntimeError("websocket not connected")
        out = b""
        while len(out) < size:
            chunk = self.sock.recv(size - len(out))
            if not chunk:
                raise RuntimeError("websocket closed")
            out += chunk
        return out

    def send_text(self, text: str) -> None:
        if not self.sock:
            raise RuntimeError("websocket not connected")
        payload = text.encode("utf-8")
        header = bytearray([0x81])
        length = len(payload)
        if length < 126:
            header.append(0x80 | length)
        elif length < 65536:
            header.extend([0x80 | 126])
            header.extend(struct.pack("!H", length))
        else:
            header.extend([0x80 | 127])
            header.extend(struct.pack("!Q", length))
        mask = os.urandom(4)
        masked = bytes(byte ^ mask[index % 4] for index, byte in enumerate(payload))
        with self._send_lock:
            self.sock.sendall(bytes(header) + mask + masked)

    def recv_text(self, timeout: float = 1.0) -> Optional[str]:
        if not self.sock:
            raise RuntimeError("websocket not connected")
        self.sock.settimeout(timeout)
        try:
            first = self._read_exact(2)
        except socket.timeout:
            return None
        b1, b2 = first
        opcode = b1 & 0x0F
        length = b2 & 0x7F
        if length == 126:
            length = struct.unpack("!H", self._read_exact(2))[0]
        elif length == 127:
            length = struct.unpack("!Q", self._read_exact(8))[0]
        mask = self._read_exact(4) if b2 & 0x80 else b""
        payload = self._read_exact(length) if length else b""
        if mask:
            payload = bytes(byte ^ mask[index % 4] for index, byte in enumerate(payload))
        if opcode == 0x8:
            raise RuntimeError("websocket closed")
        if opcode == 0x9:
            return None
        return payload.decode("utf-8", "replace") if opcode == 0x1 else None

    def recv_json(self, timeout: float = 1.0) -> dict:
        if self._backlog:
            return self._backlog.pop(0)
        text = self.recv_text(timeout)
        return json.loads(text) if text else {}

    def keep_message(self, data: dict) -> None:
        if isinstance(data, dict) and data:
            self._backlog.append(data)

    def send_cmd(self, method: str, params: Optional[dict] = None, session_id: str = "") -> int:
        self._id += 1
        command_id = self._id
        payload = {"id": command_id, "method": method, "params": params or {}}
        if session_id:
            payload["sessionId"] = session_id
        self.send_text(json.dumps(payload, ensure_ascii=False))
        return command_id


def http_json(url: str, timeout: float = 2.0) -> dict:
    request = urllib.request.Request(url, headers={"User-Agent": "Cache8701/v304"}, method="GET")
    with urllib.request.urlopen(request, timeout=timeout) as response:
        payload = json.loads(response.read().decode("utf-8"))
    if not isinstance(payload, dict):
        raise RuntimeError("JSON object 응답 아님")
    return payload


def auto8700_connection_url() -> str:
    host = str(CONFIG.get("auto8700_host") or "127.0.0.1")
    port = int(CONFIG.get("auto8700_port") or 8707)
    path = str(CONFIG.get("auto8700_connection_path") or "/api/cdp/connection")
    if not path.startswith("/"):
        path = "/" + path
    return f"http://{host}:{port}{path}"


def fetch_auto8700_connection() -> dict:
    payload = http_json(auto8700_connection_url(), timeout=2.0)
    port = int(payload.get("cdp_port") or 0)
    if port <= 0:
        raise RuntimeError("Auto8700 CDP 포트 없음")
    payload["cdp_host"] = str(payload.get("cdp_host") or "127.0.0.1")
    payload["cdp_port"] = port
    payload["owner_generation"] = int(payload.get("owner_generation") or 0)
    payload["targets"] = [item for item in payload.get("targets") or [] if isinstance(item, dict) and item.get("id")]
    return payload


def browser_websocket_url(connection: dict) -> str:
    host = str(connection.get("cdp_host") or "127.0.0.1")
    port = int(connection.get("cdp_port") or 0)
    payload = http_json(f"http://{host}:{port}/json/version", timeout=2.0)
    ws_url = str(payload.get("webSocketDebuggerUrl") or "").strip()
    if not ws_url:
        raise RuntimeError("browser WebSocket URL 없음")
    return ws_url


def wait_result(cdp: SimpleCDPWebSocket, command_id: int, timeout: float = 5.0) -> dict:
    end = time.time() + timeout
    while time.time() < end:
        text = cdp.recv_text(timeout=min(1.0, max(0.05, end - time.time())))
        if not text:
            continue
        data = json.loads(text)
        if int(data.get("id") or 0) == command_id:
            return data
        cdp.keep_message(data)
    return {"error": {"message": "command_timeout"}}


def attach_target_session(cdp: SimpleCDPWebSocket, target_id: str) -> str:
    command_id = cdp.send_cmd("Target.attachToTarget", {"targetId": target_id, "flatten": True})
    response = wait_result(cdp, command_id, timeout=5.0)
    session_id = str((response.get("result") or {}).get("sessionId") or "")
    if not session_id:
        raise RuntimeError(str((response.get("error") or {}).get("message") or "attach session 실패"))
    enable_id = cdp.send_cmd(
        "Network.enable",
        {"maxTotalBufferSize": 104857600, "maxResourceBufferSize": 10485760},
        session_id=session_id,
    )
    enabled = wait_result(cdp, enable_id, timeout=5.0)
    if enabled.get("error"):
        raise RuntimeError(str((enabled.get("error") or {}).get("message") or "Network.enable 실패"))
    return session_id


def extension_for(mime: str, raw: bytes) -> str:
    if raw.startswith(b"\xff\xd8\xff"):
        return ".jpg"
    if raw.startswith(b"\x89PNG\r\n\x1a\n"):
        return ".png"
    if raw[:4] == b"RIFF" and raw[8:12] == b"WEBP":
        return ".webp"
    if b"ftypavif" in raw[:32]:
        return ".avif"
    mime = str(mime or "").split(";", 1)[0].lower()
    return str((CONFIG.get("mime_extensions") or {}).get(mime) or "")


def save_body(obs: dict, mime: str, raw: bytes, source_url: str, cache_source: str) -> dict:
    extension = extension_for(mime, raw)
    if not extension or len(raw) < 64:
        return {"ok": False, "error": "unsupported_or_empty_image"}
    import sys
    module_root = Path(__file__).resolve().parents[1]
    if str(module_root) not in sys.path:
        sys.path.insert(0, str(module_root))
    from Shared.storage.sort_storage import build_instagram_media_path, ensure_instagram_account_index, record_downloaded_file

    try:
        account = str(obs.get("account") or "").strip().lstrip("@")
        shortcode = str(obs.get("shortcode") or "").strip()
        target = build_instagram_media_path(ROOT, account, shortcode, extension)
        target.parent.mkdir(parents=True, exist_ok=True)
        duplicate = bool(target.is_file() and target.stat().st_size > 0)
        if not duplicate:
            part = target.with_suffix(target.suffix + ".part")
            part.write_bytes(raw)
            os.replace(part, target)
        index_path = ensure_instagram_account_index(ROOT, account)
    except Exception as exc:
        return {"ok": False, "error": f"file_write_failed:{exc}"}

    try:
        event_uuid = str(uuid.uuid5(uuid.NAMESPACE_URL, f"cache8701:{account.lower()}:{shortcode}:{extension}"))
        result = record_downloaded_file(
            st_root=ROOT,
            file_path=target,
            source_mode="CACHE8701_CDP",
            shortcode=shortcode,
            owner_username=account,
            asset_kind="IMAGE",
            media_index=0,
            source_url=source_url,
            post_url=str(obs.get("post_url") or ""),
            thumbnail_url=source_url,
            published_at=(obs.get("metrics") or {}).get("published_at"),
            file_size=int(target.stat().st_size),
            sha256="",
            md5="",
            account_path=account,
            account_folder_merge=False,
            batch_key=f"cache8701-{shortcode}",
            requested_count=1,
            event_uuid=event_uuid,
            payload={
                "capture_source": "Cache8701_CDP_existing_response",
                "cache_source": cache_source,
                "metrics": obs.get("metrics") or {},
                "file_contract": "{shortcode}.{actual_extension}",
                "duplicate_policy": "shortcode_file_exists",
                "hash_generation": False,
            },
        )
        return {
            "ok": True,
            "duplicate": duplicate,
            "file_path": str(target),
            "index_path": str(index_path),
            "bytes": int(target.stat().st_size),
            "registered": True,
            "register_error": "",
            "db_path": str(result.get("db_path") or ""),
            "post_id": int(result.get("post_id") or 0),
            "asset_id": int(result.get("asset_id") or 0),
            "file_name": str(result.get("file_name") or target.name),
            "relative_path": str(result.get("relative_path") or ""),
            "idempotent": bool(result.get("idempotent")),
        }
    except Exception as exc:
        return {
            "ok": True,
            "duplicate": duplicate,
            "file_path": str(target),
            "index_path": str(index_path),
            "bytes": int(target.stat().st_size),
            "registered": False,
            "register_error": str(exc),
        }


def connection_can_watch(connection: dict) -> bool:
    # [Cache Image target 판정][회귀 금지]
    # OWNER active는 화면 포커스 상태다. CDP 연결 가능 여부는 cdp_alive와 Instagram target 존재로만 판단한다.
    return bool(connection.get("cdp_alive") and (connection.get("targets") or []))


def connection_signature(connection: dict) -> str:
    targets = sorted(str(item.get("id") or "") for item in connection.get("targets") or [] if item.get("id"))
    return json.dumps(
        {
            "host": str(connection.get("cdp_host") or "127.0.0.1"),
            "port": int(connection.get("cdp_port") or 0),
            "owner": str(connection.get("owner_target_id") or ""),
            "generation": int(connection.get("owner_generation") or 0),
            "targets": targets,
        },
        sort_keys=True,
    )


def _account_from_instagram_url(url: str) -> str:
    try:
        parts = [p for p in urllib.parse.urlsplit(str(url or "")).path.split("/") if p]
        if not parts or parts[0] in {"explore", "search", "p", "reel", "reels", "stories"}:
            return ""
        return str(parts[0]).lower()
    except Exception:
        return ""


def _observation_for_session_url(tab_id: int, url: str, session_url: str) -> Optional[dict]:
    normalized = normalize_url(url)
    with STATE_LOCK:
        exact = OBSERVED.get(f"{int(tab_id or 0)}|{normalized}")
        if exact:
            return exact
        candidates = [obs for key, obs in OBSERVED.items() if key.endswith("|" + normalized)]
    if not candidates:
        return None
    session_account = _account_from_instagram_url(session_url)
    if session_account:
        matched = [obs for obs in candidates if str(obs.get("account") or "").lower().lstrip("@") == session_account]
        if len(matched) == 1:
            return matched[0]
    return candidates[0] if len(candidates) == 1 else None


def session_watcher(connection: dict, local_stop: threading.Event) -> None:
    cdp: Optional[SimpleCDPWebSocket] = None
    sessions: Dict[str, dict] = {}
    requests: Dict[tuple[str, str], dict] = {}
    loaded: Dict[str, list[dict]] = {}
    pending: Dict[int, tuple[dict, dict]] = {}
    owner_target_id = str(connection.get("owner_target_id") or "")
    try:
        ws_url = browser_websocket_url(connection)
        cdp = SimpleCDPWebSocket(ws_url)
        cdp.connect()
        targets = list(connection.get("targets") or [])
        targets.sort(key=lambda item: 0 if str(item.get("id") or "") == owner_target_id else 1)
        for target in targets:
            target_id = str(target.get("id") or "")
            if not target_id:
                continue
            try:
                session_id = attach_target_session(cdp, target_id)
                sessions[session_id] = {"target_id": target_id, "url": str(target.get("url") or ""), "tab_id": int(target.get("tab_id") or target.get("tabId") or 0)}
                check_log(
                    "CI-02",
                    "체크완료",
                    f"독립 CDP 세션 연결 / tab_id={sessions[session_id].get('tab_id') or '-'} target={target_id} session={session_id[:12]} owner={target_id == owner_target_id} url={sessions[session_id].get('url') or '-'}",
                    target_id,
                )
            except Exception as exc:
                check_log("CI-02", "체크실패", f"CDP 세션 연결 실패 / target={target_id} error={exc}", target_id)
        if not sessions:
            raise RuntimeError("연결된 Instagram CDP session 없음")
        with STATE_LOCK:
            RUNTIME["sessions"] = dict(sessions)
            RUNTIME["last_error"] = ""
            LOG_ONCE.pop("CI-02|session_watcher", None)
            LOG_ONCE.pop("CI-02-E01|EXC|session_watcher", None)
        while not STOP.is_set() and not local_stop.is_set():
            data = cdp.recv_json(1.0)
            method = str(data.get("method") or "")
            params = data.get("params") or {}
            session_id = str(data.get("sessionId") or "")
            if method == "Network.responseReceived" and session_id in sessions:
                response = params.get("response") or {}
                mime = str(response.get("mimeType") or "").split(";", 1)[0].lower()
                request_id = str(params.get("requestId") or "")
                if request_id and str(params.get("type") or "") == "Image" and mime in (CONFIG.get("mime_extensions") or {}):
                    url = normalize_url(response.get("url"))
                    requests[(session_id, request_id)] = {
                        "url": url,
                        "mime": mime,
                        "source": "disk_cache" if response.get("fromDiskCache") else "prefetch_cache" if response.get("fromPrefetchCache") else "service_worker" if response.get("fromServiceWorker") else "network_or_memory",
                        "ts": time.time(),
                        "session_id": session_id,
                    }
            elif method == "Network.loadingFinished" and session_id in sessions:
                request_id = str(params.get("requestId") or "")
                item = requests.pop((session_id, request_id), None)
                if item:
                    item["request_id"] = request_id
                    item["loaded_at"] = time.time()
                    loaded.setdefault(item["url"], []).append(item)
            elif int(data.get("id") or 0) in pending:
                item, obs = pending.pop(int(data.get("id") or 0))
                result = data.get("result") or {}
                if not isinstance(result.get("body"), str):
                    failed = {"ok": False, "error": "response_body_missing"}
                    cache_run_record(obs, failed)
                    check_log(
                        "CI-04",
                        "체크실패",
                        f"응답 body 없음 / tab_id={obs.get('tab_id') or '-'} shortcode={obs['shortcode']} error={(data.get('error') or {}).get('message', 'unknown')}",
                        item["url"],
                    )
                    continue
                raw = base64.b64decode(result["body"], validate=False) if result.get("base64Encoded") else result["body"].encode("latin1")
                check_log("CI-04", "체크완료", f"이미지 body 확보 / tab_id={obs.get('tab_id') or '-'} shortcode={obs['shortcode']} bytes={len(raw)}", item["url"])
                saved = save_body(obs, item["mime"], raw, item["url"], item["source"])
                cache_run_record(obs, saved)
                if saved.get("ok"):
                    save_key = f"{obs.get('shortcode')}|{item['url']}"
                    with STATE_LOCK:
                        SAVED[save_key] = time.time()
                    check_log("CI-05", "체크완료", f"파일 저장 / tab_id={obs.get('tab_id') or '-'} path={saved.get('file_path')} duplicate={saved.get('duplicate', False)}", item["url"])
                    if saved.get("registered"):
                        check_log(
                            "CI-06",
                            "체크완료",
                            f"SQLite 등록 / tab_id={obs.get('tab_id') or '-'} db={Path(str(saved.get('db_path') or '')).name} post_id={saved.get('post_id')} asset_id={saved.get('asset_id')} file={saved.get('file_name')}",
                            item["url"],
                        )
                    elif saved.get("register_error"):
                        check_log("CI-06", "체크실패", f"SQLite 등록 실패 / {saved['register_error']}", item["url"])
                else:
                    check_log("CI-05", "체크실패", f"파일 저장 실패 / {saved.get('error')}", item["url"])
            cutoff = time.time() - float(CONFIG.get("ttl_sec") or 120)
            for url, items in list(loaded.items()):
                items[:] = [item for item in items if float(item.get("loaded_at") or 0) >= cutoff]
                if not items:
                    loaded.pop(url, None)
                    continue
                session_meta = sessions.get(str(items[0].get("session_id") or ""), {}) if items else {}
                observation = _observation_for_session_url(int(session_meta.get("tab_id") or 0), url, str(session_meta.get("url") or ""))
                if not observation:
                    continue
                item = items.pop(0)
                if not items:
                    loaded.pop(url, None)
                check_log("CI-03", "체크완료", f"CDP URL 매칭 / tab_id={observation.get('tab_id') or '-'} shortcode={observation['shortcode']} source={item['source']} run_id={observation.get('run_id') or '-'}", url)
                command_id = cdp.send_cmd(
                    "Network.getResponseBody",
                    {"requestId": item["request_id"]},
                    session_id=item["session_id"],
                )
                pending[command_id] = (item, observation)
    except Exception as exc:
        with STATE_LOCK:
            RUNTIME["last_error"] = str(exc)
        check_log("CI-02", "체크실패", f"독립 CDP 세션 종료 / error={exc}", "session_watcher")
        log_exception("CI-02-E01", "독립 CDP 세션 예외", exc, "session_watcher")
    finally:
        if cdp:
            cdp.close()
        with STATE_LOCK:
            RUNTIME["sessions"] = {}


def stop_runtime_watcher() -> None:
    with STATE_LOCK:
        event = RUNTIME.get("watcher_stop")
        thread = RUNTIME.get("watcher")
        RUNTIME["watcher_stop"] = None
        RUNTIME["watcher"] = None
        RUNTIME["watcher_signature"] = ""
    if isinstance(event, threading.Event):
        event.set()
    if isinstance(thread, threading.Thread) and thread.is_alive():
        thread.join(timeout=2.0)


def target_manager() -> None:
    check_log("CI-00", "체크시작", f"Cache8701 시작 / Auto8700={auto8700_connection_url()}")
    while not STOP.is_set():
        purge_old()
        try:
            connection = fetch_auto8700_connection()
            with STATE_LOCK:
                RUNTIME["connection"] = dict(connection)
            port = int(connection.get("cdp_port") or 0)
            owner = str(connection.get("owner_target_id") or "")
            generation = int(connection.get("owner_generation") or 0)
            if not connection.get("cdp_alive"):
                raise RuntimeError(f"Chrome CDP 준비 안 됨 / port={port}")
            # [Cache Image 기존 연결 보호][회귀 금지]
            # OWNER active=false는 사용자가 다른 탭을 보고 있다는 뜻일 뿐 CDP target 단절이 아니다.
            # Instagram target이 존재하면 비활성 OWNER를 포함한 전체 target 감시를 계속한다.
            if not connection_can_watch(connection):
                raise RuntimeError(f"Instagram target 없음 / port={port}")
            signature = connection_signature(connection)
            with STATE_LOCK:
                LOG_ONCE.pop("CI-00|connection_wait", None)
                LOG_ONCE.pop("CI-00-E01|EXC|connection_wait", None)
            check_log(
                "CI-00",
                "체크완료",
                f"연결정보 확인 / cdp_port={port} owner={owner} owner_active={bool(connection.get('owner_active'))} generation={generation} targets={len(connection.get('targets') or [])}",
                signature,
            )
            with STATE_LOCK:
                current_signature = str(RUNTIME.get("watcher_signature") or "")
                current_thread = RUNTIME.get("watcher")
            if signature != current_signature or not isinstance(current_thread, threading.Thread) or not current_thread.is_alive():
                stop_runtime_watcher()
                local_stop = threading.Event()
                thread = threading.Thread(
                    target=session_watcher,
                    args=(connection, local_stop),
                    daemon=True,
                    name=f"Cache8701-CDP-{generation}",
                )
                with STATE_LOCK:
                    RUNTIME["watcher_signature"] = signature
                    RUNTIME["watcher_stop"] = local_stop
                    RUNTIME["watcher"] = thread
                thread.start()
        except Exception as exc:
            stop_runtime_watcher()
            with STATE_LOCK:
                RUNTIME["last_error"] = str(exc)
            check_log("CI-00", "체크실패", f"Auto8700 CDP 연결정보 대기 / {exc}", "connection_wait")
            log_exception("CI-00-E01", "Auto8700 연결정보 또는 target manager 예외", exc, "connection_wait")
        STOP.wait(float(CONFIG.get("poll_sec") or 1.0))
    stop_runtime_watcher()


class Handler(BaseHTTPRequestHandler):
    def log_message(self, fmt, *args):
        return

    def cors(self) -> None:
        self.send_header("Access-Control-Allow-Origin", "*")
        self.send_header("Access-Control-Allow-Headers", "Content-Type")
        self.send_header("Access-Control-Allow-Methods", "GET,POST,OPTIONS")

    def send_json(self, status: int, payload: dict) -> None:
        raw = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self.send_response(status)
        self.send_header("Content-Type", "application/json; charset=utf-8")
        self.send_header("Content-Length", str(len(raw)))
        self.cors()
        self.end_headers()
        self.wfile.write(raw)

    def do_OPTIONS(self) -> None:
        self.send_response(204)
        self.cors()
        self.end_headers()

    def do_GET(self) -> None:
        if self.path.startswith("/api/status"):
            query = urllib.parse.parse_qs(urllib.parse.urlsplit(self.path).query)
            run_id = str((query.get("run_id") or [""])[0]).strip()
            return self.send_json(200, {"ok": True, "cache": cache_run_snapshot(run_id)})
        if self.path.startswith("/api/health"):
            with STATE_LOCK:
                connection = dict(RUNTIME.get("connection") or {})
                payload = {
                    "ok": True,
                    "version": CONFIG.get("version"),
                    "port": CONFIG.get("port"),
                    "observed": len(OBSERVED),
                    "sessions": len(RUNTIME.get("sessions") or {}),
                    "saved_recent": len(SAVED),
                    "cdp_port": connection.get("cdp_port"),
                    "owner_target_id": connection.get("owner_target_id"),
                    "owner_generation": connection.get("owner_generation"),
                    "last_error": RUNTIME.get("last_error") or "",
                    "file_logging": bool(CONFIG.get("file_logging", False)),
                    "log_file": LOG_FILE_PATH,
                }
            return self.send_json(200, payload)
        return self.send_json(404, {"ok": False, "error": "not_found"})

    def do_POST(self) -> None:
        try:
            length = int(self.headers.get("Content-Length") or 0)
            payload = json.loads(self.rfile.read(length).decode("utf-8") or "{}")
        except Exception as exc:
            return self.send_json(400, {"ok": False, "error": str(exc)})
        if self.path.startswith("/api/meta"):
            observations = list(iter_observations(payload))
            accepted = sum(1 for observation in observations if accept_observation(observation))
            run_id = str(payload.get("run_id") or payload.get("cache_run_id") or "").strip()
            return self.send_json(200, {"ok": True, "accepted": accepted, "cache": cache_run_snapshot(run_id)})
        if self.path.startswith("/api/settings"):
            settings = payload.get("settings") if isinstance(payload.get("settings"), dict) else payload
            CONFIG["enabled"] = bool(settings.get("auto_thumbnail_enabled", CONFIG.get("enabled", True)))
            logging_value = bool(settings.get("cache8701_file_logging", CONFIG.get("file_logging", False)))
            CONFIG["page_types"] = {
                "explore": bool(settings.get("auto_thumbnail_explore", False)),
                "single_post": bool(settings.get("auto_thumbnail_single_post", True)),
                "profile": bool(settings.get("auto_thumbnail_profile", True)),
                "profile_reels": bool(settings.get("auto_thumbnail_profile_reels", True)),
            }
            set_file_logging(logging_value, "settings_change")
            save_config()
            check_log("CI-00", "체크완료", f"설정 반영 / enabled={CONFIG['enabled']} logging={CONFIG['file_logging']} log_file={LOG_FILE_PATH or '-'} pages={CONFIG['page_types']}", "settings")
            return self.send_json(200, {"ok": True, "settings": {"enabled": CONFIG["enabled"], "file_logging": CONFIG["file_logging"], "log_file": LOG_FILE_PATH, "page_types": CONFIG["page_types"]}})
        return self.send_json(404, {"ok": False, "error": "not_found"})


def main() -> None:
    global CONFIG
    parser = argparse.ArgumentParser()
    parser.add_argument("--host")
    parser.add_argument("--port", type=int)
    parser.add_argument("--cdp-port", type=int, help="진단용 강제값. 정상 실행은 Auto8700 연결정보를 사용한다.")
    args = parser.parse_args()
    CONFIG = load_config()
    if args.host:
        CONFIG["host"] = args.host
    if args.port:
        CONFIG["port"] = args.port
    if args.cdp_port:
        CONFIG["cdp_port"] = args.cdp_port
    if bool(CONFIG.get("file_logging", False)):
        set_file_logging(True, "program_start")
    threading.Thread(target=target_manager, daemon=True, name="Cache8701TargetManager").start()
    server = ThreadingHTTPServer((str(CONFIG.get("host") or "127.0.0.1"), int(CONFIG.get("port") or 8701)), Handler)
    check_log("CI-00", "체크완료", f"HTTP 준비 / http://{CONFIG['host']}:{CONFIG['port']}")
    try:
        server.serve_forever()
    except KeyboardInterrupt:
        pass
    finally:
        STOP.set()
        server.server_close()
        _close_file_log()


if __name__ == "__main__":
    main()
