diff --git a/admin/notifications/forms.py b/admin/notifications/forms.py
index 98ad5c467a2..2e8f6b96029 100644
--- a/admin/notifications/forms.py
+++ b/admin/notifications/forms.py
@@ -1,5 +1,6 @@
from django import forms
from osf.models import NotificationType, NotificationCampaign
+from website import settings
import json
@@ -24,12 +25,18 @@ class NotificationCampaignCreateForm(forms.ModelForm):
batch_size = forms.IntegerField(
min_value=1,
- initial=1000,
+ initial=settings.DEFAULT_CAMPAIGN_BATCH_SIZE,
)
max_retries = forms.IntegerField(
min_value=0,
- initial=3,
+ initial=settings.DEFAULT_CAMPAIGN_MAX_RETRIES,
+ )
+
+ activity_threshold = forms.IntegerField(
+ min_value=0,
+ initial=settings.DEFAULT_CAMPAIGN_ACTIVITY_THRESHOLD,
+ help_text='Non-spam users at or above this activity total are sent in the high-activity phase.',
)
class Meta:
diff --git a/admin/notifications/views.py b/admin/notifications/views.py
index 6815d77534c..9f58d9be79e 100644
--- a/admin/notifications/views.py
+++ b/admin/notifications/views.py
@@ -508,6 +508,7 @@ def form_valid(self, form):
'execution': {
'batch_size': form.cleaned_data['batch_size'],
'max_retries': form.cleaned_data['max_retries'],
+ 'activity_threshold': form.cleaned_data['activity_threshold'],
},
}
diff --git a/admin/templates/notifications/notification_campaing_create.html b/admin/templates/notifications/notification_campaing_create.html
index bc274bdf4cc..9b092673351 100644
--- a/admin/templates/notifications/notification_campaing_create.html
+++ b/admin/templates/notifications/notification_campaing_create.html
@@ -152,7 +152,7 @@
Execution
class="form-control"
type="number"
name="batch_size"
- value="1000"
+ value="{{ form.batch_size.initial }}"
min="1"
>
@@ -165,11 +165,27 @@ Execution
class="form-control"
type="number"
name="max_retries"
- value="3"
+ value="{{ form.max_retries.initial }}"
min="0"
>
+
+
+ | Activity Threshold |
+
+
+
+ Non-spam users at or above this activity total are sent first
+
+ |
+
diff --git a/osf/email/notification_campaign.py b/osf/email/notification_campaign.py
index e4d4aa81406..40695761687 100644
--- a/osf/email/notification_campaign.py
+++ b/osf/email/notification_campaign.py
@@ -1,14 +1,17 @@
import logging
+
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
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 datetime import timedelta
from osf.models.notification_campaign import NotificationCampaign, NotificationCampaignRecipient, NotificationCampaignStatus, NotificationCampaignRecipientStatus
from osf.email import send_email_with_send_grid
from framework import sentry
+from website import settings
logger = logging.getLogger(__name__)
@@ -48,10 +51,27 @@ 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=settings.DEFAULT_CAMPAIGN_BATCH_SIZE,
+ 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)
+
+ 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.annotate(activity_total=Coalesce(Subquery(counter_subquery), 0)).order_by('-activity_total', '-date_registered', '-id')
+ qs = qs.order_by('-activity_total', '-date_registered', '-id')
last_total = None
last_date = None
@@ -68,24 +88,48 @@ def get_filtered_batches(filters, batch_size=1000, campaign_id=None, restart_fai
)
batch = batch_qs[:batch_size]
-
- if not batch:
- break
-
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_total, last_date = rows[-1]
+
+ yield batch_ids
+
- last_id, last_total, last_date = (
- rows[-1][0],
- rows[-1][1],
- rows[-1][2],
+def build_campaign_group(
+ filters,
+ batch_size=settings.DEFAULT_CAMPAIGN_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)
- yield batch_ids
+ return group(tasks), total_recipients
FILTER_PRESETS = {
@@ -96,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)
+ max_retries = campaign.metadata.get('execution', {}).get('max_retries', settings.DEFAULT_CAMPAIGN_MAX_RETRIES)
+ batch_size = campaign.metadata.get('execution', {}).get('batch_size', settings.DEFAULT_CAMPAIGN_BATCH_SIZE)
failed_recipients_count = failed_recipients.count()
if not failed_recipients_count:
campaign.status = NotificationCampaignStatus.COMPLETED
@@ -157,26 +200,50 @@ 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)
+
+ execution = campaign.metadata.get('execution', {})
+ batch_size = execution.get('batch_size', settings.DEFAULT_CAMPAIGN_BATCH_SIZE)
+ activity_threshold = execution.get('activity_threshold', settings.DEFAULT_CAMPAIGN_ACTIVITY_THRESHOLD)
+ 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
+ high_activity_tasks, high_activity_count = build_campaign_group(
+ filters=filters,
+ **batch_task_kwargs,
+ min_activity=activity_threshold,
+ )
+
+ # Phase 2: non-spam users below threshold (includes zero activity)
+ low_activity_tasks, low_activity_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_users_tasks, spam_users_count = build_campaign_group(
+ filters={**filters, 'spam_status': SpamStatus.SPAM},
+ **batch_task_kwargs,
+ exclude_spam=False,
+ )
+
+ total_recipients = high_activity_count + low_activity_count + spam_users_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(
+ high_activity_tasks,
+ low_activity_tasks,
+ spam_users_tasks,
+ process_campaign_retry.si(campaign_id=campaign_id)
+ ).apply_async()
@celery_app.task(name='email.send_campaign_batch', ignore_result=False)
diff --git a/osf/models/notification_campaign.py b/osf/models/notification_campaign.py
index 84e635e2368..89597638b25 100644
--- a/osf/models/notification_campaign.py
+++ b/osf/models/notification_campaign.py
@@ -57,6 +57,7 @@ class NotificationCampaign(models.Model):
# "execution": {
# "batch_size": ,
# "max_retries": ,
+ # "activity_threshold": ,
# },
# "template": ,
# }
diff --git a/website/settings/defaults.py b/website/settings/defaults.py
index f87c293f245..32bddc153af 100644
--- a/website/settings/defaults.py
+++ b/website/settings/defaults.py
@@ -191,6 +191,11 @@ def parent_dir(path):
NOTIFICATIONS_CLEANUP_AGE = timedelta(weeks=12) # 3 months to clean up old notifications and email tasks
NOTIFICATIONS_CLEANUP_BATCH_SIZE = 10000 # Batch size for notifications and email tasks cleanup
+# Notification campaign execution defaults (overridable per campaign in admin metadata)
+DEFAULT_CAMPAIGN_ACTIVITY_THRESHOLD = 3 # Users at/above this activity total are scheduled in the high-activity phase
+DEFAULT_CAMPAIGN_BATCH_SIZE = 1000
+DEFAULT_CAMPAIGN_MAX_RETRIES = 3
+
# Configuration for "We miss you at OSF" email (`NotificationTypeEnum.USER_NO_LOGIN`)
# Note: 1) we can gradually increase `MAX_DAILY_NO_LOGIN_EMAILS` to 10000, 100000, etc. or set it to `None` after we
# have verified that users are not spammed by this email after NR release. 2) If we want to clean up database for those