from datetime import timedelta

from django.conf import settings
from django.db import transaction
from django.test import TransactionTestCase, override_settings
from django.utils import timezone

from api.models import IntegrationOutbox, IntegrationPublisherState, TtLocation, TtStaff
from api.services.integration.readiness import phase1_readiness
from api.services.integration.tracking import mutation_tracking


SNAPSHOT_SETTINGS = {
    **settings.RESOURCE_BOOKING_INTEGRATION,
    "SOURCE_SCOPE": "campus-a",
    "DEPLOYMENT_ID": "campus-a",
    "SNAPSHOT_SERVICE_TOKEN": "snapshot-test-token",
    "SNAPSHOT_MAX_PAGE_SIZE": 2,
    "SNAPSHOT_MAX_BYTES": 2097152,
}


@override_settings(RESOURCE_BOOKING_INTEGRATION=SNAPSHOT_SETTINGS)
class ResourceBookingSnapshotTests(TransactionTestCase):
    def get_snapshot(self, **query):
        return self.client.get(
            "/api/integration/resource-booking/v1/snapshot",
            {"source_scope": "campus-a", **query},
            HTTP_AUTHORIZATION="Bearer snapshot-test-token",
        )

    def test_snapshot_requires_service_authentication_and_exact_scope(self):
        unauthorized = self.client.get(
            "/api/integration/resource-booking/v1/snapshot",
            {"source_scope": "campus-a"},
        )
        self.assertEqual(unauthorized.status_code, 401)
        wrong_scope = self.client.get(
            "/api/integration/resource-booking/v1/snapshot",
            {"source_scope": "other"},
            HTTP_AUTHORIZATION="Bearer snapshot-test-token",
        )
        self.assertEqual(wrong_scope.status_code, 404)

    def test_fixed_watermark_returns_latest_canonical_state_at_that_boundary(self):
        with transaction.atomic():
            with mutation_tracking():
                staff = TtStaff.objects.create(code="S1", name="Before", status=1)
        first_watermark = IntegrationOutbox.objects.get().transport_sequence
        with transaction.atomic():
            with mutation_tracking():
                staff.name = "After"
                staff.save(update_fields=["name", "updated_at"])

        response = self.get_snapshot(snapshot_watermark=first_watermark)
        self.assertEqual(response.status_code, 200)
        body = response.json()
        self.assertEqual(body["snapshot_watermark"], str(first_watermark))
        self.assertTrue(body["complete"])
        self.assertEqual(len(body["envelopes"]), 1)
        self.assertEqual(
            body["envelopes"][0]["events"][0]["committed_state"]["replacement"]["current"]["name"],
            "Before",
        )
        self.assertEqual(body["envelopes"][0]["source_sequence"], first_watermark)
        self.assertEqual(body["envelopes"][0]["transaction_count"], 1)

    def test_snapshot_pages_without_duplication_under_one_watermark(self):
        for code, name, model in (
            ("S1", "One", TtStaff),
            ("S2", "Two", TtStaff),
            ("L1", "Room", TtLocation),
        ):
            with transaction.atomic():
                with mutation_tracking():
                    model.objects.create(code=code, name=name, status=1)
        first = self.get_snapshot(page_size=2)
        self.assertEqual(first.status_code, 200)
        first_body = first.json()
        self.assertFalse(first_body["complete"])
        second = self.get_snapshot(
            page_size=2,
            cursor=first_body["next_cursor"],
            snapshot_watermark=first_body["snapshot_watermark"],
        )
        self.assertEqual(second.status_code, 200)
        second_body = second.json()
        self.assertTrue(second_body["complete"])
        event_ids = [
            envelope["event_id"]
            for envelope in first_body["envelopes"] + second_body["envelopes"]
        ]
        self.assertEqual(len(event_ids), 3)
        self.assertEqual(len(set(event_ids)), 3)
        sequences = [
            envelope["source_sequence"]
            for envelope in first_body["envelopes"] + second_body["envelopes"]
        ]
        self.assertEqual(sequences, sorted(set(sequences)))
        self.assertTrue(all(sequence <= int(first_body["snapshot_watermark"]) for sequence in sequences))

    def test_snapshot_groups_latest_members_by_originating_transaction_sequence(self):
        with transaction.atomic():
            with mutation_tracking():
                TtStaff.objects.create(code="S1", name="One", status=1)
                TtLocation.objects.create(code="L1", name="Room", status=1)
        response = self.get_snapshot(page_size=2)
        self.assertEqual(response.status_code, 200)
        envelopes = response.json()["envelopes"]
        self.assertEqual(len(envelopes), 1)
        self.assertEqual(envelopes[0]["transaction_count"], 2)
        self.assertTrue(envelopes[0]["complete"])
        self.assertEqual(len(envelopes[0]["events"]), 2)
        self.assertEqual(
            set(envelopes[0]),
            {
                "event_id",
                "schema_version",
                "event_type",
                "source_system",
                "source_scope",
                "source_sequence",
                "ordering_key",
                "source_transaction_id",
                "transaction_count",
                "complete",
                "committed_at",
                "events_hash",
                "events",
            },
        )

    def test_health_is_safe_while_publication_remains_disabled(self):
        response = self.client.get(
            "/api/integration/resource-booking/v1/health",
            HTTP_AUTHORIZATION="Bearer snapshot-test-token",
        )
        self.assertEqual(response.status_code, 200)
        body = response.json()
        self.assertFalse(body["publish_enabled"])
        self.assertFalse(body["reverse_delivery_enabled"])
        self.assertNotIn("snapshot-test-token", str(body))

    def test_enabled_preflight_allows_supervisor_to_replace_stale_publisher(self):
        enabled = {
            **SNAPSHOT_SETTINGS,
            "PUBLISH_ENABLED": True,
            "ACTIVATION_APPROVED": True,
            "TRANSPORT": "kafka",
            "KAFKA_TOPIC": "phase1-test",
            "KAFKA_BOOTSTRAP_SERVERS": "broker.invalid:9092",
        }
        IntegrationPublisherState.objects.create(
            source_scope="campus-a",
            status="running",
            heartbeat_at=timezone.now()
            - timedelta(seconds=enabled["PUBLISHER_HEARTBEAT_TIMEOUT_SECONDS"] * 2),
        )
        with override_settings(RESOURCE_BOOKING_INTEGRATION=enabled):
            result = phase1_readiness(require_enabled=True)
        self.assertTrue(result["ready"])
        self.assertFalse(result["publisher_live"])
