#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
Domain Security Observer — Watchdog Python 1.1.0
================================================

Lapisan pemulih utama agent untuk VPS/server yang bisa menjalankan Python
via cron. Ditempatkan DI LUAR webroot (tak terlihat WAF) dan dipicu cron
(entry cron bukan file — tak bisa dihapus scanner file).

POOL 5 FILE SALING-REGENERASI: setiap eksekusi memastikan kelima salinan
(dso_watchdog.py + dso_watchdog2..5.py) ada di direktori yang sama; yang
hilang ditulis ulang dari salinan yang sedang berjalan. Cron mencoba
semuanya, jadi selama 1 dari 5 masih hidup, pool kembali menjadi 5.

Tugas TUNGGAL per domain:
  1. Pastikan file agent ada dan isinya PERSIS seperti sumber kanonik
     (template resmi dari Cloudflare Pages + secret domain) — hilang,
     rusak, atau versi lama langsung dibangun ulang secara atomik.
  2. (Opsional, interval default 180 dtk) Heartbeat ke Worker: dashboard
     menampilkan chip "♥ watchdog" dan verifikasi menandai Python AKTIF;
     Worker menaikkan insiden bila heartbeat berhenti > 12 menit.

Dedup: dengan 5 salinan hidup, pemrosesan domain + heartbeat hanya
efektif SEKALI per menit (gerbang file dso_last_run 50 dtk).

Tanpa eval, tanpa parameter, tanpa kemampuan lain. Secret dibaca dari
file config JSON terpisah (chmod 600). Komentar ini bagian dari desain —
jangan dihapus/obfuscate.

Cara pasang (VPS):
  1. mkdir -p /opt/dso && taruh dso_watchdog.py di sana
  2. Salin dso_watchdog.example.json -> /opt/dso/dso_watchdog.json
     lalu isi (agent_path, secret, domain_id, worker_base per domain).
     chmod 600 /opt/dso/dso_watchdog.json
  3. Uji:     python3 /opt/dso/dso_watchdog.py --doctor
  4. Cron (mencoba kelima salinan — tahan penghapusan):
     * * * * * for f in /opt/dso/dso_watchdog*.py; do [ -f "$f" ] && /usr/bin/python3 "$f" >/dev/null 2>&1; done
"""

import hashlib
import hmac
import json
import os
import re
import sys
import tempfile
import time
import urllib.request

PLACEHOLDER = "const AGENT_SECRET = '';"
SECRET_LINE = "const AGENT_SECRET = '{}';"
TIME_WINDOW = 300
LOG_MAX_LINES = 2000
USER_AGENT = "dso-watchdog-py/1.1"
POOL_COUNT_DEFAULT = 5
RUN_GATE_SECONDS = 50      # dedupe: proses domain maks sekali per menit
HB_INTERVAL_DEFAULT = 180  # kirim heartbeat maks tiap 3 menit (hemat kuota)


def pool_paths(me_path, count):
    """Nama pool TETAP: <basis>.py + <basis>2..5.py. Bila file yang berjalan
    adalah salinan bernomor (mis. dso_watchdog3.py), basis tetap dso_watchdog
    — sehingga salinan mana pun yang dieksekusi, pool yang sama dipulihkan."""
    d, base = os.path.split(me_path)
    stem, ext = os.path.splitext(base)
    m = re.search(r"^(.*?)[2-9]$", stem)
    if m:
        stem = m.group(1)
    return [os.path.join(d, stem + ext)] + [os.path.join(d, "{}{}{}".format(stem, i, ext)) for i in range(2, count + 1)]


def ensure_pool(me_path, count, log_file):
    """Salinan yang hilang ditulis ulang dari file ini — pool selalu 5."""
    with open(me_path, "rb") as f:
        src = f.read()
    restored = 0
    for p in pool_paths(me_path, count):
        if p == me_path:
            continue
        if not os.path.isfile(p):
            try:
                atomic_write(p, src, 0o755)
                restored += 1
                log("POOL diregenerasi: {}".format(os.path.basename(p)), log_file)
            except OSError as e:
                log("ERROR gagal regenerasi pool {}: {}".format(p, e), log_file)
    return restored


def gate_ok(gate_file, seconds):
    """True bila sudah lewat `seconds` sejak tanda terakhir (atau pertama kali)."""
    now = time.time()
    try:
        if os.path.isfile(gate_file):
            with open(gate_file, "r", encoding="utf-8") as f:
                last = float(f.read().strip() or 0)
            if now - last < seconds:
                return False
    except (OSError, ValueError):
        pass
    try:
        atomic_write(gate_file, str(now).encode(), 0o600)
    except OSError:
        pass
    return True


def log(msg, log_file):
    line = time.strftime("%Y-%m-%dT%H:%M:%S") + " " + msg
    try:
        old = []
        if os.path.isfile(log_file):
            with open(log_file, "r", encoding="utf-8", errors="replace") as f:
                old = f.read().splitlines()
        old.append(line)
        if len(old) > LOG_MAX_LINES:
            old = old[-LOG_MAX_LINES:]
        with open(log_file, "w", encoding="utf-8") as f:
            f.write("\n".join(old) + "\n")
    except OSError:
        pass
    print(line)


def sha256_bytes(data):
    return hashlib.sha256(data).hexdigest()


def fetch(url, timeout=20):
    """GET sumber template (mendukung file:// untuk pengujian)."""
    if url.startswith("file://"):
        path_part = url[7:]
        # Windows: file:///C:/... => '/C:/...' — buang slash di depan drive.
        if re.match(r"^/[A-Za-z]:[/_]", path_part):
            path_part = path_part[1:]
        path_part = path_part.replace("/", os.sep)
        with open(path_part, "rb") as f:
            return f.read()
    req = urllib.request.Request(url, headers={"User-Agent": USER_AGENT})
    with urllib.request.urlopen(req, timeout=timeout) as r:
        return r.read()


def post_signed(url, secret, canonical_body, timeout=15, payload=None):
    """POST heartbeat dengan header HMAC ala agent (X-DSO-*)."""
    ts = str(int(time.time()))
    canonical = "hb|" + canonical_body + "|" + ts
    sig = hmac.new(secret.encode(), canonical.encode(), hashlib.sha256).hexdigest()
    body = json.dumps(payload or {}).encode()
    req = urllib.request.Request(
        url,
        data=body,
        method="POST",
        headers={
            "User-Agent": USER_AGENT,
            "Content-Type": "application/json",
            "X-DSO-Timestamp": ts,
            "X-DSO-Signature": sig,
        },
    )
    with urllib.request.urlopen(req, timeout=timeout) as r:
        return json.loads(r.read().decode("utf-8", "replace"))


def atomic_write(path, data_bytes, mode=0o644):
    """Tulis file secara atomik (tempfile di direktori yang sama + rename)."""
    d = os.path.dirname(path) or "."
    fd, tmp = tempfile.mkstemp(prefix=".dso-tmp-", dir=d)
    try:
        with os.fdopen(fd, "wb") as f:
            f.write(data_bytes)
        os.chmod(tmp, mode)
        os.replace(tmp, path)
    except BaseException:
        try:
            os.unlink(tmp)
        except OSError:
            pass
        raise


def expected_content(cfg):
    """Isi agent yang seharusna: template resmi + secret domain."""
    tpl = fetch(cfg["template_url"])
    text = tpl.decode("utf-8", "replace")
    pin = cfg.get("pin_sha256") or ""
    if pin:
        h = sha256_bytes(tpl)
        if h.lower() != pin.lower():
            raise RuntimeError("hash template ({}) tidak cocok pin_sha256".format(h[:16]))
    if PLACEHOLDER not in text or "AGENT_VERSION" not in text:
        raise RuntimeError("template tidak dikenali (placeholder/versi tidak ada)")
    return text.replace(PLACEHOLDER, SECRET_LINE.format(cfg["secret"]))


def process_domain(cfg, log_file, hb_gate_file=None, hb_extra=None):
    name = cfg.get("name") or cfg.get("domain_id", "?")
    agent_path = cfg["agent_path"]

    # 1. Pastikan agent ada & identik dengan sumber kanonik.
    try:
        expected = expected_content(cfg).encode("utf-8")
    except Exception as e:  # gagal ambil sumber => jangan sentuh agent lama
        log("WARN [{}] gagal membangun konten yang diharapkan: {} — agent dibiarkan".format(name, e), log_file)
        return False

    current = None
    try:
        if os.path.isfile(agent_path):
            with open(agent_path, "rb") as f:
                current = f.read()
    except OSError:
        current = None

    if current == expected:
        log("OK [{}] agent utuh".format(name), log_file)
    else:
        reason = "hilang" if current is None else (
            "berubah (sha256={})".format(sha256_bytes(current)[:16]))
        try:
            atomic_write(agent_path, expected)
            log("RESTORED [{}] agent {} — dibangun ulang dari sumber kanonik (sha256={})".format(
                name, reason, sha256_bytes(expected)[:16]), log_file)
        except OSError as e:
            log("ERROR [{}] gagal menulis {}: {}".format(name, agent_path, e), log_file)
            return False

    # 2. Heartbeat ke Worker (interval dibatasi — hemat kuota KV Worker).
    interval = int(cfg.get("hb_interval") or HB_INTERVAL_DEFAULT)
    if (cfg.get("heartbeat") and cfg.get("worker_base") and cfg.get("domain_id")
            and gate_ok(hb_gate_file, interval)):
        url = cfg["worker_base"].rstrip("/") + "/hb/" + cfg["domain_id"]
        try:
            r = post_signed(url, cfg["secret"], cfg["domain_id"], payload=hb_extra or {})
            if not r.get("ok"):
                log("WARN [{}] heartbeat ditolak: {}".format(name, r), log_file)
        except Exception as e:
            log("WARN [{}] heartbeat gagal: {}".format(name, e), log_file)
    return True


def main():
    argv = sys.argv[1:]
    me_path = os.path.abspath(__file__)
    base = os.path.dirname(me_path)
    config_path = os.path.join(base, "dso_watchdog.json")
    if "--config" in argv:
        config_path = argv[argv.index("--config") + 1]

    try:
        with open(config_path, "r", encoding="utf-8") as f:
            config = json.load(f)
    except Exception as e:
        print("ERROR membaca config {}: {}".format(config_path, e))
        return 1

    domains = config.get("domains") or [config]
    log_file = config.get("log_file") or os.path.join(base, "dso_watchdog.log")
    pool_count = int(config.get("pool_count") or POOL_COUNT_DEFAULT)

    if config.get("require_private_config", True):
        try:
            mode = os.stat(config_path).st_mode & 0o777
            if mode != 0o600:
                print("PERINGATAN: config {} permission {:o} — disarankan chmod 600".format(config_path, mode))
        except OSError:
            pass

    if "--doctor" in argv:
        print("dso_watchdog 1.1.0 — doctor")
        print("config       :", config_path)
        print("pool         :", len(pool_paths(me_path, pool_count)), "salinan di", base)
        print("jumlah domain:", len(domains))
        for cfg in domains:
            print(" - {:<24} agent: {}".format(cfg.get("name", "?"), cfg.get("agent_path", "?")))
            try:
                exp = expected_content(cfg)
                print("   sumber kanonik OK, sha256={}".format(sha256_bytes(exp.encode())[:16]))
            except Exception as e:
                print("   SUMBER GAGAL:", e)
        print('cron: * * * * * for f in {}/*.py; do [ -f "$f" ] && /usr/bin/python3 "$f" >/dev/null 2>&1; done'.format(base))
        return 0

    # 0. Pool saling-regenerasi: selalu jalan di SETIAP eksekusi (murah).
    ensure_pool(me_path, pool_count, log_file)
    pool_alive = sum(1 for p in pool_paths(me_path, pool_count) if os.path.isfile(p))

    # Dedupe: bila salinan lain baru saja memproses domain, cukup pool check.
    run_gate = os.path.join(base, "dso_last_run")
    run_gate_s = float(config.get("run_gate_seconds", RUN_GATE_SECONDS))
    if not gate_ok(run_gate, run_gate_s):
        return 0

    ok = True
    for cfg in domains:
        try:
            hb_gate = os.path.join(base, "dso_last_hb_" + str(cfg.get("domain_id", "x")))
            if not process_domain(cfg, log_file, hb_gate,
                                  hb_extra={"pool": pool_alive, "dir": base}):
                ok = False
        except Exception as e:
            log("ERROR [{}] tak terduga: {}".format(cfg.get("name", "?"), e), log_file)
            ok = False
    return 0 if ok else 1


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