from celery import shared_task
from celery.utils.log import get_task_logger
from celery.exceptions import Retry
from apps.business.distributor.models import DistributorID
from apps.business.distributor.services.ambassador_service import refresh_brand_ambassadors
from apps.business.rewards.services.level_service import calculate_level_progress, process_auto_id_generation
from apps.business.schemes.models import OpportunityBundleLevel

logger = get_task_logger(__name__)

@shared_task(bind=True, max_retries=3, default_retry_delay=10)
def process_auto_id_generation_task(self, distributor_id):
    logger.info(f"[START] process_auto_id_generation_task distributor_id={distributor_id}")
    try:
        # Idempotency: skip if already processed (e.g., auto_id_generated flag)
        distributor = DistributorID.objects.filter(id=distributor_id).first()
        if not distributor:
            logger.error(f"Distributor not found: distributor_id={distributor_id}")
            return
        if getattr(distributor, 'auto_id_generated', False):
            logger.info(f"[SKIP] Already processed: distributor_id={distributor_id}")
            return

        progress_result = calculate_level_progress(distributor)
        current_level = int(progress_result.get("current_level") or 0)
        if current_level <= 0:
            result = {"auto_ids_created": 0, "total_cost": 0}
        else:
            levels = OpportunityBundleLevel.objects.filter(
                opportunity_bundle=distributor.product.opportunity_bundle,
                level_number__lte=current_level,
                auto_id_enabled=True,
            ).order_by("level_number")

            total_auto_ids = 0
            total_cost = 0
            for level in levels:
                level_result = process_auto_id_generation(distributor, level)
                total_auto_ids += int(level_result.get("auto_ids_created") or 0)
                total_cost += level_result.get("total_cost") or 0

            result = {"auto_ids_created": total_auto_ids, "total_cost": total_cost}
        # Optionally set a flag to mark as processed (if model supports it)
        # distributor.auto_id_generated = True
        # distributor.save(update_fields=["auto_id_generated"])
        logger.info(f"[SUCCESS] process_auto_id_generation_task distributor_id={distributor_id} result={result}")
        return result
    except Exception as exc:
        logger.error(f"[FAIL] process_auto_id_generation_task distributor_id={distributor_id} error={exc}")
        try:
            self.retry(exc=exc)
        except Retry:
            logger.error(f"[RETRY-FAIL] process_auto_id_generation_task distributor_id={distributor_id}")


@shared_task(bind=True, max_retries=2, default_retry_delay=30)
def refresh_brand_ambassadors_task(self, threshold=100, opportunity_bundle_id=None):
    """Periodic-safe task for recalculating ambassador eligibility snapshots."""
    try:
        result = refresh_brand_ambassadors(min_distributor_ids=threshold, opportunity_bundle_id=opportunity_bundle_id)
        logger.info(f"[SUCCESS] refresh_brand_ambassadors_task result={result}")
        return result
    except Exception as exc:
        logger.error(f"[FAIL] refresh_brand_ambassadors_task error={exc}")
        try:
            self.retry(exc=exc)
        except Retry:
            logger.error("[RETRY-FAIL] refresh_brand_ambassadors_task")
