"""Bridge between the event bus and the workflow runtime.

When a workflow registers a transition with ``trigger="some.event"``, we
subscribe a single dispatcher to that event. On fire, the dispatcher looks
up active instances of any definition whose subject matches the payload
and runs every legal transition.

Payload routing keys (any of):
    * ``workflow_instance_id``       — UUID/str of a specific instance
    * ``subject_type`` + ``object_id`` — locate by generic FK
"""

from __future__ import annotations

from collections.abc import Mapping
from typing import Any

import structlog

from simorgh.apps.events.bus import EventError, subscribe
from simorgh.apps.workflow.registry import (
    WorkflowDefinition,
    workflows_for_event,
)
from simorgh.apps.workflow.registry import (
    register_workflow as _register_definition,
)

_log = structlog.get_logger("simorgh.workflow.triggers")
_WIRED_EVENTS: set[str] = set()


def _route(event_name: str, payload: Mapping[str, Any]) -> None:
    # Lazy imports to avoid touching models during app startup.
    from simorgh.apps.workflow.engine import fire_for_event
    from simorgh.apps.workflow.models import (
        WorkflowInstance,
        WorkflowInstanceStatus,
    )

    instances = _resolve_instances(event_name, payload, WorkflowInstance, WorkflowInstanceStatus)
    if not instances:
        return
    for instance in instances:
        try:
            fire_for_event(instance, event_name, payload=dict(payload))
        except Exception as exc:
            _log.warning(
                "workflow.trigger_routing_failed",
                event_name=event_name,
                instance=str(instance.public_id),
                error=str(exc),
            )


def _resolve_instances(
    event_name: str,
    payload: Mapping[str, Any],
    instance_model: Any,
    status_choices: Any,
) -> list[Any]:
    from django.contrib.contenttypes.models import ContentType

    instance_id = payload.get("workflow_instance_id")
    if instance_id:
        qs = instance_model.objects.filter(
            public_id=instance_id,
            status=status_choices.ACTIVE,
        )
        return list(qs)

    subject_type = payload.get("subject_type")
    object_id = payload.get("object_id")
    if subject_type and object_id is not None:
        names = {wf.name for wf in workflows_for_event(event_name)}
        if not names:
            return []
        try:
            app_label, model_name = str(subject_type).split(".", 1)
        except ValueError:
            return []
        try:
            ct = ContentType.objects.get(app_label=app_label, model=model_name)
        except ContentType.DoesNotExist:
            return []
        qs = instance_model.objects.filter(
            definition_name__in=names,
            content_type=ct,
            object_id=str(object_id),
            status=status_choices.ACTIVE,
        )
        return list(qs)

    return []


def attach_trigger(event_name: str) -> None:
    """Idempotently subscribe the workflow router to ``event_name``."""

    if event_name in _WIRED_EVENTS:
        return
    _WIRED_EVENTS.add(event_name)

    @subscribe(event_name)
    def _handler(payload: dict[str, Any], _event_name: str = event_name) -> None:
        try:
            _route(_event_name, payload)
        except EventError:
            raise
        except Exception as exc:
            _log.warning(
                "workflow.trigger_handler_failed",
                event_name=_event_name,
                error=str(exc),
            )


def register_workflow(definition: WorkflowDefinition) -> WorkflowDefinition:
    """Wrap ``registry.register_workflow`` and wire its triggers."""

    wf = _register_definition(definition)
    for transition in wf.transitions:
        if transition.trigger:
            attach_trigger(transition.trigger)
    return wf


def reset_triggers_for_tests() -> None:
    _WIRED_EVENTS.clear()
