"""
Workflow Services.

Business logic services for the workflow engine.
"""
import logging
from typing import Any, Dict, List, Optional
from datetime import timedelta

from django.db import transaction
from django.utils import timezone
from django.contrib.auth import get_user_model

from .models import (
    ProcessDefinition,
    ProcessInstance,
    Task,
    TaskAssignment,
    ProcessVariable,
    ProcessHistory,
    SLADefinition,
    SLAViolation,
    InstanceStatus,
    TaskStatus,
    TaskType,
    AssignmentType,
    HistoryEventType,
    DurationUnit,
    SLAAction,
)
from .signals import (
    process_started,
    process_completed,
    process_cancelled,
    process_failed,
    task_created,
    task_assigned,
    task_completed,
    sla_violated,
)


User = get_user_model()
logger = logging.getLogger(__name__)


class ProcessDefinitionService:
    """Service for managing process definitions."""
    
    @staticmethod
    def get_startable_definition(tenant, slug: str) -> Optional[ProcessDefinition]:
        """Get the latest startable (deployed and active) process definition by slug."""
        return ProcessDefinition.objects.filter(
            tenant=tenant,
            slug=slug,
            is_deployed=True,
            is_active=True
        ).order_by('-version').first()

    @staticmethod
    def get_all_versions(tenant, slug: str) -> List[ProcessDefinition]:
        """Get all versions of a process definition."""
        return list(ProcessDefinition.objects.filter(
            tenant=tenant,
            slug=slug
        ).order_by('-version'))

    @staticmethod
    @transaction.atomic
    def create_new_version(
        definition: ProcessDefinition,
        user,
        copy_sla: bool = True
    ) -> ProcessDefinition:
        """Create a new version of a process definition."""
        latest_version = ProcessDefinition.objects.filter(
            tenant=definition.tenant,
            slug=definition.slug
        ).order_by('-version').values_list('version', flat=True).first() or 0

        new_definition = ProcessDefinition.objects.create(
            tenant=definition.tenant,
            slug=definition.slug,
            name=definition.name,
            description=definition.description,
            bpmn_xml=definition.bpmn_xml,
            metadata=definition.metadata,
            version=latest_version + 1,
            is_deployed=False,
            is_active=True,
            created_by=user,
        )

        if copy_sla:
            for sla in definition.sla_definitions.filter(is_active=True):
                SLADefinition.objects.create(
                    definition=new_definition,
                    task_definition_key=sla.task_definition_key,
                    duration_value=sla.duration_value,
                    duration_unit=sla.duration_unit,
                    action_on_violation=sla.action_on_violation,
                    escalation_config=sla.escalation_config,
                    is_active=True,
                )

        return new_definition

    @staticmethod
    @transaction.atomic
    def deploy(definition: ProcessDefinition) -> ProcessDefinition:
        """Deploy a process definition."""
        if definition.is_deployed:
            raise ValueError("Process definition is already deployed")

        if not definition.bpmn_xml:
            raise ValueError("Cannot deploy without BPMN XML content")

        definition.is_deployed = True
        definition.deployed_at = timezone.now()
        definition.save(update_fields=['is_deployed', 'deployed_at'])

        logger.info(
            f"Deployed process definition: {definition.name} v{definition.version}"
        )

        return definition


class ProcessInstanceService:
    """Service for managing process instances."""
    
    @staticmethod
    @transaction.atomic
    def start(
        definition: ProcessDefinition,
        user,
        business_key: str = None,
        variables: Dict[str, Any] = None
    ) -> ProcessInstance:
        """Start a new process instance."""
        if not definition.is_deployed:
            raise ValueError("Cannot start instance from non-deployed definition")

        if not definition.is_active:
            raise ValueError("Cannot start instance from inactive definition")

        instance = ProcessInstance.objects.create(
            tenant=definition.tenant,
            definition=definition,
            business_key=business_key or '',
            started_by=user,
            context=variables or {},
            status=InstanceStatus.RUNNING,
        )

        # Create initial variables
        if variables:
            for name, value in variables.items():
                var_type = ProcessInstanceService._infer_type(value)
                ProcessVariable.objects.create(
                    instance=instance,
                    name=name,
                    value=value,
                    type=var_type,
                )

        logger.info(
            f"Started process instance: {instance.id} "
            f"(definition: {definition.name}, business_key: {business_key})"
        )

        # Emit signal
        process_started.send(
            sender=ProcessInstance,
            instance=instance,
            user=user,
            variables=variables or {}
        )

        return instance

    @staticmethod
    def _infer_type(value: Any) -> str:
        """Infer the type of a variable value."""
        if isinstance(value, bool):
            return 'boolean'
        elif isinstance(value, int):
            return 'number'
        elif isinstance(value, float):
            return 'number'
        elif isinstance(value, dict):
            return 'object'
        elif isinstance(value, list):
            return 'array'
        return 'string'

    @staticmethod
    @transaction.atomic
    def cancel(
        instance: ProcessInstance,
        user,
        reason: str = ''
    ) -> ProcessInstance:
        """Cancel a process instance."""
        if instance.status != InstanceStatus.RUNNING:
            raise ValueError("Only running instances can be cancelled")

        instance.status = InstanceStatus.CANCELLED
        instance.cancelled_at = timezone.now()
        instance.error_message = reason
        instance.save(update_fields=['status', 'cancelled_at', 'error_message'])

        # Cancel all active tasks
        instance.tasks.filter(
            status__in=[TaskStatus.CREATED, TaskStatus.ASSIGNED, TaskStatus.CLAIMED]
        ).update(status=TaskStatus.CANCELLED, completed_at=timezone.now())

        logger.info(f"Cancelled process instance: {instance.id}")

        # Emit signal
        process_cancelled.send(
            sender=ProcessInstance,
            instance=instance,
            user=user,
            reason=reason
        )

        return instance

    @staticmethod
    @transaction.atomic
    def complete(instance: ProcessInstance) -> ProcessInstance:
        """Complete a process instance."""
        instance.status = InstanceStatus.COMPLETED
        instance.completed_at = timezone.now()
        instance.save(update_fields=['status', 'completed_at'])

        logger.info(f"Completed process instance: {instance.id}")

        # Emit signal
        process_completed.send(sender=ProcessInstance, instance=instance)

        return instance

    @staticmethod
    @transaction.atomic
    def fail(instance: ProcessInstance, error_message: str) -> ProcessInstance:
        """Mark a process instance as failed."""
        instance.status = InstanceStatus.FAILED
        instance.error_message = error_message
        instance.save(update_fields=['status', 'error_message'])

        logger.error(f"Process instance failed: {instance.id} - {error_message}")

        # Emit signal
        process_failed.send(
            sender=ProcessInstance,
            instance=instance,
            error=error_message
        )

        return instance

    @staticmethod
    def suspend(instance: ProcessInstance) -> ProcessInstance:
        """Suspend a process instance."""
        if instance.status != InstanceStatus.RUNNING:
            raise ValueError("Only running instances can be suspended")

        instance.status = InstanceStatus.SUSPENDED
        instance.save(update_fields=['status'])

        logger.info(f"Suspended process instance: {instance.id}")

        return instance

    @staticmethod
    def resume(instance: ProcessInstance) -> ProcessInstance:
        """Resume a suspended process instance."""
        if instance.status != InstanceStatus.SUSPENDED:
            raise ValueError("Only suspended instances can be resumed")

        instance.status = InstanceStatus.RUNNING
        instance.save(update_fields=['status'])

        logger.info(f"Resumed process instance: {instance.id}")

        return instance

    @staticmethod
    def get_variable(instance: ProcessInstance, name: str, default=None):
        """Get a process variable value."""
        try:
            variable = instance.variables.get(name=name, scope__isnull=True)
            return variable.value
        except ProcessVariable.DoesNotExist:
            return default

    @staticmethod
    def set_variable(instance: ProcessInstance, name: str, value: Any, var_type: str = None):
        """Set a process variable value."""
        if var_type is None:
            var_type = ProcessInstanceService._infer_type(value)

        variable, created = ProcessVariable.objects.update_or_create(
            instance=instance,
            name=name,
            scope=None,
            defaults={'value': value, 'type': var_type}
        )

        return variable

    @staticmethod
    def get_all_variables(instance: ProcessInstance) -> Dict[str, Any]:
        """Get all global process variables as a dictionary."""
        return {
            var.name: var.value
            for var in instance.variables.filter(scope__isnull=True)
        }


class TaskService:
    """Service for managing tasks."""

    @staticmethod
    @transaction.atomic
    def create_task(
        instance: ProcessInstance,
        task_definition_key: str,
        element_id: str,
        name: str,
        task_type: str = TaskType.USER_TASK,
        description: str = '',
        form_data: Dict = None,
        input_variables: Dict = None,
        candidate_users: List = None,
        candidate_groups: List = None,
        due_date=None,
        priority: int = 50,
    ) -> Task:
        """Create a new task."""
        task = Task.objects.create(
            tenant=instance.tenant,
            instance=instance,
            task_definition_key=task_definition_key,
            element_id=element_id,
            name=name,
            task_type=task_type,
            description=description,
            form_data=form_data or {},
            input_variables=input_variables or {},
            priority=priority,
            due_date=due_date,
            status=TaskStatus.CREATED,
        )

        # Add candidate users
        if candidate_users:
            for user_id in candidate_users:
                TaskAssignment.objects.create(
                    task=task,
                    user_id=user_id,
                    assignment_type=AssignmentType.CANDIDATE,
                    is_active=True,
                )
            task.status = TaskStatus.ASSIGNED
            task.save(update_fields=['status'])

        # Add candidate groups
        if candidate_groups:
            for group_id in candidate_groups:
                TaskAssignment.objects.create(
                    task=task,
                    group_id=group_id,
                    assignment_type=AssignmentType.CANDIDATE,
                    is_active=True,
                )
            task.status = TaskStatus.ASSIGNED
            task.save(update_fields=['status'])

        # Check for SLA
        TaskService._apply_sla(task)

        logger.info(f"Created task: {task.name} (instance: {instance.id})")

        # Emit signal
        task_created.send(sender=Task, task=task, instance=instance)

        return task

    @staticmethod
    def _apply_sla(task: Task):
        """Apply SLA to a task if defined."""
        sla = SLADefinition.objects.filter(
            definition=task.instance.definition,
            task_definition_key=task.task_definition_key,
            is_active=True
        ).first()

        if sla and not task.due_date:
            # Calculate due date based on SLA
            multipliers = {
                DurationUnit.MINUTES: 1,
                DurationUnit.HOURS: 60,
                DurationUnit.DAYS: 60 * 24,
                DurationUnit.WEEKS: 60 * 24 * 7,
            }
            minutes = sla.duration_value * multipliers.get(sla.duration_unit, 60)
            task.due_date = timezone.now() + timedelta(minutes=minutes)
            task.save(update_fields=['due_date'])

    @staticmethod
    @transaction.atomic
    def claim(task: Task, user) -> Task:
        """Claim a task."""
        if not task.can_claim(user):
            raise ValueError("User cannot claim this task")

        task.assignee = user
        task.status = TaskStatus.CLAIMED
        task.claimed_at = timezone.now()
        task.save(update_fields=['assignee', 'status', 'claimed_at'])

        logger.info(f"Task claimed: {task.id} by {user}")

        return task

    @staticmethod
    @transaction.atomic
    def unclaim(task: Task) -> Task:
        """Release a claimed task."""
        if task.status != TaskStatus.CLAIMED:
            raise ValueError("Task is not claimed")

        previous_assignee = task.assignee
        task.assignee = None
        task.claimed_at = None
        task.status = TaskStatus.ASSIGNED if task.assignments.filter(is_active=True).exists() else TaskStatus.CREATED
        task.save(update_fields=['assignee', 'status', 'claimed_at'])

        logger.info(f"Task unclaimed: {task.id} by {previous_assignee}")

        return task

    @staticmethod
    @transaction.atomic
    def complete(task: Task, user, output_variables: Dict = None) -> Task:
        """Complete a task."""
        if not task.can_complete(user):
            raise ValueError("User cannot complete this task")

        task.status = TaskStatus.COMPLETED
        task.completed_at = timezone.now()
        task.output_variables = output_variables or {}
        task.save(update_fields=['status', 'completed_at', 'output_variables'])

        # Mark any SLA violations as resolved
        task.sla_violations.filter(is_resolved=False).update(
            is_resolved=True,
            resolved_at=timezone.now()
        )

        logger.info(f"Task completed: {task.id} by {user}")

        # Emit signal
        task_completed.send(
            sender=Task,
            task=task,
            user=user,
            variables=output_variables or {}
        )

        return task

    @staticmethod
    @transaction.atomic
    def delegate(task: Task, from_user, to_user) -> Task:
        """Delegate a task to another user."""
        if task.assignee != from_user:
            raise ValueError("Only the assignee can delegate this task")

        if task.status != TaskStatus.CLAIMED:
            raise ValueError("Only claimed tasks can be delegated")

        task.assignee = to_user
        task.save(update_fields=['assignee'])

        # Create delegation assignment record
        TaskAssignment.objects.create(
            task=task,
            user=to_user,
            assignment_type=AssignmentType.DELEGATED,
            is_active=True,
        )

        logger.info(f"Task delegated: {task.id} from {from_user} to {to_user}")

        return task

    @staticmethod
    @transaction.atomic
    def reassign(task: Task, user_ids: List = None, group_ids: List = None) -> Task:
        """Reassign task candidates (admin action)."""
        # Deactivate existing assignments
        task.assignments.filter(is_active=True).update(
            is_active=False,
            unassigned_at=timezone.now()
        )

        # Add new user assignments
        if user_ids:
            for user_id in user_ids:
                TaskAssignment.objects.create(
                    task=task,
                    user_id=user_id,
                    assignment_type=AssignmentType.CANDIDATE,
                    is_active=True,
                )

        # Add new group assignments
        if group_ids:
            for group_id in group_ids:
                TaskAssignment.objects.create(
                    task=task,
                    group_id=group_id,
                    assignment_type=AssignmentType.CANDIDATE,
                    is_active=True,
                )

        # Reset task status if claimed
        if task.status == TaskStatus.CLAIMED:
            task.assignee = None
            task.claimed_at = None
            task.status = TaskStatus.ASSIGNED if (user_ids or group_ids) else TaskStatus.CREATED
            task.save(update_fields=['assignee', 'status', 'claimed_at'])

        logger.info(f"Task reassigned: {task.id}")

        # Emit signal
        task_assigned.send(
            sender=Task,
            task=task,
            users=user_ids or [],
            groups=group_ids or []
        )

        return task

    @staticmethod
    def get_inbox(user) -> List[Task]:
        """Get user's task inbox (available tasks)."""
        return Task.objects.filter(
            tenant=user.tenant
        ).available_for(user).active().select_related(
            'instance',
            'instance__definition'
        ).order_by('-priority', '-created_at')

    @staticmethod
    def get_my_tasks(user) -> List[Task]:
        """Get tasks assigned to/claimed by user."""
        return Task.objects.filter(
            tenant=user.tenant,
            assignee=user
        ).active().select_related(
            'instance',
            'instance__definition'
        ).order_by('-priority', '-created_at')


class SLAService:
    """Service for SLA management."""

    @staticmethod
    def check_violations():
        """Check for SLA violations across all active tasks."""
        now = timezone.now()
        
        # Find tasks that are overdue but don't have an unresolved violation
        overdue_tasks = Task.objects.filter(
            status__in=[TaskStatus.CREATED, TaskStatus.ASSIGNED, TaskStatus.CLAIMED],
            due_date__lt=now
        ).exclude(
            sla_violations__is_resolved=False
        ).select_related('instance', 'instance__definition')

        violations = []
        for task in overdue_tasks:
            # Find SLA definition
            sla = SLADefinition.objects.filter(
                definition=task.instance.definition,
                task_definition_key=task.task_definition_key,
                is_active=True
            ).first()

            if sla:
                violation = SLAViolation.objects.create(
                    task=task,
                    sla_definition=sla,
                    expected_completion=task.due_date,
                )
                violations.append(violation)

                # Emit signal
                sla_violated.send(
                    sender=SLAViolation,
                    violation=violation,
                    task=task,
                    sla_definition=sla
                )

                # Handle action
                SLAService._handle_violation_action(violation, sla)

                logger.warning(
                    f"SLA violation: task {task.id}, "
                    f"expected {task.due_date}, now {now}"
                )

        return violations

    @staticmethod
    def _handle_violation_action(violation: SLAViolation, sla: SLADefinition):
        """Handle the action configured for SLA violation."""
        if sla.action_on_violation == SLAAction.NOTIFY:
            # TODO: Send notification to configured recipients
            pass
        elif sla.action_on_violation == SLAAction.ESCALATE:
            # TODO: Escalate to configured users/groups
            pass
        elif sla.action_on_violation == SLAAction.REASSIGN:
            # TODO: Reassign to configured users/groups
            pass

    @staticmethod
    def get_task_sla_status(task: Task) -> Dict:
        """Get SLA status for a task."""
        sla = SLADefinition.objects.filter(
            definition=task.instance.definition,
            task_definition_key=task.task_definition_key,
            is_active=True
        ).first()

        if not sla or not task.due_date:
            return {'has_sla': False}

        now = timezone.now()
        remaining = task.due_date - now
        is_overdue = remaining.total_seconds() < 0

        violations = task.sla_violations.filter(is_resolved=False)

        return {
            'has_sla': True,
            'due_date': task.due_date,
            'remaining_seconds': remaining.total_seconds(),
            'is_overdue': is_overdue,
            'has_violation': violations.exists(),
            'sla_definition': {
                'duration_value': sla.duration_value,
                'duration_unit': sla.duration_unit,
                'action_on_violation': sla.action_on_violation,
            }
        }
