import hashlib
import importlib.util
from io import StringIO
import json
import os
from pathlib import Path
from types import SimpleNamespace
from unittest import mock, TestCase


ROOT = Path(__file__).resolve().parents[2]
SCRIPT = ROOT / "deploy" / "phase1_kafka_staging_diagnostic.py"
SPEC = importlib.util.spec_from_file_location("phase1_kafka_diagnostic", SCRIPT)
MODULE = importlib.util.module_from_spec(SPEC)
SPEC.loader.exec_module(MODULE)


class Phase1KafkaStagingDiagnosticTests(TestCase):
    def test_host_fingerprint_emits_only_sha256(self):
        host = "private-staging-host.example.internal"
        output = StringIO()
        with mock.patch.dict(os.environ, {"TT_STAGING_HOST_VALUE": host}):
            with mock.patch("sys.stdout", output):
                self.assertEqual(MODULE.host_fingerprint(), 0)
        result = json.loads(output.getvalue())
        self.assertEqual(
            result,
            {"staging_host_sha256": hashlib.sha256(host.encode()).hexdigest()},
        )
        self.assertNotIn(host, output.getvalue())

    def test_host_classification_never_returns_input(self):
        cases = {
            None: "absent",
            "localhost": "loopback_hostname",
            "127.0.0.1": "loopback_address",
            "10.20.30.40": "private_address",
            "8.8.8.8": "public_address",
            "secret-broker.internal": "hostname",
        }
        for value, expected in cases.items():
            with self.subTest(value=value):
                classification = MODULE.classify_host(value)
                self.assertEqual(classification, expected)
                if value:
                    self.assertNotEqual(classification, value)

    def test_metadata_summary_reports_counts_and_expected_topic_only(self):
        topic_metadata = SimpleNamespace(error=None)
        metadata = SimpleNamespace(
            brokers={1: object(), 2: object()},
            topics={
                MODULE.EXPECTED_TOPIC: topic_metadata,
                "private.topic.name": topic_metadata,
            },
        )
        client = mock.Mock()
        client.list_topics.return_value = metadata
        with mock.patch(
            "confluent_kafka.admin.AdminClient", return_value=client
        ) as admin_client:
            result = MODULE.kafka_metadata_summary(True)

        self.assertEqual(result["broker_metadata_classification"], "available")
        self.assertEqual(result["broker_count"], 2)
        self.assertEqual(result["metadata_topic_count"], 2)
        self.assertTrue(result["expected_topic_exists"])
        self.assertNotIn("private.topic.name", json.dumps(result))
        configuration = admin_client.call_args.args[0]
        self.assertEqual(
            set(configuration),
            {"bootstrap.servers", "socket.timeout.ms", "log_level"},
        )
        self.assertNotIn("sasl", json.dumps(configuration).lower())

    def test_metadata_failure_is_classified_without_error_or_endpoint(self):
        client = mock.Mock()
        client.list_topics.side_effect = RuntimeError(
            "secret-broker.example:9092 credential payload-value"
        )
        with mock.patch("confluent_kafka.admin.AdminClient", return_value=client):
            result = MODULE.kafka_metadata_summary(True)
        encoded = json.dumps(result)
        self.assertEqual(
            result["broker_metadata_classification"], "metadata_unavailable"
        )
        self.assertNotIn("secret-broker", encoded)
        self.assertNotIn("credential", encoded)
        self.assertNotIn("payload-value", encoded)

    def test_remote_report_contains_classifications_and_default_off_flags_only(self):
        settings = SimpleNamespace(
            KAFKA_CONFIG={"HOST": "legacy.secret-broker.internal"},
            RESOURCE_BOOKING_INTEGRATION={
                "CAPTURE_ENABLED": True,
                "PUBLISH_ENABLED": False,
                "ACTIVATION_APPROVED": False,
                "REVERSE_DELIVERY_ENABLED": False,
                "PHASE": "phase1",
                "TRANSPORT": "disabled",
                "KAFKA_BOOTSTRAP_SERVERS": "private-broker.internal:9092",
                "KAFKA_TOPIC": MODULE.EXPECTED_TOPIC,
            },
        )
        metadata = {
            "broker_metadata_classification": "available",
            "broker_count": 1,
            "metadata_topic_count": 3,
            "expected_topic_exists": True,
        }
        with mock.patch.object(MODULE, "tcp_accepts_loopback", return_value=True):
            with mock.patch.object(
                MODULE, "kafka_metadata_summary", return_value=metadata
            ):
                report = MODULE.build_remote_report(settings, publisher_stopped=True)
        encoded = json.dumps(report, sort_keys=True)
        self.assertTrue(MODULE.report_is_safe(report))
        self.assertEqual(report["legacy_kafka_host_classification"], "hostname")
        self.assertFalse(report["bootstrap_requested"])
        self.assertFalse(report["projection_revision_requested"])
        self.assertNotIn("legacy.secret-broker", encoded)
        self.assertNotIn("private-broker", encoded)

    def test_unsafe_runtime_flags_fail_the_diagnostic_gate(self):
        report = {
            "publisher_stopped": True,
            "rb_publish_enabled": True,
            "rb_activation_approved": False,
            "rb_reverse_delivery_enabled": False,
            "bootstrap_requested": False,
            "projection_revision_requested": False,
        }
        self.assertFalse(MODULE.report_is_safe(report))

    def test_workflow_uses_exact_host_secret_without_transport_secrets(self):
        workflow = (ROOT / ".github" / "workflows" / "deploy.yml").read_text()
        self.assertIn("workflow_dispatch:", workflow)
        diagnostic_job = workflow.split("  diagnose-staging-kafka:\n", 1)[1].split(
            "\n  phase1-e2e:", 1
        )[0]
        self.assertEqual(diagnostic_job.count("${{ secrets.HOST }}"), 2)
        self.assertIn("inputs.kafka_diagnostic", diagnostic_job)
        self.assertIn("host-fingerprint", diagnostic_job)
        self.assertIn("phase1_attest_no_publisher_processes", diagnostic_job)
        self.assertIn("TT_PHASE1_DIAGNOSTIC_PUBLISHER_STOPPED=True", diagnostic_job)
        self.assertNotIn("RB_INTEGRATION_KAFKA_BOOTSTRAP_SERVERS", diagnostic_job)
        self.assertNotIn("RB_INTEGRATION_KAFKA_SASL", diagnostic_job)
        self.assertNotIn("secrets.RB_INTEGRATION", diagnostic_job)
        self.assertNotIn("KAFKA_HOST", diagnostic_job)
        self.assertNotIn("bootstrap_resource_booking_resources", diagnostic_job)
        self.assertNotIn("revise_resource_booking_resource_projection", diagnostic_job)
        self.assertNotIn("publish_resource_booking_outbox --", diagnostic_job)
