#!/usr/bin/env python3
import argparse
from contextlib import contextmanager
import gc
import hashlib
import ipaddress
import json
import os
from pathlib import Path
import socket
import sys


EXPECTED_TOPIC = "timetabler.activity.events"
LOOPBACK_BROKER = ("127.0.0.1", 9092)


def classify_host(value: str | None) -> str:
    if not value:
        return "absent"
    candidate = value.strip().lower()
    if candidate == "localhost":
        return "loopback_hostname"
    try:
        address = ipaddress.ip_address(candidate)
    except ValueError:
        return "hostname"
    if address.is_loopback:
        return "loopback_address"
    if address.is_private:
        return "private_address"
    if address.is_link_local:
        return "link_local_address"
    if address.is_multicast:
        return "multicast_address"
    if address.is_unspecified:
        return "unspecified_address"
    if address.is_reserved:
        return "reserved_address"
    return "public_address"


def tcp_accepts_loopback(timeout_seconds: float = 2.0) -> bool:
    try:
        with socket.create_connection(LOOPBACK_BROKER, timeout=timeout_seconds):
            return True
    except OSError:
        return False


@contextmanager
def suppress_process_stderr():
    saved_stderr = os.dup(2)
    null_stderr = os.open(os.devnull, os.O_WRONLY)
    try:
        os.dup2(null_stderr, 2)
        yield
    finally:
        os.dup2(saved_stderr, 2)
        os.close(null_stderr)
        os.close(saved_stderr)


def kafka_metadata_summary(tcp_available: bool) -> dict[str, bool | int | str]:
    result: dict[str, bool | int | str] = {
        "broker_metadata_classification": "skipped_tcp_unavailable",
        "broker_count": 0,
        "metadata_topic_count": 0,
        "expected_topic_exists": False,
    }
    if not tcp_available:
        return result
    try:
        from confluent_kafka.admin import AdminClient
    except ImportError:
        result["broker_metadata_classification"] = "client_unavailable"
        return result
    try:
        with suppress_process_stderr():
            client = AdminClient(
                {
                    "bootstrap.servers": "127.0.0.1:9092",
                    "socket.timeout.ms": 2000,
                    "log_level": 0,
                }
            )
            metadata = client.list_topics(timeout=5)
            del client
            gc.collect()
        valid_topics = {
            name
            for name, topic_metadata in metadata.topics.items()
            if topic_metadata.error is None
        }
        result.update(
            {
                "broker_metadata_classification": (
                    "available" if metadata.brokers else "empty"
                ),
                "broker_count": len(metadata.brokers),
                "metadata_topic_count": len(valid_topics),
                "expected_topic_exists": EXPECTED_TOPIC in valid_topics,
            }
        )
    except Exception:
        result["broker_metadata_classification"] = "metadata_unavailable"
    return result


def build_remote_report(settings_module, *, publisher_stopped: bool) -> dict:
    integration = settings_module.RESOURCE_BOOKING_INTEGRATION
    legacy_host = settings_module.KAFKA_CONFIG.get("HOST")
    tcp_available = tcp_accepts_loopback()
    report = {
        "legacy_kafka_host_classification": classify_host(legacy_host),
        "loopback_9092_tcp_accepts": tcp_available,
        "rb_capture_enabled": bool(integration.get("CAPTURE_ENABLED")),
        "rb_publish_enabled": bool(integration.get("PUBLISH_ENABLED")),
        "rb_activation_approved": bool(integration.get("ACTIVATION_APPROVED")),
        "rb_reverse_delivery_enabled": bool(
            integration.get("REVERSE_DELIVERY_ENABLED")
        ),
        "rb_phase_is_phase1": integration.get("PHASE") == "phase1",
        "rb_transport_disabled": integration.get("TRANSPORT") == "disabled",
        "rb_transport_kafka": integration.get("TRANSPORT") == "kafka",
        "rb_kafka_bootstrap_configured": bool(
            integration.get("KAFKA_BOOTSTRAP_SERVERS")
        ),
        "rb_kafka_topic_configured": bool(integration.get("KAFKA_TOPIC")),
        "rb_kafka_topic_is_expected": integration.get("KAFKA_TOPIC") == EXPECTED_TOPIC,
        "bootstrap_requested": False,
        "projection_revision_requested": False,
        "publisher_stopped": publisher_stopped,
    }
    report.update(kafka_metadata_summary(tcp_available))
    return report


def report_is_safe(report: dict) -> bool:
    return bool(
        report["publisher_stopped"]
        and not report["rb_publish_enabled"]
        and not report["rb_activation_approved"]
        and not report["rb_reverse_delivery_enabled"]
        and not report["bootstrap_requested"]
        and not report["projection_revision_requested"]
    )


def host_fingerprint() -> int:
    host = os.environ.get("TT_STAGING_HOST_VALUE")
    if not host:
        print("Staging host secret is unavailable", file=sys.stderr)
        return 2
    fingerprint = hashlib.sha256(host.encode("utf-8")).hexdigest()
    print(json.dumps({"staging_host_sha256": fingerprint}, sort_keys=True))
    return 0


def remote_diagnostic() -> int:
    repository_root = Path(__file__).resolve().parents[1]
    sys.path.insert(0, str(repository_root))
    os.environ.setdefault("DJANGO_SETTINGS_MODULE", "backend.settings")
    from django.conf import settings

    publisher_stopped = (
        os.environ.get("TT_PHASE1_DIAGNOSTIC_PUBLISHER_STOPPED") == "True"
    )
    report = build_remote_report(settings, publisher_stopped=publisher_stopped)
    print(json.dumps(report, sort_keys=True))
    return 0 if report_is_safe(report) else 3


def main() -> int:
    parser = argparse.ArgumentParser()
    parser.add_argument("mode", choices=("host-fingerprint", "remote"))
    args = parser.parse_args()
    if args.mode == "host-fingerprint":
        return host_fingerprint()
    return remote_diagnostic()


if __name__ == "__main__":
    raise SystemExit(main())
