from collections import Counter, defaultdict
import datetime
import hashlib
import json
from zoneinfo import ZoneInfo

from django.conf import settings
from django.db import migrations, models


def backfill_transport_contract(apps, schema_editor):
    Outbox = apps.get_model("api", "IntegrationOutbox")
    AggregateVersion = apps.get_model("api", "IntegrationAggregateVersion")
    Cursor = apps.get_model("api", "IntegrationTransportCursor")

    events = list(Outbox.objects.all().order_by("created_at", "id"))
    scopes = {}
    counts = Counter()
    for event in events:
        payload = event.payload or {}
        deployment_id = (payload.get("aggregate") or {}).get("deployment_id")
        # Pre-v2 aggregate IDs were deployment-prefixed numeric primary keys.
        # Prefer the canonical payload fact so deployment IDs containing ':' are
        # not truncated; rsplit is the safe legacy fallback.
        scope = str(deployment_id or event.aggregate_id.rsplit(":", 1)[0] or "default")
        scopes[event.pk] = scope
        counts[(scope, event.change_set_id)] += 1

    sequences = defaultdict(int)
    change_sequences = {}
    indexes = defaultdict(int)
    aggregate_versions = defaultdict(int)
    for event in events:
        scope = scopes[event.pk]
        key = (scope, event.change_set_id)
        if key not in change_sequences:
            sequences[scope] += 1
            change_sequences[key] = sequences[scope]
        indexes[key] += 1
        event.source_scope = scope
        event.transport_sequence = change_sequences[key]
        event.transaction_index = indexes[key]
        event.transaction_count = counts[key]
        event.transaction_finalized = True
        event.ordering_key = scope
        payload = dict(event.payload or {})
        payload["schema_version"] = "2"
        payload["ordering_key"] = scope
        payload["transport"] = {
            "source_scope": scope,
            "sequence": event.transport_sequence,
            "transaction_id": str(event.change_set_id),
            "transaction_index": event.transaction_index,
            "transaction_count": event.transaction_count,
            "member_only": True,
        }
        committed_at = event.created_at
        if committed_at.tzinfo is None:
            committed_at = committed_at.replace(tzinfo=ZoneInfo(settings.TIME_ZONE))
        payload["committed_at"] = committed_at.astimezone(datetime.UTC).isoformat().replace(
            "+00:00", "Z"
        )
        committed_state = payload.get("committed_state", {})
        payload["committed_state_hash"] = hashlib.sha256(
            json.dumps(
                committed_state,
                sort_keys=True,
                separators=(",", ":"),
                ensure_ascii=False,
                allow_nan=False,
                default=str,
            ).encode("utf-8")
        ).hexdigest()
        event.schema_version = "2"
        event.payload = payload
        aggregate_key = (event.aggregate_type, event.aggregate_id)
        aggregate_versions[aggregate_key] = max(
            aggregate_versions[aggregate_key], event.event_version
        )

    if events:
        Outbox.objects.bulk_update(
            events,
            [
                "source_scope",
                "transport_sequence",
                "transaction_index",
                "transaction_count",
                "transaction_finalized",
                "ordering_key",
                "payload",
                "schema_version",
            ],
            batch_size=500,
        )
    for scope, sequence in sequences.items():
        Cursor.objects.update_or_create(
            source_scope=scope,
            defaults={"last_sequence": sequence},
        )
    # The outbox is the immutable version evidence. Repair a missing/stale cursor
    # before new writers can allocate a colliding aggregate version (for example
    # after a legacy restore/import that retained events but omitted cursor rows).
    for (aggregate_type, aggregate_id), version in aggregate_versions.items():
        cursor, created = AggregateVersion.objects.get_or_create(
            aggregate_type=aggregate_type,
            aggregate_id=aggregate_id,
            defaults={"version": version},
        )
        if not created and cursor.version < version:
            cursor.version = version
            cursor.save(update_fields=["version", "updated_at"])


class Migration(migrations.Migration):
    dependencies = [("api", "0100_integration_outbox_and_engine_receipt")]

    operations = [
        migrations.CreateModel(
            name="IntegrationTransportCursor",
            fields=[
                (
                    "id",
                    models.BigAutoField(
                        auto_created=True,
                        primary_key=True,
                        serialize=False,
                        verbose_name="ID",
                    ),
                ),
                ("source_scope", models.CharField(max_length=190, unique=True)),
                ("last_sequence", models.PositiveBigIntegerField(default=0)),
                ("updated_at", models.DateTimeField(auto_now=True)),
            ],
            options={"db_table": "integration_transport_cursor"},
        ),
        migrations.CreateModel(
            name="IntegrationPublisherState",
            fields=[
                (
                    "id",
                    models.BigAutoField(
                        auto_created=True,
                        primary_key=True,
                        serialize=False,
                        verbose_name="ID",
                    ),
                ),
                ("source_scope", models.CharField(max_length=190, unique=True)),
                ("status", models.CharField(default="starting", max_length=32)),
                ("heartbeat_at", models.DateTimeField(blank=True, null=True)),
                ("last_success_at", models.DateTimeField(blank=True, null=True)),
                ("last_published_sequence", models.PositiveBigIntegerField(default=0)),
                ("last_error", models.TextField(blank=True, null=True)),
                ("updated_at", models.DateTimeField(auto_now=True)),
            ],
            options={"db_table": "integration_publisher_state"},
        ),
        migrations.AddField(
            model_name="integrationoutbox",
            name="source_scope",
            field=models.CharField(db_index=True, default="default", max_length=190),
        ),
        migrations.AlterField(
            model_name="integrationoutbox",
            name="schema_version",
            field=models.CharField(default="2", max_length=32),
        ),
        migrations.AddField(
            model_name="integrationoutbox",
            name="transport_sequence",
            field=models.PositiveBigIntegerField(db_index=True, default=0),
        ),
        migrations.AddField(
            model_name="integrationoutbox",
            name="transaction_index",
            field=models.PositiveIntegerField(default=1),
        ),
        migrations.AddField(
            model_name="integrationoutbox",
            name="transaction_count",
            field=models.PositiveIntegerField(default=1),
        ),
        migrations.AddField(
            model_name="integrationoutbox",
            name="transaction_finalized",
            field=models.BooleanField(db_index=True, default=True),
        ),
        migrations.RunPython(backfill_transport_contract, migrations.RunPython.noop),
        migrations.AddConstraint(
            model_name="integrationoutbox",
            constraint=models.UniqueConstraint(
                fields=("source_scope", "change_set_id", "transaction_index"),
                name="uq_outbox_scope_change_index",
            ),
        ),
        migrations.AddIndex(
            model_name="integrationoutbox",
            index=models.Index(
                fields=["source_scope", "transport_sequence", "status"],
                name="ix_outbox_transport_order",
            ),
        ),
    ]
