from collections import defaultdict from collections.abc import Sequence from dataclasses import dataclass from datetime import datetime, timedelta from anydi import singleton from fansifter_common.utils import timezone from email_campaigns.audiences.repositories import AudienceRepository from email_campaigns.campaigns.enums import EmailCampaignStatus from email_campaigns.campaigns.models import ( CampaignBatch, DeliveryQuota, EmailCampaign, ) from email_campaigns.campaigns.repositories.campaign import EmailCampaignRepository from email_campaigns.campaigns.repositories.campaign_batch import ( CampaignBatchRepository, ) from email_campaigns.campaigns.repositories.delivery_quota import ( DeliveryQuotaRepository, ) from email_campaigns.campaigns.services.campaign import EmailCampaignService MAX_SIMULATION_MINUTES = 24 * 60 * 365 # Safety limit: 1 year MIN_HOURLY_QUOTA = 20.0 # Default minimum hourly quota @dataclass class CampaignState: """State of a campaign during simulation.""" campaign_id: str scheduled_at: datetime | None send_at: datetime | None created_at: datetime | None remaining_recipients: float def is_eligible(self, current_time: datetime) -> bool: """Check if campaign is eligible to send at current_time.""" if self.send_at is None: return True return self.send_at <= current_time def fcfs_sort_key(self) -> tuple[bool, datetime | None, datetime | None, str]: """Sort key for FCFS ordering.""" return ( self.scheduled_at is not None, self.scheduled_at or self.send_at, self.created_at, self.campaign_id, ) @singleton class SendTimeEstimatorService: def __init__( self, audience_repository: AudienceRepository, campaign_repository: EmailCampaignRepository, batch_repository: CampaignBatchRepository, quota_repository: DeliveryQuotaRepository, campaign_service: EmailCampaignService, ) -> None: self.audience_repository = audience_repository self.campaign_repository = campaign_repository self.batch_repository = batch_repository self.quota_repository = quota_repository self.campaign_service = campaign_service def estimate_send_duration( self, campaign: EmailCampaign, send_at: datetime, ) -> tuple[int, int]: if not campaign.email_domain_id: return (0, 0) audience = self.audience_repository.get(campaign.audience_id) if not audience: return (0, 0) # Determine recipients per provider if audience.recipients_count_by_email_provider: recipients_by_provider = audience.recipients_count_by_email_provider elif audience.fan_count > 0: recipients_by_provider = None else: return (0, 0) pending_campaigns = self.campaign_repository.find_pending_campaigns( domain_id=campaign.email_domain_id ) self.campaign_service.prefetch_audience(pending_campaigns) pending_batches = self.batch_repository.find_active_batches( campaign.email_domain_id ) # Load provider quotas domain_quotas = self.quota_repository.find_by_domain(campaign.email_domain_id) delivery_quotas = {quota.provider: quota for quota in domain_quotas} # If no per-provider counts exist, split equally among available providers if recipients_by_provider is None: providers = set(delivery_quotas.keys()) if not providers: return (0, 0) recipients_by_provider = dict.fromkeys(providers, audience.fan_count) if not recipients_by_provider: return (0, 0) # Sort pending campaigns by FCFS rules sorted_campaigns = self._sort_campaigns_by_fcfs( list(pending_campaigns), campaign ) # Build batch index once for all campaigns in domain batch_index = self._build_batch_index(pending_batches) max_duration_seconds = 0.0 # Calculate duration separately for each provider for provider, recipient_count in recipients_by_provider.items(): duration = self._estimate_for_provider( provider=provider, recipient_count=recipient_count, campaign=campaign, sorted_campaigns=sorted_campaigns, batch_index=batch_index, delivery_quotas=delivery_quotas, send_at=send_at, ) max_duration_seconds = max(max_duration_seconds, duration) emails_to_send = self._calculate_emails_to_send( campaign=campaign, pending_batches=pending_batches, audience_fan_count=audience.fan_count, ) return (int(max_duration_seconds), emails_to_send) @staticmethod def _calculate_emails_to_send( campaign: EmailCampaign, pending_batches: Sequence[CampaignBatch], audience_fan_count: int, ) -> int: """Calculate remaining emails to send for campaign.""" if campaign.status != EmailCampaignStatus.IN_PROGRESS: return audience_fan_count # Calculate progress from already loaded pending_batches campaign_batches = [b for b in pending_batches if b.campaign_id == campaign.id] if campaign_batches: total_batch_size = sum(b.batch_size for b in campaign_batches) total_batch_offset = sum(b.batch_offset for b in campaign_batches) return total_batch_size - total_batch_offset return audience_fan_count @staticmethod def _build_batch_index( batches: Sequence[CampaignBatch], ) -> dict[str, dict[str, CampaignBatch]]: """ Build O(1) lookup index for batches. Returns: {campaign_id: {provider: batch}} Complexity: O(B) to build, O(1) per lookup """ index: dict[str, dict[str, CampaignBatch]] = defaultdict(dict) for batch in batches: index[batch.campaign_id][batch.provider] = batch return index @staticmethod def _sort_campaigns_by_fcfs( pending_campaigns: list[EmailCampaign], target_campaign: EmailCampaign, ) -> list[EmailCampaign]: campaigns = list(pending_campaigns) if target_campaign not in campaigns: campaigns.append(target_campaign) return sorted( campaigns, key=lambda c: ( c.scheduled_at is not None, c.scheduled_at or c.send_at, c.created_at, c.id, ), ) def _estimate_for_provider( self, provider: str, recipient_count: int, campaign: EmailCampaign, sorted_campaigns: list[EmailCampaign], batch_index: dict[str, dict[str, CampaignBatch]], delivery_quotas: dict[str, DeliveryQuota], send_at: datetime, ) -> float: """Estimate send duration for a single provider.""" if recipient_count <= 0: return 0.0 quota = delivery_quotas.get(provider) hourly_quota = ( quota.last_hourly_quota if quota and quota.last_hourly_quota else MIN_HOURLY_QUOTA ) if hourly_quota <= 0: return 0.0 # Build state list for this provider campaign_states = self._build_campaign_states_for_provider( sorted_campaigns=sorted_campaigns, target_campaign_id=campaign.id, target_recipients=recipient_count, provider=provider, batch_index=batch_index, ) if not campaign_states: return 0.0 target_state = next( (s for s in campaign_states if s.campaign_id == campaign.id), None ) if not target_state: return 0.0 return self._simulate_fcfs( campaign_states=campaign_states, target_campaign_id=campaign.id, hourly_quota=hourly_quota, start_time=send_at, ) @staticmethod def _build_campaign_states_for_provider( sorted_campaigns: list[EmailCampaign], target_campaign_id: str, target_recipients: int, provider: str, batch_index: dict[str, dict[str, CampaignBatch]], ) -> list[CampaignState]: states: list[CampaignState] = [] for campaign in sorted_campaigns: if campaign.id == target_campaign_id: recipient_count = target_recipients else: matching_batch = batch_index.get(campaign.id, {}).get(provider) if matching_batch: remaining = matching_batch.batch_size - matching_batch.batch_offset recipient_count = max(remaining, 0) elif campaign.audience: if campaign.audience.recipients_count_by_email_provider: recipient_count = ( campaign.audience.recipients_count_by_email_provider.get( provider, 0 ) ) else: recipient_count = campaign.audience.fan_count else: recipient_count = 0 if recipient_count <= 0: continue states.append( CampaignState( campaign_id=campaign.id, scheduled_at=campaign.scheduled_at, send_at=campaign.send_at, created_at=campaign.created_at, remaining_recipients=float(recipient_count), ) ) return states @staticmethod def _simulate_fcfs( # noqa: C901 campaign_states: list[CampaignState], target_campaign_id: str, hourly_quota: float, start_time: datetime, ) -> float: """ Simulate FCFS dispatcher until target campaign completes. Returns total seconds from start_time until target campaign finishes. """ if hourly_quota <= 0: return 0.0 minute_quota = hourly_quota / 60.0 if minute_quota <= 0: return 0.0 current_time = start_time target_state = next( (s for s in campaign_states if s.campaign_id == target_campaign_id), None, ) if not target_state or target_state.remaining_recipients <= 0: return 0.0 # campaign_states are already in FCFS order from sorted_campaigns ordered_states = campaign_states max_duration = timedelta(minutes=MAX_SIMULATION_MINUTES) for _ in range(MAX_SIMULATION_MINUTES): if target_state.remaining_recipients <= 0: return (current_time - start_time).total_seconds() quota_remaining = minute_quota any_eligible_with_work = False # Distribute quota in FCFS order for state in ordered_states: if quota_remaining <= 0: break if not state.is_eligible(current_time): continue if state.remaining_recipients <= 0: continue any_eligible_with_work = True # Amount of emails sent this minute for this campaign emails_to_send = min(quota_remaining, state.remaining_recipients) state.remaining_recipients -= emails_to_send quota_remaining -= emails_to_send # Target finished inside this minute if ( state.campaign_id == target_campaign_id and state.remaining_recipients <= 0 ): emails_sent_this_minute = minute_quota - quota_remaining fraction_of_minute = ( emails_sent_this_minute / minute_quota if minute_quota > 0 else 0.0 ) exact_time = current_time + timedelta(minutes=fraction_of_minute) return (exact_time - start_time).total_seconds() # No campaign received any quota — jump to the next earliest send_at if not any_eligible_with_work: future_times = [ s.send_at for s in ordered_states if ( s.remaining_recipients > 0 and s.send_at is not None and s.send_at > current_time ) ] if future_times: next_time = min(future_times) if next_time - start_time > max_duration: return max_duration.total_seconds() current_time = next_time continue # Advance to the next minute current_time += timedelta(minutes=1) if current_time - start_time > max_duration: return max_duration.total_seconds() # Reached simulation limit return max_duration.total_seconds() def estimate_send_duration_bulk( # noqa: C901 self, campaigns: list[EmailCampaign] ) -> list[tuple[str, int, int]]: if not campaigns: return [] campaigns_by_domain: dict[str, list[EmailCampaign]] = defaultdict(list) for campaign in campaigns: if campaign.email_domain_id: campaigns_by_domain[campaign.email_domain_id].append(campaign) all_domain_ids = list(campaigns_by_domain.keys()) audience_ids = [c.audience_id for c in campaigns if c.audience_id] audiences_by_id = {} if audience_ids: audience_list = self.audience_repository.find_by_ids(audience_ids) audiences_by_id = {a.id: a for a in audience_list} all_pending_campaigns: Sequence[EmailCampaign] = [] if all_domain_ids: all_pending_campaigns = ( self.campaign_repository.find_pending_campaigns_for_domains( all_domain_ids ) ) if all_pending_campaigns: self.campaign_service.prefetch_audience(all_pending_campaigns) pending_by_domain: dict[str, list[EmailCampaign]] = defaultdict(list) for pending in all_pending_campaigns: if pending.email_domain_id: pending_by_domain[pending.email_domain_id].append(pending) all_batches: Sequence[CampaignBatch] = [] if all_domain_ids: all_batches = self.batch_repository.find_active_batches_for_domains( all_domain_ids ) batches_by_campaign: dict[str, list[CampaignBatch]] = defaultdict(list) for batch in all_batches: batches_by_campaign[batch.campaign_id].append(batch) all_quotas: Sequence[DeliveryQuota] = [] if all_domain_ids: all_quotas = self.quota_repository.find_by_domains(all_domain_ids) quotas_by_domain: dict[str, list[DeliveryQuota]] = defaultdict(list) for quota in all_quotas: quotas_by_domain[quota.domain_id].append(quota) results: list[tuple[str, int, int]] = [] is_single_domain = len(campaigns_by_domain) == 1 for domain_id, domain_campaigns in campaigns_by_domain.items(): # Optimize for single domain case (most clients requests will be single domain) if is_single_domain: domain_pending = list(all_pending_campaigns) domain_quotas = {q.provider: q for q in all_quotas} domain_batches = all_batches else: domain_pending = pending_by_domain.get(domain_id, []) domain_quotas = { q.provider: q for q in quotas_by_domain.get(domain_id, []) } domain_batches = [b for b in all_batches if b.domain_id == domain_id] batch_index = self._build_batch_index(domain_batches) sorted_pending = sorted( domain_pending, key=lambda c: ( c.scheduled_at is not None, c.scheduled_at or c.send_at, c.created_at, c.id, ), ) for campaign in domain_campaigns: try: if campaign.status == EmailCampaignStatus.IN_PROGRESS: campaign_send_at = timezone.now() else: campaign_send_at = campaign.send_at or timezone.now() if not campaign.audience_id: results.append((campaign.id, 0, 0)) continue audience = audiences_by_id.get(campaign.audience_id) if not audience: results.append((campaign.id, 0, 0)) continue if audience.recipients_count_by_email_provider: recipients_by_provider = ( audience.recipients_count_by_email_provider ) elif audience.fan_count > 0: recipients_by_provider = None else: results.append((campaign.id, 0, 0)) continue if recipients_by_provider is None: providers = set(domain_quotas.keys()) if not providers: results.append((campaign.id, 0, 0)) continue recipients_by_provider = dict.fromkeys( providers, audience.fan_count ) if not recipients_by_provider: results.append((campaign.id, 0, 0)) continue max_duration = 0.0 for provider, recipient_count in recipients_by_provider.items(): duration = self._estimate_for_provider( provider=provider, recipient_count=recipient_count, campaign=campaign, sorted_campaigns=sorted_pending, batch_index=batch_index, delivery_quotas=domain_quotas, send_at=campaign_send_at, ) max_duration = max(max_duration, duration) emails_to_send = self._calculate_emails_to_send( campaign=campaign, pending_batches=domain_batches, audience_fan_count=audience.fan_count, ) results.append((campaign.id, int(max_duration), emails_to_send)) except Exception: results.append((campaign.id, 0, 0)) return results