import uuid
import json
import importlib
from concurrent.futures import ThreadPoolExecutor
from datetime import date, time, timedelta
from unittest import mock, skipUnless

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

from api.models import (
    AppliedEngineResponse,
    EngineResponseQuarantine,
    IntegrationOutbox,
    PostCommitDelivery,
    TtAcademicTerm,
    TtActivity,
    TtLocation,
    TtStaff,
    TtStaffResourceMap,
    TtWeek,
)
from api.services.integration.outbox import (
    IntegrationAtomicityError,
    append_outbox_event,
    capture_activity_state,
)
from api.services.integration.publisher import (
    KafkaOutboxTransport,
    claim_change_sets,
    publish_once,
)
from api.services.integration.quarantine import quarantine_engine_response
from api.services.integration.receipts import EngineResponseConflict, claim_engine_response
from api.services.integration.replay import replay_event
from api.services.integration.tracking import mutation_tracking, track_activity_ids, track_resource_queryset


class IntegrationOutboxTests(TransactionTestCase):
    reset_sequences = True

    def test_append_requires_authoritative_transaction(self):
        with self.assertRaises(IntegrationAtomicityError):
            append_outbox_event(
                event_type="timetabler.staff.updated",
                aggregate_type="staff",
                source_id=1,
                payload={"value": "unsafe"},
            )

    def test_resource_mutation_and_event_commit_together_with_monotonic_version(self):
        with transaction.atomic():
            with mutation_tracking(actor_id=7):
                staff = TtStaff.objects.create(code="S1", name="One", status=1)
        first = IntegrationOutbox.objects.get(aggregate_type="staff")
        self.assertEqual(first.event_version, 1)
        self.assertEqual(first.actor_id, 7)
        self.assertEqual(first.payload["committed_state"]["replacement"]["current"]["code"], "S1")

        with transaction.atomic():
            with mutation_tracking(actor_id=8):
                staff.name = "Two"
                staff.save(update_fields=["name", "updated_at"])
        versions = list(
            IntegrationOutbox.objects.filter(aggregate_type="staff").order_by("event_version")
            .values_list("event_version", flat=True)
        )
        self.assertEqual(versions, [1, 2])

    def test_multi_aggregate_change_set_has_commit_order_and_complete_boundaries(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)
        events = list(IntegrationOutbox.objects.order_by("transport_sequence"))
        self.assertEqual([event.transport_sequence for event in events], [1, 1])
        self.assertEqual(len({event.change_set_id for event in events}), 1)
        self.assertEqual([event.transaction_index for event in events], [1, 2])
        self.assertEqual([event.transaction_count for event in events], [2, 2])
        self.assertTrue(all(event.transaction_finalized for event in events))
        self.assertTrue(all(event.ordering_key == "default" for event in events))
        self.assertTrue(all(event.payload["transport"]["member_only"] for event in events))

    def test_composite_bound_failure_rolls_back_every_source_member_and_cursor(self):
        constrained = {
            **settings.RESOURCE_BOOKING_INTEGRATION,
            "MAX_EVENTS_PER_CHANGE_SET": 1,
        }
        with override_settings(RESOURCE_BOOKING_INTEGRATION=constrained):
            with self.assertRaises(IntegrationAtomicityError):
                with transaction.atomic():
                    with mutation_tracking():
                        TtStaff.objects.create(code="S1", name="One", status=1)
                        TtLocation.objects.create(code="L1", name="Room", status=1)
        self.assertFalse(TtStaff.objects.exists())
        self.assertFalse(TtLocation.objects.exists())
        self.assertFalse(IntegrationOutbox.objects.exists())

    def test_failure_rolls_back_mutation_and_outbox(self):
        with self.assertRaises(RuntimeError):
            with transaction.atomic():
                with mutation_tracking():
                    TtLocation.objects.create(code="L1", name="Room", status=1)
                    raise RuntimeError("injected after database stage")
        self.assertFalse(TtLocation.objects.filter(code="L1").exists())
        self.assertFalse(IntegrationOutbox.objects.exists())

    def test_bulk_resource_update_reloads_true_previous_state(self):
        staff = TtStaff.objects.create(code="S1", name="Before", status=1)
        with transaction.atomic():
            with mutation_tracking():
                staff.name = "After"
                track_resource_queryset([staff])
                TtStaff.objects.bulk_update([staff], ["name"])
        event = IntegrationOutbox.objects.get(aggregate_type="staff")
        replacement = event.payload["committed_state"]["replacement"]
        self.assertEqual(replacement["previous"]["name"], "Before")
        self.assertEqual(replacement["current"]["name"], "After")

    def test_activity_replacement_contains_old_and_final_scope(self):
        old_staff = TtStaff.objects.create(code="S1", name="Old", status=1)
        new_staff = TtStaff.objects.create(code="S2", name="New", status=1)
        location = TtLocation.objects.create(code="L1", name="Room", status=1)
        activity = TtActivity.objects.create(
            code="A1", name="Activity", status=1, duration=60, scheduled=1, scheduled_start_slot=4
        )
        activity.staff.add(old_staff)
        activity.location.add(location)

        with transaction.atomic():
            with mutation_tracking(request_id="manual-1"):
                track_activity_ids([activity.id])
                activity.staff.set([new_staff.id])
                activity.scheduled_start_slot = 12
                activity.save(update_fields=["scheduled_start_slot", "updated_at"])

        event = IntegrationOutbox.objects.get(aggregate_type="activity")
        replacement = event.payload["committed_state"]["replacement"]
        self.assertEqual(replacement["previous"]["staff_ids"], [old_staff.id])
        self.assertEqual(replacement["current"]["staff_ids"], [new_staff.id])
        self.assertEqual(replacement["previous"]["scheduled_start_slot"], 4)
        self.assertEqual(replacement["current"]["scheduled_start_slot"], 12)

    def test_advisory_requirements_are_invisible_until_one_final_allocation(self):
        staff = TtStaff.objects.create(code="S1", name="Staff", status=1)
        location = TtLocation.objects.create(code="L1", name="Room", status=1)
        activity = TtActivity.objects.create(
            code="A1", name="Activity", status=1, duration=60, scheduled=0
        )
        canonical_before = capture_activity_state(activity)

        # This mirrors resources/update-requirement: it prepares engine inputs
        # but does not allocate either shared resource.
        TtActivity.objects.filter(pk=activity.id).update(
            staff_requirement=1,
            location_requirement=1,
            staff_requirement_type=TtActivity.STAFF_REQUIREMENT_TYPE["preset"],
            location_requirement_type=TtActivity.LOCATION_REQUIREMENT_TYPE["preset"],
        )
        activity.staff_preset.add(staff)
        activity.location_preset.add(location)
        activity.refresh_from_db()

        self.assertFalse(activity.scheduled)
        self.assertEqual(capture_activity_state(activity), canonical_before)
        self.assertFalse(IntegrationOutbox.objects.exists())

        with transaction.atomic():
            with mutation_tracking(request_id="final-allocation-1"):
                track_activity_ids([activity.id])
                activity.staff.set([staff.id])
                activity.location.set([location.id])
                activity.scheduled = 1
                activity.scheduled_start_slot = 12
                activity.save(
                    update_fields=["scheduled", "scheduled_start_slot", "updated_at"]
                )

        events = list(IntegrationOutbox.objects.all())
        self.assertEqual(len(events), 1)
        event = events[0]
        current = event.payload["committed_state"]["replacement"]["current"]
        self.assertEqual(current["staff_ids"], [staff.id])
        self.assertEqual(current["location_ids"], [location.id])
        self.assertTrue(current["scheduled"])
        self.assertEqual(event.request_id, "final-allocation-1")
        self.assertTrue(event.transaction_finalized)
        self.assertEqual(event.transaction_count, 1)

        change_sets = claim_change_sets()
        self.assertEqual(len(change_sets), 1)
        self.assertEqual(change_sets[0].payload["transaction_count"], 1)
        self.assertTrue(change_sets[0].payload["complete"])
        self.assertEqual(len(change_sets[0].payload["events"]), 1)

    def test_scheduled_activity_contains_deterministic_zoned_absolute_occurrences(self):
        term = TtAcademicTerm.objects.create(
            code="2026-T1",
            name="Term 1",
            start_date=date(2026, 1, 5),
            end_date=date(2026, 4, 5),
            start_time=time(8),
            end_time=time(18),
            status=1,
        )
        week = TtWeek.objects.create(week=1, start_date=date(2026, 1, 5))
        activity = TtActivity.objects.create(
            code="A1",
            name="Activity",
            academic_term=term,
            status=1,
            duration=90,
            scheduled=1,
            scheduled_day=1,
            scheduled_start_time=time(9, 30),
            scheduled_start_slot=67,
        )
        activity.week.add(week)
        with transaction.atomic():
            with mutation_tracking():
                track_activity_ids([activity.id])
                activity.name = "Final Activity"
                activity.save(update_fields=["name", "updated_at"])
        current = IntegrationOutbox.objects.get().payload["committed_state"]["replacement"]["current"]
        self.assertEqual(current["calendar"]["weeks"][0]["week_start_date"], "2026-01-05")
        occurrence = current["occurrences"][0]
        self.assertEqual(occurrence["local_start"], "2026-01-06T09:30:00+08:00")
        self.assertEqual(occurrence["local_end"], "2026-01-06T11:00:00+08:00")
        self.assertEqual(occurrence["utc_start"], "2026-01-06T01:30:00Z")
        self.assertEqual(occurrence["utc_end"], "2026-01-06T03:00:00Z")

    def test_resource_delete_emits_tombstone_and_affected_activity_replacement(self):
        staff = TtStaff.objects.create(code="S1", name="Old", status=1)
        activity = TtActivity.objects.create(
            code="A1", name="Activity", status=1, duration=60, scheduled=1, scheduled_start_slot=4
        )
        activity.staff.add(staff)
        with transaction.atomic():
            with mutation_tracking():
                staff.delete()
        resource_event = IntegrationOutbox.objects.get(aggregate_type="staff")
        activity_event = IntegrationOutbox.objects.get(aggregate_type="activity")
        self.assertTrue(resource_event.payload["committed_state"]["tombstone"])
        self.assertEqual(
            activity_event.payload["committed_state"]["replacement"]["current"]["staff_ids"], []
        )
        change_sets = claim_change_sets()
        self.assertEqual(len(change_sets), 1)
        self.assertEqual(change_sets[0].payload["transaction_count"], 2)
        self.assertEqual(len(change_sets[0].payload["events"]), 2)
        self.assertEqual(
            {event["aggregate"]["aggregate_type"] for event in change_sets[0].payload["events"]},
            {"activity", "staff"},
        )

    def test_matching_engine_redelivery_is_noop_and_conflicting_hash_is_rejected(self):
        response = {"request_id": "engine-1", "status": "success", "schedule": "[]"}
        change_set_id = uuid.uuid4()
        with transaction.atomic():
            first = claim_engine_response(
                request_id="engine-1", response_data=response, change_set_id=change_set_id
            )
            self.assertFalse(first.duplicate)
        with transaction.atomic():
            duplicate = claim_engine_response(
                request_id="engine-1", response_data=response, change_set_id=uuid.uuid4()
            )
            self.assertTrue(duplicate.duplicate)
        with self.assertRaises(EngineResponseConflict):
            with transaction.atomic():
                claim_engine_response(
                    request_id="engine-1",
                    response_data={**response, "status": "failed"},
                    change_set_id=uuid.uuid4(),
                )
        self.assertEqual(AppliedEngineResponse.objects.count(), 1)

    def test_dead_letter_replay_retries_same_transport_identity_and_unblocks_ordering(self):
        with transaction.atomic():
            original = append_outbox_event(
                event_type="timetabler.staff.updated",
                aggregate_type="staff",
                source_id=42,
                payload={"replacement": {"current": {"code": "S42"}}},
            )
        original.status = IntegrationOutbox.Status.DEAD_LETTER
        original.save(update_fields=["status"])
        replay = replay_event(event=original, actor_id=9, reason="adapter outage repaired")
        original.refresh_from_db()
        self.assertEqual(original.status, IntegrationOutbox.Status.RETRY)
        self.assertEqual(replay.pk, original.pk)
        self.assertEqual(replay.event_id, original.event_id)
        self.assertEqual(replay.transport_sequence, original.transport_sequence)
        self.assertEqual(replay.event_version, original.event_version)
        self.assertEqual(replay.replay_requested_by, 9)

    def test_kafka_side_effect_is_durable_before_commit_and_triggered_after_commit(self):
        from backend.kafka import send_request

        with mock.patch("api.services.integration.delivery.deliver_by_id") as deliver:
            with transaction.atomic():
                send_request("legacy-topic", {"session_id": "test", "activity": [1]}, None, "schedule")
                delivery = PostCommitDelivery.objects.get()
                self.assertEqual(delivery.status, PostCommitDelivery.Status.PENDING)
                deliver.assert_not_called()
            deliver.assert_called_once_with(delivery.pk)

    def test_invalid_engine_response_quarantine_is_deduplicated_and_counted(self):
        payload = {"request_id": "bad-1", "schedule": "not-json"}
        first = quarantine_engine_response(payload=payload, headers={"method": "schedule"}, reason="protocol")
        second = quarantine_engine_response(payload=payload, headers={"method": "schedule"}, reason="protocol")
        self.assertEqual(first.pk, second.pk)
        self.assertEqual(EngineResponseQuarantine.objects.count(), 1)
        second.refresh_from_db()
        self.assertEqual(second.occurrences, 2)

    @skipUnless(connection.vendor == "postgresql", "PostgreSQL advisory-lock concurrency")
    def test_concurrent_version_allocation_is_gap_free(self):
        def append(index):
            close_old_connections()
            try:
                with transaction.atomic():
                    append_outbox_event(
                        event_type="timetabler.activity.updated",
                        aggregate_type="activity",
                        source_id=99,
                        payload={"index": index},
                    )
            finally:
                close_old_connections()

        with ThreadPoolExecutor(max_workers=4) as executor:
            list(executor.map(append, range(12)))
        versions = list(
            IntegrationOutbox.objects.filter(aggregate_type="activity")
            .order_by("event_version")
            .values_list("event_version", flat=True)
        )
        self.assertEqual(versions, list(range(1, 13)))


class EngineAtomicApplyTests(TransactionTestCase):
    @classmethod
    def setUpClass(cls):
        super().setUpClass()
        from api.models import TtSetting

        for param, value in {
            "minute_per_slot": "30",
            "slot_per_week": "336",
            "slot_per_day": "48",
        }.items():
            TtSetting.objects.update_or_create(param=param, defaults={"value": value})
        cls.consumer_module = importlib.import_module("kafka_consumer.tt_response")

    def test_engine_result_commits_activity_receipt_and_outbox_once(self):
        activity = TtActivity.objects.create(code="A1", name="One", status=1, duration=60, scheduled=0)
        response = {
            "session_id": "test",
            "request_id": "engine-atomic-1",
            "status": "success",
            "schedule": json.dumps(
                [
                    {
                        "activity": str(activity.id),
                        "teaching_staff": "[]",
                        "location": "[]",
                        "start_slot": "12",
                    }
                ]
            ),
        }
        kafka_log = type("Log", (), {"socket_id": None})()
        first = self.consumer_module.schedule(response, kafka_log)
        activity.refresh_from_db()
        self.assertFalse(first["duplicate"])
        self.assertTrue(activity.scheduled)
        self.assertEqual(activity.scheduled_start_slot, 12)
        self.assertEqual(AppliedEngineResponse.objects.count(), 1)
        self.assertEqual(IntegrationOutbox.objects.filter(aggregate_type="activity").count(), 1)

        duplicate = self.consumer_module.schedule(response, kafka_log)
        self.assertTrue(duplicate["duplicate"])
        self.assertEqual(AppliedEngineResponse.objects.count(), 1)
        self.assertEqual(IntegrationOutbox.objects.filter(aggregate_type="activity").count(), 1)

    def test_engine_result_preserves_initiating_request_correlation(self):
        activity = TtActivity.objects.create(
            code="A-CORRELATION",
            name="Correlation",
            status=1,
            duration=60,
            scheduled=0,
        )
        delivery = PostCommitDelivery.objects.create(
            method="schedule",
            target="timetabler.request",
            payload={"activities": [activity.id]},
            correlation_id="phase1-load-test-load-02-000001",
        )
        response = {
            "session_id": "test",
            "request_id": str(delivery.delivery_id),
            "status": "success",
            "schedule": json.dumps(
                [
                    {
                        "activity": str(activity.id),
                        "teaching_staff": "[]",
                        "location": "[]",
                        "start_slot": "12",
                    }
                ]
            ),
        }

        self.consumer_module.schedule(
            response, type("Log", (), {"socket_id": None})()
        )

        event = IntegrationOutbox.objects.get(aggregate_type="activity")
        self.assertEqual(
            event.payload["correlation_id"],
            "phase1-load-test-load-02-000001",
        )
        self.assertEqual(event.payload["causation_id"], str(delivery.delivery_id))
        self.assertEqual(event.payload["request_id"], str(delivery.delivery_id))

    def test_partial_no_slot_result_does_not_mutate_failed_item(self):
        scheduled = TtActivity.objects.create(code="A1", name="One", status=1, duration=60, scheduled=0)
        no_slot = TtActivity.objects.create(code="A2", name="Two", status=1, duration=60, scheduled=0)
        response = {
            "session_id": "test",
            "request_id": "engine-partial-1",
            "status": "success",
            "schedule": json.dumps(
                [
                    {
                        "activity": str(scheduled.id),
                        "teaching_staff": "[]",
                        "location": "[]",
                        "start_slot": "8",
                    },
                    {
                        "activity": str(no_slot.id),
                        "teaching_staff": "[]",
                        "location": "[]",
                        "start_slot": "",
                    },
                ]
            ),
        }
        self.consumer_module.schedule(response, type("Log", (), {"socket_id": None})())
        scheduled.refresh_from_db()
        no_slot.refresh_from_db()
        self.assertTrue(scheduled.scheduled)
        self.assertFalse(no_slot.scheduled)
        self.assertEqual(
            list(IntegrationOutbox.objects.filter(aggregate_type="activity").values_list("aggregate_id", flat=True)),
            [f"default:{scheduled.id}"],
        )

    def test_bulk_engine_result_produces_one_complete_transport_record(self):
        activities = [
            TtActivity.objects.create(
                code=f"A{index}", name=f"Activity {index}", status=1, duration=60, scheduled=0
            )
            for index in (1, 2)
        ]
        response = {
            "session_id": "test",
            "request_id": "engine-bulk-transport-1",
            "status": "success",
            "schedule": json.dumps(
                [
                    {
                        "activity": str(activity.id),
                        "teaching_staff": "[]",
                        "location": "[]",
                        "start_slot": str(8 + index),
                    }
                    for index, activity in enumerate(activities)
                ]
            ),
        }
        self.consumer_module.schedule(response, type("Log", (), {"socket_id": None})())
        rows = list(IntegrationOutbox.objects.order_by("transaction_index"))
        self.assertEqual(len(rows), 2)
        self.assertEqual(len({row.transport_sequence for row in rows}), 1)
        change_sets = claim_change_sets()
        self.assertEqual(len(change_sets), 1)
        envelope = change_sets[0].payload
        self.assertTrue(envelope["complete"])
        self.assertEqual(envelope["transaction_count"], 2)
        self.assertEqual(len(envelope["events"]), 2)
        self.assertEqual(
            [event["aggregate"]["source_id"] for event in envelope["events"]],
            [str(activity.id) for activity in activities],
        )

    def test_processing_failure_rolls_back_state_receipt_outbox_and_durable_effect(self):
        activity = TtActivity.objects.create(code="A1", name="One", status=1, duration=60, scheduled=0)
        response = {
            "session_id": "test",
            "request_id": "engine-failure-1",
            "status": "success",
            "schedule": json.dumps(
                [
                    {
                        "activity": str(activity.id),
                        "teaching_staff": "[]",
                        "location": "[]",
                        "start_slot": "8",
                    }
                ]
            ),
        }

        def fail_after_mutation(response_data, kafka_log):
            TtActivity.objects.filter(pk=activity.pk).update(scheduled=1, scheduled_start_slot=8)
            from backend.kafka import send_request

            send_request("legacy-topic", {"session_id": "test", "activity": [activity.id]}, None, "schedule")
            raise RuntimeError("injected database-stage failure")

        with mock.patch.object(self.consumer_module, "_apply_schedule", side_effect=fail_after_mutation):
            with self.assertRaises(RuntimeError):
                self.consumer_module.schedule(response, type("Log", (), {"socket_id": None})())
        activity.refresh_from_db()
        self.assertFalse(activity.scheduled)
        self.assertFalse(AppliedEngineResponse.objects.exists())
        self.assertFalse(IntegrationOutbox.objects.exists())
        self.assertFalse(PostCommitDelivery.objects.exists())

    @skipUnless(connection.vendor == "postgresql", "PostgreSQL resource-map SQL")
    def test_reschedule_replaces_old_resource_map_and_event_scope(self):
        week = TtWeek.objects.create(week=1, start_date=date(2026, 1, 5))
        old_staff = TtStaff.objects.create(code="S1", name="Old", status=1)
        new_staff = TtStaff.objects.create(code="S2", name="New", status=1)
        TtStaffResourceMap.objects.create(staff=old_staff, week=week, pattern="0" * 336)
        TtStaffResourceMap.objects.create(staff=new_staff, week=week, pattern="0" * 336)
        activity = TtActivity.objects.create(
            code="A1",
            name="One",
            status=1,
            duration=60,
            slot_required=2,
            scheduled=1,
            scheduled_start_slot=4,
        )
        activity.week.add(week)
        activity.staff.add(old_staff)
        response = {
            "session_id": "test",
            "request_id": "engine-reschedule-1",
            "status": "success",
            "schedule": json.dumps(
                [
                    {
                        "activity": str(activity.id),
                        "teaching_staff": json.dumps([new_staff.id]),
                        "location": "[]",
                        "start_slot": "8",
                    }
                ]
            ),
        }
        with mock.patch.object(self.consumer_module, "bulk_sync_to_redis"):
            self.consumer_module.schedule(response, type("Log", (), {"socket_id": None})())
        activity.refresh_from_db()
        old_map = TtStaffResourceMap.objects.get(staff=old_staff, week=week)
        new_map = TtStaffResourceMap.objects.get(staff=new_staff, week=week)
        self.assertEqual(list(activity.staff.values_list("id", flat=True)), [new_staff.id])
        self.assertEqual(old_map.pattern, "0" * 336)
        self.assertEqual(new_map.pattern[8:10], "11")
        event = IntegrationOutbox.objects.get(aggregate_type="activity")
        replacement = event.payload["committed_state"]["replacement"]
        self.assertEqual(replacement["previous"]["staff_ids"], [old_staff.id])
        self.assertEqual(replacement["current"]["staff_ids"], [new_staff.id])


@override_settings(
    RESOURCE_BOOKING_INTEGRATION={
        "PUBLISH_ENABLED": True,
        "PUBLISH_BATCH_SIZE": 100,
        "MAX_ATTEMPTS": 2,
        "RETRY_BASE_SECONDS": 1,
        "RETRY_MAX_SECONDS": 10,
        "CLAIM_TIMEOUT_SECONDS": 30,
        "HTTP_AUTH_TOKEN": "secret",
        "KAFKA_TOPIC": "tt-phase1-test",
        "KAFKA_BOOTSTRAP_SERVERS": "127.0.0.1:9092",
        "KAFKA_SECURITY_PROTOCOL": "PLAINTEXT",
        "MAX_EVENT_BYTES": 2097152,
        "SNAPSHOT_SERVICE_TOKEN": "snapshot-secret",
    }
)
class PublisherTests(TransactionTestCase):
    def create_event(self, source_id, version, status=IntegrationOutbox.Status.PENDING):
        sequence = IntegrationOutbox.objects.count() + 1
        return IntegrationOutbox.objects.create(
            event_id=uuid.uuid4(),
            event_type="timetabler.activity.updated",
            aggregate_type="activity",
            aggregate_id=f"test:{source_id}",
            event_version=version,
            source_scope="default",
            transport_sequence=sequence,
            transaction_index=1,
            transaction_count=1,
            transaction_finalized=True,
            ordering_key="default",
            change_set_id=uuid.uuid4(),
            payload={"version": version},
            status=status,
        )

    def test_claim_preserves_global_commit_order(self):
        first = self.create_event(1, 1)
        second = self.create_event(1, 2)
        other = self.create_event(2, 1)
        claimed = claim_change_sets()
        self.assertEqual(
            [change_set.events[0].pk for change_set in claimed],
            [first.pk, second.pk, other.pk],
        )

    def test_multi_member_change_set_is_one_transport_publish_and_one_ack_unit(self):
        change_set_id = uuid.uuid4()
        events = []
        for index, aggregate_type in enumerate(("staff", "activity"), start=1):
            events.append(
                IntegrationOutbox.objects.create(
                    event_id=uuid.uuid4(),
                    schema_version="2",
                    event_type=f"timetabler.{aggregate_type}.updated",
                    aggregate_type=aggregate_type,
                    aggregate_id=f"default:{index}",
                    event_version=1,
                    source_scope="default",
                    transport_sequence=1,
                    transaction_index=index,
                    transaction_count=2,
                    transaction_finalized=True,
                    ordering_key="default",
                    change_set_id=change_set_id,
                    payload={
                        "event_id": str(uuid.uuid4()),
                        "schema_version": "2",
                        "event_type": f"timetabler.{aggregate_type}.updated",
                        "aggregate": {
                            "deployment_id": "default",
                            "aggregate_type": aggregate_type,
                            "source_id": str(index),
                        },
                        "committed_at": "2026-08-06T00:00:00Z",
                        "committed_state": {"index": index},
                    },
                )
            )

        published = []

        class RecordingTransport:
            def publish(self, change_set):
                published.append(change_set.payload)

        result = publish_once(transport=RecordingTransport())
        self.assertEqual(result.published, 1)
        self.assertEqual(len(published), 1)
        self.assertEqual(published[0]["source_sequence"], 1)
        self.assertEqual(published[0]["transaction_count"], 2)
        self.assertTrue(published[0]["complete"])
        self.assertEqual(len(published[0]["events"]), 2)
        self.assertEqual(
            set(IntegrationOutbox.objects.values_list("status", flat=True)),
            {IntegrationOutbox.Status.PUBLISHED},
        )

    def test_kafka_transport_emits_one_record_for_one_composite_change_set(self):
        change_set_id = uuid.uuid4()
        for index in (1, 2):
            IntegrationOutbox.objects.create(
                event_id=uuid.uuid4(),
                schema_version="2",
                event_type="timetabler.activity.updated",
                aggregate_type="activity",
                aggregate_id=f"default:{index}",
                event_version=1,
                source_scope="default",
                transport_sequence=1,
                transaction_index=index,
                transaction_count=2,
                transaction_finalized=True,
                ordering_key="default",
                change_set_id=change_set_id,
                payload={
                    "event_id": str(uuid.uuid4()),
                    "schema_version": "2",
                    "event_type": "timetabler.activity.updated",
                    "aggregate": {
                        "deployment_id": "default",
                        "aggregate_type": "activity",
                        "source_id": str(index),
                    },
                    "committed_at": "2026-08-06T00:00:00Z",
                    "committed_state": {"index": index},
                },
            )
        change_set = claim_change_sets()[0]
        fake_producer = mock.Mock()
        fake_producer.flush.return_value = 0
        with mock.patch("api.services.integration.publisher.Producer", return_value=fake_producer):
            KafkaOutboxTransport().publish(change_set)
        fake_producer.produce.assert_called_once()
        args, kwargs = fake_producer.produce.call_args
        envelope = json.loads(args[1])
        self.assertEqual(args[0], "tt-phase1-test")
        self.assertEqual(kwargs["key"], b"default")
        self.assertEqual(envelope["source_sequence"], 1)
        self.assertEqual(envelope["transaction_count"], 2)
        self.assertTrue(envelope["complete"])

    def test_failed_sequence_releases_and_blocks_later_sequences(self):
        first = self.create_event(1, 1)
        second = self.create_event(2, 1)

        class FailureTransport:
            def publish(self, event):
                raise RuntimeError("adapter unavailable")

        publish_once(transport=FailureTransport())
        first.refresh_from_db()
        second.refresh_from_db()
        self.assertEqual(first.status, IntegrationOutbox.Status.RETRY)
        self.assertEqual(second.status, IntegrationOutbox.Status.RETRY)
        self.assertEqual(claim_change_sets(), [])

    def test_publish_ack_and_retry_dead_letter_state(self):
        event = self.create_event(1, 1)

        class SuccessTransport:
            def publish(self, event):
                return None

        result = publish_once(transport=SuccessTransport())
        event.refresh_from_db()
        self.assertEqual(result.published, 1)
        self.assertEqual(event.status, IntegrationOutbox.Status.PUBLISHED)

        failed = self.create_event(2, 1)

        class FailureTransport:
            def publish(self, event):
                raise RuntimeError(
                    "token=secret snapshot=snapshot-secret broker=127.0.0.1:9092 unavailable"
                )

        publish_once(transport=FailureTransport())
        failed.refresh_from_db()
        self.assertEqual(failed.status, IntegrationOutbox.Status.RETRY)
        self.assertNotIn("secret", failed.last_error)
        self.assertNotIn("127.0.0.1:9092", failed.last_error)
        self.assertIn("[REDACTED]", failed.last_error)
        self.assertGreater(failed.next_attempt_at, timezone.now() - timedelta(seconds=1))
        failed.next_attempt_at = timezone.now() - timedelta(seconds=1)
        failed.save(update_fields=["next_attempt_at"])
        publish_once(transport=FailureTransport())
        failed.refresh_from_db()
        self.assertEqual(failed.status, IntegrationOutbox.Status.DEAD_LETTER)
