import ast
from pathlib import Path

from django.test import SimpleTestCase


ROOT = Path(__file__).resolve().parents[2]


class IntegrationWriterGuardTests(SimpleTestCase):
    def test_legacy_consumer_has_no_authoritative_model_writes(self):
        for name in ("schedule_response.py", "tt_response-backup.py"):
            source = (ROOT / "kafka_consumer" / name).read_text()
            self.assertNotIn("TtActivity.objects", source)
            self.assertIn("consume_forever", source)

    def test_simulation_routes_to_atomic_engine_apply(self):
        source = (ROOT / "api" / "views" / "admin" / "simulate_schedule_response.py").read_text()
        self.assertIn("from kafka_consumer.tt_response import schedule", source)
        self.assertNotIn("bulk_update", source)

    def test_expired_management_writers_fail_closed(self):
        for name in (
            "20260705-manual-schedule-activity.py",
            "import-manual-schedule-activity.py",
        ):
            source = (ROOT / "api" / "management" / "commands" / name).read_text()
            tree = ast.parse(source)
            calls = [node for node in ast.walk(tree) if isinstance(node, ast.Call)]
            self.assertTrue(any(getattr(call.func, "id", None) == "CommandError" for call in calls))
            self.assertNotIn("TtActivity", source)

    def test_preschedule_handlers_do_not_append_integration_events(self):
        source = (ROOT / "kafka_consumer" / "tt_response.py").read_text()
        start = source.index("def preschedule(")
        end = source.index("def preschedule_details(")
        self.assertNotIn("append_outbox_event", source[start:end])
        self.assertNotIn("mutation_tracking", source[start:end])

    def test_response_consumer_is_manual_commit_and_age_is_not_a_discard(self):
        source = (ROOT / "kafka_consumer" / "tt_response.py").read_text()
        self.assertIn("'enable.auto.commit': False", source)
        self.assertIn("'enable.auto.offset.store': False", source)
        self.assertIn("active_consumer.store_offsets(msg)", source)
        self.assertIn("active_consumer.commit(asynchronous=False)", source)
        self.assertIn("quarantine_engine_response", source)
        self.assertIn("active_consumer.pause(assigned)", source)
        self.assertIn("active_consumer.resume(assigned)", source)
        late_block = source[source.index("if timezone.now() > kafka_log.request_at") :]
        self.assertIn("late_engine_response", late_block[:500])
        self.assertNotIn("continue", late_block[:500])
