From 9a7570d9dacf6fd110976e2cf44129780236e377 Mon Sep 17 00:00:00 2001 From: Anton Krytskyi Date: Fri, 17 Jul 2026 15:48:32 +0300 Subject: [PATCH 1/2] sort recipients by activity priority --- osf/email/notification_campaign.py | 55 +++++++++++++++++++++--------- 1 file changed, 38 insertions(+), 17 deletions(-) diff --git a/osf/email/notification_campaign.py b/osf/email/notification_campaign.py index 31a041a5642..c97daf4fd1a 100644 --- a/osf/email/notification_campaign.py +++ b/osf/email/notification_campaign.py @@ -1,6 +1,9 @@ import logging +from enum import IntEnum + from osf.models import NotificationType, NotificationTypeEnum, OSFUser, UserActivityCounter, Email -from django.db.models import OuterRef, Subquery, Exists, F, Q, Case, When, CharField +from osf.models.spam import SpamStatus +from django.db.models import OuterRef, Subquery, Exists, F, Q, Case, When, CharField, IntegerField, Value from django.db.models.functions import Coalesce from framework.celery_tasks import app as celery_app from celery import chord @@ -10,6 +13,17 @@ logger = logging.getLogger(__name__) +ACTIVITY_ACTIVE_THRESHOLD = 5000 +ACTIVITY_SPAM_THRESHOLD = 100000 + + +class ActivityPriority(IntEnum): + """Numeric priorities for ordering recipients (higher = send first).""" + SPAM_NON_ACTIVE = 0 + NON_ACTIVE = 1 + SPAM_ACTIVE = 2 + ACTIVE = 3 + first_email_subquery = ( Email.objects @@ -49,8 +63,23 @@ def filter_users(filters, campaign_id=None, restart_failed=False): def get_filtered_batches(filters, batch_size=1000, campaign_id=None, restart_failed=False): qs = filter_users(filters, campaign_id, restart_failed=restart_failed) - qs = qs.annotate(activity_total=Coalesce(Subquery(counter_subquery), 0)).order_by('-activity_total', '-date_registered', '-id') + is_spam = Q(spam_status=SpamStatus.SPAM) + qs = ( + qs + .annotate( + activity_total=Coalesce(Subquery(counter_subquery), 0), + ).annotate( + activity_priority=Case( + When(~is_spam & Q(activity_total__gte=ACTIVITY_ACTIVE_THRESHOLD), then=Value(ActivityPriority.ACTIVE)), + When(is_spam & Q(activity_total__gte=ACTIVITY_SPAM_THRESHOLD), then=Value(ActivityPriority.SPAM_ACTIVE)), + When(~is_spam, then=Value(ActivityPriority.NON_ACTIVE)), + default=Value(ActivityPriority.SPAM_NON_ACTIVE), + output_field=IntegerField(), + ), + ).order_by('-activity_priority', '-activity_total', '-date_registered', '-id') + ) + last_priority = None last_total = None last_date = None last_id = None @@ -58,30 +87,22 @@ def get_filtered_batches(filters, batch_size=1000, campaign_id=None, restart_fai while True: batch_qs = qs - if last_total is not None: + if last_priority is not None: batch_qs = batch_qs.filter( - Q(activity_total__lt=last_total) | - Q(activity_total=last_total, date_registered__lt=last_date) | - Q(activity_total=last_total, date_registered=last_date, id__lt=last_id) + Q(activity_priority__lt=last_priority) | + Q(activity_priority=last_priority, activity_total__lt=last_total) | + Q(activity_priority=last_priority, activity_total=last_total, date_registered__lt=last_date) | + Q(activity_priority=last_priority, activity_total=last_total, date_registered=last_date, id__lt=last_id) ) batch = batch_qs[:batch_size] - - if not batch: - break - - rows = list(batch.values_list('id', 'activity_total', 'date_registered')) + rows = list(batch.values_list('id', 'activity_priority', 'activity_total', 'date_registered')) if not rows: break batch_ids = [r[0] for r in rows] - - last_id, last_total, last_date = ( - rows[-1][0], - rows[-1][1], - rows[-1][2], - ) + last_id, last_priority, last_total, last_date = rows[-1] yield batch_ids From 20d1916d64d8af038559b4bcfa3d0f5959667d5e Mon Sep 17 00:00:00 2001 From: Anton Krytskyi Date: Fri, 17 Jul 2026 20:42:30 +0300 Subject: [PATCH 2/2] simplify filter; send campaign email tasks in 3 phases --- osf/email/notification_campaign.py | 156 +++++++++++++++++++---------- 1 file changed, 101 insertions(+), 55 deletions(-) diff --git a/osf/email/notification_campaign.py b/osf/email/notification_campaign.py index c97daf4fd1a..4b68904c56c 100644 --- a/osf/email/notification_campaign.py +++ b/osf/email/notification_campaign.py @@ -1,28 +1,19 @@ import logging -from enum import IntEnum from osf.models import NotificationType, NotificationTypeEnum, OSFUser, UserActivityCounter, Email from osf.models.spam import SpamStatus -from django.db.models import OuterRef, Subquery, Exists, F, Q, Case, When, CharField, IntegerField, Value +from django.db.models import OuterRef, Subquery, Exists, F, Q, Case, When, CharField from django.db.models.functions import Coalesce from framework.celery_tasks import app as celery_app -from celery import chord +from celery import chord, group, chain from django.utils import timezone from osf.models.notification_campaign import NotificationCampaign, NotificationCampaignRecipient, NotificationCampaignStatus, NotificationCampaignRecipientStatus from osf.email import send_email_with_send_grid logger = logging.getLogger(__name__) -ACTIVITY_ACTIVE_THRESHOLD = 5000 -ACTIVITY_SPAM_THRESHOLD = 100000 - - -class ActivityPriority(IntEnum): - """Numeric priorities for ordering recipients (higher = send first).""" - SPAM_NON_ACTIVE = 0 - NON_ACTIVE = 1 - SPAM_ACTIVE = 2 - ACTIVE = 3 +ACTIVITY_THRESHOLD = 5000 +DEFAULT_BATCH_SIZE = 1000 first_email_subquery = ( @@ -60,26 +51,28 @@ def filter_users(filters, campaign_id=None, restart_failed=False): return qs -def get_filtered_batches(filters, batch_size=1000, campaign_id=None, restart_failed=False): +def get_filtered_batches( + filters, + batch_size=1000, + campaign_id=None, + restart_failed=False, + min_activity=None, + max_activity=None, + exclude_spam=False, +): qs = filter_users(filters, campaign_id, restart_failed=restart_failed) + if exclude_spam: + qs = qs.exclude(spam_status=SpamStatus.SPAM) - is_spam = Q(spam_status=SpamStatus.SPAM) - qs = ( - qs - .annotate( - activity_total=Coalesce(Subquery(counter_subquery), 0), - ).annotate( - activity_priority=Case( - When(~is_spam & Q(activity_total__gte=ACTIVITY_ACTIVE_THRESHOLD), then=Value(ActivityPriority.ACTIVE)), - When(is_spam & Q(activity_total__gte=ACTIVITY_SPAM_THRESHOLD), then=Value(ActivityPriority.SPAM_ACTIVE)), - When(~is_spam, then=Value(ActivityPriority.NON_ACTIVE)), - default=Value(ActivityPriority.SPAM_NON_ACTIVE), - output_field=IntegerField(), - ), - ).order_by('-activity_priority', '-activity_total', '-date_registered', '-id') - ) + qs = qs.annotate(activity_total=Coalesce(Subquery(counter_subquery), 0)) + + if min_activity is not None: + qs = qs.filter(activity_total__gte=min_activity) + if max_activity is not None: + qs = qs.filter(activity_total__lt=max_activity) + + qs = qs.order_by('-activity_total', '-date_registered', '-id') - last_priority = None last_total = None last_date = None last_id = None @@ -87,26 +80,58 @@ def get_filtered_batches(filters, batch_size=1000, campaign_id=None, restart_fai while True: batch_qs = qs - if last_priority is not None: + if last_total is not None: batch_qs = batch_qs.filter( - Q(activity_priority__lt=last_priority) | - Q(activity_priority=last_priority, activity_total__lt=last_total) | - Q(activity_priority=last_priority, activity_total=last_total, date_registered__lt=last_date) | - Q(activity_priority=last_priority, activity_total=last_total, date_registered=last_date, id__lt=last_id) + Q(activity_total__lt=last_total) | + Q(activity_total=last_total, date_registered__lt=last_date) | + Q(activity_total=last_total, date_registered=last_date, id__lt=last_id) ) batch = batch_qs[:batch_size] - rows = list(batch.values_list('id', 'activity_priority', 'activity_total', 'date_registered')) + rows = list(batch.values_list('id', 'activity_total', 'date_registered')) if not rows: break batch_ids = [r[0] for r in rows] - last_id, last_priority, last_total, last_date = rows[-1] + last_id, last_total, last_date = rows[-1] yield batch_ids +def build_campaign_group( + filters, + batch_size=DEFAULT_BATCH_SIZE, + campaign_id=None, + restart_failed=False, + min_activity=None, + max_activity=None, + exclude_spam=True, + **send_kwargs, +): + tasks = [] + total_recipients = 0 + for batch in get_filtered_batches( + filters, + batch_size=batch_size, + campaign_id=campaign_id, + restart_failed=restart_failed, + min_activity=min_activity, + max_activity=max_activity, + exclude_spam=exclude_spam, + ): + tasks.append( + send_campaign_batch.si( + recipients_ids=batch, + campaign_id=campaign_id, + **send_kwargs, + ) + ) + total_recipients += len(batch) + + return group(tasks), total_recipients + + FILTER_PRESETS = { 'all': {}, 'active': {'is_active': True}, @@ -115,12 +140,11 @@ def get_filtered_batches(filters, batch_size=1000, campaign_id=None, restart_fai @celery_app.task(name='email.process_campaign_retry') def process_campaign_retry(*args, **kwargs): - campaign_id = kwargs.get('campaign_id') campaign = NotificationCampaign.objects.get(id=campaign_id) failed_recipients = NotificationCampaignRecipient.objects.filter(campaign=campaign, status=NotificationCampaignRecipientStatus.FAILED) max_retries = campaign.metadata.get('execution', {}).get('max_retries', 3) - batch_size = campaign.metadata.get('execution', {}).get('batch_size', 1000) + batch_size = campaign.metadata.get('execution', {}).get('batch_size', DEFAULT_BATCH_SIZE) failed_recipients_count = failed_recipients.count() if not failed_recipients_count: campaign.status = NotificationCampaignStatus.COMPLETED @@ -174,26 +198,48 @@ def start_notification_campaign(campaign_id, restart_failed=False): else: manual_filters[f'{item["field"]}__{item["lookup"]}'] = [value.strip() for value in item['value'].split(',')] filters = manual_filters - tasks = [] - total_recipients = 0 - batch_size = campaign.metadata.get('execution', {}).get('batch_size', 1000) - for batch in get_filtered_batches(filters=filters, batch_size=batch_size, campaign_id=campaign_id, restart_failed=restart_failed): - tasks.append( - send_campaign_batch.s( - notification_type_name=notification_type_name, - recipients_ids=batch, - context=context, - campaign_id=campaign_id, - ) - ) - total_recipients += len(batch) + + batch_size = campaign.metadata.get('execution', {}).get('batch_size', DEFAULT_BATCH_SIZE) + batch_task_kwargs = dict( + batch_size=batch_size, + campaign_id=campaign_id, + restart_failed=restart_failed, + notification_type_name=notification_type_name, + context=context, + ) + + # Phase 1: non-spam users at/above activity threshold + active_tasks, active_count = build_campaign_group( + filters=filters, + **batch_task_kwargs, + min_activity=ACTIVITY_THRESHOLD, + ) + + # Phase 2: non-spam users below threshold (includes zero activity) + inactive_tasks, inactive_count = build_campaign_group( + filters=filters, + **batch_task_kwargs, + max_activity=ACTIVITY_THRESHOLD, + ) + + # Phase 3: confirmed spam (scheduled only after non-spam phases finish) + spam_tasks, spam_count = build_campaign_group( + filters={**filters, 'spam_status': SpamStatus.SPAM}, + **batch_task_kwargs, + exclude_spam=False, + ) + + total_recipients = active_count + inactive_count + spam_count if not restart_failed: campaign.recipient_count = total_recipients campaign.save(update_fields=['recipient_count']) - chord(tasks)( - process_campaign_retry.s(campaign_id=campaign_id) - ) + chain( + active_tasks, + inactive_tasks, + spam_tasks, + process_campaign_retry.si(campaign_id=campaign_id) + ).apply_async() @celery_app.task(name='email.send_campaign_batch', ignore_result=False)