import concurrent.futures import logging from collections.abc import Sequence from datetime import datetime from anydi import singleton from fansifter_common.adapters.twilio.client import TwilioClient from fansifter_common.adapters.twilio.exceptions import TwilioClientError from fansifter_common.utils import timezone from app.adapters.db import Database from app.adapters.ows_text_campaigns import ( OwsTextCampaignsClient, PersonalizedAttributes, ) from app.config import Settings from app.enums import CampaignCancelReason, MessageSendQueueStatus from app.exceptions import MessageSendValidationError from app.models import Campaign, MessageSendQueue from app.repositories import ( AudienceFanRepository, CampaignRepository, MessageSendQueueRepository, ) from app.types import ( MessageSend, MessageSendResult, PreparedMessage, ProcessableSendQueue, ) logger = logging.getLogger(__name__) @singleton class TwilioSendMessagesHandler: max_messages_to_send = 100 def __init__( self, db: Database, campaign_repository: CampaignRepository, send_queue_repository: MessageSendQueueRepository, audience_fan_repository: AudienceFanRepository, ows_text_campaigns_client: OwsTextCampaignsClient, twilio_client: TwilioClient, settings: Settings, ) -> None: self.db = db self.campaign_repository = campaign_repository self.send_queue_repository = send_queue_repository self.audience_fan_repository = audience_fan_repository self.ows_text_campaigns_client = ows_text_campaigns_client self.twilio_client = twilio_client self.settings = settings def handle(self) -> list[MessageSendResult]: if self.settings.environment == "prod": logger.info("Skip for now.") return [] # Get the schedules with self.db.session_factory(): campaigns: list[Campaign] = [] if not campaigns: logger.info("No schedules found.") return [] # Create queues for the schedules with self.db.session_factory(), self.db.transaction(): self.create_queues_for_campaigns(campaigns) # Process the queues return self.process_queues() def create_queues_for_campaigns(self, campaigns: list[Campaign]) -> None: """Create queues for the given campaigns.""" # Find the ready queues for the campaigns in the schedules queues = self.send_queue_repository.find_ready_by_campaign_ids( campaign_ids={campaign.id for campaign in campaigns} ) # Create a map of queues by campaign id queue_by_campaign_id = {queue.campaign_id: queue for queue in queues} # Create queues for the schedules for campaign in campaigns: if queue := queue_by_campaign_id.get(campaign.id): # Update the send_at if different if campaign.send_at is not None and campaign.send_at != queue.send_at: queue.send_at = campaign.send_at self.send_queue_repository.save(queue) continue # Create queue for the schedule self.create_queue_for_campaign(campaign) def create_queue_for_campaign(self, campaign: Campaign) -> None: """Create a queue for the given campaign.""" logger.info( "Creating queue for the campaign `%s` ...", campaign.id, extra={ "text_campaign": { "id": campaign.id, }, }, ) # Check if the campaign is processable if not campaign.is_processable or campaign.send_at is None: logger.error( "Campaign is not processable.", extra={ "email_campaign": { "id": campaign.id, "status": campaign.status, "send_at": campaign.send_at, } }, ) return None # Get the audience snapshot if campaign.audience_snapshot_id is None: logger.error( "Campaign audience snapshot is not set.", extra={ "text_campaign": { "id": campaign.id, } }, ) return None # Get the total recipients total = self.audience_fan_repository.count_by_snapshot_id( campaign.audience_snapshot_id ) # Update the recipients count if different if campaign.recipients_count != total and total > 0: campaign.recipients_count = total self.campaign_repository.save(campaign) # Check if the total is zero if total == 0: campaign.recipients_count = 0 campaign.cancel(reason=CampaignCancelReason.NO_FANS) self.campaign_repository.save(campaign) logger.error( "Campaign has no fans.", extra={ "text_campaign": { "id": campaign.id, } }, ) return None queue = MessageSendQueue( campaign_id=campaign.id, total=total, send_at=campaign.send_at, ) self.send_queue_repository.save(queue) return None def process_queues( self, ready_queues: Sequence[MessageSendQueue] | None = None ) -> list[MessageSendResult]: """Process the given queues.""" # Get the processable queues with self.db.session_factory(): processable_queues = self.get_processable_queues(ready_queues) # Check if there are no processable queues if not processable_queues: logger.info("No processable queues found.") return [] # Process the queues if self.settings.process_in_parallel: return self.process_in_parallel(processable_queues) return self.process(processable_queues) def get_processable_queues( self, ready_queues: Sequence[MessageSendQueue] | None ) -> list[ProcessableSendQueue]: """Get the processable queues for the given ready queues.""" # Get the ready queues ready_queues = ready_queues or self.send_queue_repository.find_ready() processable_queues: list[ProcessableSendQueue] = [] for queue in ready_queues: messages_to_send = min( queue.total - queue.offset, self.max_messages_to_send, ) processable_queues.append( ProcessableSendQueue( queue=queue, messages_to_send=messages_to_send, ) ) return processable_queues def process( self, processable_queues: list[ProcessableSendQueue] ) -> list[MessageSendResult]: """Process the given queues.""" send_results: list[MessageSendResult] = [] for processable_queue in processable_queues: try: send_results.append(self.process_queue(processable_queue)) except Exception as exc: logger.error( str(exc), exc_info=exc, extra={ "message_send_queue": { "id": processable_queue.queue.id, }, "text_campaign": { "id": processable_queue.queue.campaign_id, }, }, ) return send_results def process_in_parallel( self, processable_queues: list[ProcessableSendQueue] ) -> list[MessageSendResult]: """Process the given queues in parallel.""" send_results: list[MessageSendResult] = [] with concurrent.futures.ThreadPoolExecutor( max_workers=len(processable_queues) ) as executor: future_to_queue = { executor.submit( self.process_queue, processable_queue ): processable_queue.queue for processable_queue in processable_queues } for future in concurrent.futures.as_completed(future_to_queue): queue = future_to_queue[future] try: send_results.append(future.result()) except Exception as exc: logger.error( str(exc), exc_info=exc, extra={ "message_send_queue": { "id": queue.id, }, "text_campaign": { "id": queue.campaign_id, }, }, ) return send_results def process_queue( self, processable_queue: ProcessableSendQueue ) -> MessageSendResult: """Process the given queue.""" with self.db.session_factory(): # Create message send for the queue message_send = self.create_message_send(processable_queue) # Send messages for the queue return self.send_messages(message_send=message_send) def create_message_send( self, processable_queue: ProcessableSendQueue ) -> MessageSend: """Create send for the given queue.""" send_at = timezone.now() # Get the queue queue, messages_to_send = processable_queue # Get the campaign campaign = self.campaign_repository.get(queue.campaign_id) # Check if the campaign is processable if campaign is None or campaign.is_deleted or campaign.is_sent: self.cancel_queue(queue) raise MessageSendValidationError("Campaign is already sent or not found.") # Check if the content is missing if not campaign.content: raise MessageSendValidationError("Campaign subject or content is missing.") # Check if the audience snapshot is set if campaign.audience_snapshot_id is None: raise MessageSendValidationError("Campaign audience snapshot is not set.") # Get the fans data messages = self.prepare_messages( campaign_id=campaign.id, snapshot_id=campaign.audience_snapshot_id, limit=messages_to_send, offset=queue.offset, ) return MessageSend( queue=queue, campaign=campaign, messages=messages, messages_to_send=messages_to_send, send_at=send_at, ) def send_messages(self, message_send: MessageSend) -> MessageSendResult: """Send messages for the given queue.""" logger.info( "Sending `%s` campaign messages ...", message_send.campaign.name, extra={ "message_send_queue": { "id": message_send.queue.id, "messages_to_send": message_send.messages_to_send, }, "text_campaign": { "id": message_send.queue.campaign_id, }, }, ) with self.db.transaction(): # Get the queue for update queue = self.send_queue_repository.get_for_update(message_send.queue.id) # Set the campaign in progress self.campaign_repository.save_as_in_progress(message_send.campaign) # Sent messages counter sent_messages = 0 queue_offset = message_send.messages_to_send # Send Twilio message in batches for message in message_send.messages: try: logger.info( "Sending Twilio message `%s` to `%s` using `%s` channel ...", message.sender, message.recipient, message.channel, extra={ "message_send_queue": { "id": queue.id, }, "text_campaign": { "id": message_send.campaign.id, }, }, ) sent_messages += 1 except TwilioClientError as exc: logger.error( exc.log_msg, exc_info=exc, extra={ "message_send_queue": { "id": queue.id, }, "text_campaign": { "id": message_send.campaign.id, }, **exc.log_extra, }, ) break # Update the campaign queue if queue_offset > 0: queue = self.update_queue( queue, offset=queue_offset, send_at=message_send.send_at, ) logger.info( "Successfully sent messages for the `%s` campaign queue.", message_send.campaign.name, extra={ "message_send_result": { "campaign_id": message_send.campaign.id, "campaign_name": message_send.campaign.name, "queue_id": queue.id, "queue_offset": queue_offset, "sent_messages": sent_messages, "messages_to_send": message_send.messages_to_send, } }, ) # Check if the queue is complete if queue.is_completed: # Set the campaign as sent self.campaign_repository.save_as_sent(message_send.campaign) logger.info( "Sending messages is completed for the `%s` campaign queue.", message_send.campaign.name, extra={ "message_send_queue": { "id": queue.id, }, "text_campaign": { "id": message_send.campaign.id, }, }, ) return MessageSendResult( campaign_id=queue.campaign_id, queue_id=queue.id, sent_messages=sent_messages, ) def prepare_messages( self, campaign_id: str, snapshot_id: str, limit: int, offset: int, ) -> list[PreparedMessage]: """Get fans data for the given snapshot id.""" fans = self.audience_fan_repository.find_by_snapshot_id( snapshot_id=snapshot_id, limit=limit, offset=offset, ) attributes = [ PersonalizedAttributes( fan_id=fan.fan_id, phone_number=fan.phone_number, channel=fan.channel, ) for fan in fans ] result = self.ows_text_campaigns_client.render_messages( campaign_id=campaign_id, attributes=attributes, ) return [ PreparedMessage( sender=item.sender, recipient=item.recipient, channel=item.channel, message=item.message, media_url=item.media_url, ) for item in result ] def cancel_queue(self, queue: MessageSendQueue) -> None: """Cancel the given queue.""" queue.status = MessageSendQueueStatus.CANCELED self.send_queue_repository.save(queue) logger.warning( "Campaign queue is canceled.", extra={ "message_send_queue": { "id": queue.id, }, "text_campaign": { "id": queue.campaign_id, }, }, ) def update_queue( self, queue: MessageSendQueue, *, offset: int, send_at: datetime, ) -> MessageSendQueue: """Update the given queue.""" last_sent_at = send_at # Move the offset queue.offset += offset # Update first sent at if queue.first_sent_at is None: queue.first_sent_at = last_sent_at # Update last sent at queue.last_sent_at = last_sent_at # Set the queue status if queue.offset >= queue.total: queue.status = MessageSendQueueStatus.COMPLETED else: queue.status = MessageSendQueueStatus.PROCESSING # Save the updated queue self.send_queue_repository.save(queue) logger.info( "Message campaign queue `%s` is updated.", queue.id, extra={ "message_send_queue": { "id": queue.id, "offset": offset, }, "text_campaign": { "id": queue.campaign_id, }, }, ) return queue