diff --git a/admin/notifications/forms.py b/admin/notifications/forms.py index 946754415bb..9b576a0d3e9 100644 --- a/admin/notifications/forms.py +++ b/admin/notifications/forms.py @@ -1,8 +1,72 @@ from django import forms -from osf.models import NotificationType +from osf.models import NotificationType, NotificationCampaign +from website import settings +import json class NotificationTypeForm(forms.ModelForm): class Meta: model = NotificationType fields = '__all__' + + +class NotificationCampaignCreateForm(forms.ModelForm): + context = forms.CharField( + required=False, + widget=forms.Textarea(attrs={'rows': 8}), + initial='{}', + ) + + filters = forms.CharField( + required=False, + widget=forms.HiddenInput(), + initial='{}', + ) + + batch_size = forms.IntegerField( + min_value=1, + initial=settings.DEFAULT_CAMPAIGN_BATCH_SIZE, + ) + + max_retries = forms.IntegerField( + min_value=0, + 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.', + ) + + time_window = forms.IntegerField( + min_value=1, + initial=8 * 60 * 60, # 8 hours + help_text='The time in seconds before the developer reminder is sent.', + ) + + sendgrid_bulk = forms.BooleanField( + required=False, + initial=False, + ) + + class Meta: + model = NotificationCampaign + fields = ( + 'name', + 'notification_type', + ) + + def clean_context(self): + value = self.cleaned_data['context'] or '{}' + try: + return json.loads(value) + except Exception as e: + raise forms.ValidationError(e) + + def clean_filters(self): + value = self.cleaned_data['filters'] or '{}' + try: + return json.loads(value) + except Exception as e: + raise forms.ValidationError(e) diff --git a/admin/notifications/urls.py b/admin/notifications/urls.py index 236059a577e..c9c03343fd8 100644 --- a/admin/notifications/urls.py +++ b/admin/notifications/urls.py @@ -11,4 +11,10 @@ re_path(r'types_preview/(?P\d+)/$', views.NotificationTypePreview.as_view(), name='types_preview'), re_path(r'subscriptions/$', views.NotificationSubscriptionsList.as_view(), name='subscriptions_list'), re_path(r'email_tasks/$', views.EmailTasksList.as_view(), name='email_tasks_list'), + re_path(r'notification_campaigns_list/$', views.NotificationCampaignsList.as_view(), name='notification_campaigns_list'), + re_path(r'notification_campaigns_detail/(?P\d+)/$', views.NotificationCampaignDetail.as_view(), name='notification_campaigns_detail'), + re_path(r'notification_campaigns_create/$', views.NotificationCampaignCreateView.as_view(), name='notification_campaigns_create'), + re_path(r'notification_campaigns_recipients_preview/$', views.NotificationCampaignsRecipientsPreview.as_view(), name='notification_campaigns_recipients_preview'), + re_path(r'notification_campaigns_recipients_list/$', views.NotificationCampaignsRecipientsView.as_view(), name='notification_campaigns_recipients_list'), + re_path(r'notification_campaigns_start/(?P\d+)/$', views.StartNotificationCampaign.as_view(), name='notification_campaigns_start'), ] diff --git a/admin/notifications/views.py b/admin/notifications/views.py index e1c55c05f47..b22b5dd8346 100644 --- a/admin/notifications/views.py +++ b/admin/notifications/views.py @@ -1,16 +1,28 @@ +import re +import json +from collections import defaultdict +from datetime import timedelta from django.urls import reverse_lazy -from django.db.models import Q -from osf.models import NotificationSubscription, NotificationType, Notification, EmailTask -from django.views.generic import ListView, DetailView, UpdateView +from django.utils import timezone +from django.db.models import Q, F, Subquery +from django.db.models.functions import Coalesce +from django.db import models +from django.shortcuts import get_object_or_404, redirect +from django.views.generic import ListView, DetailView, UpdateView, CreateView, View +from django.contrib import messages from django.contrib.auth.mixins import PermissionRequiredMixin +from osf.models import NotificationSubscription, NotificationType, Notification, EmailTask, NotificationCampaign, OSFUser, NotificationCampaignRecipient +from osf.models.notification_campaign import NotificationCampaignStatus, NotificationCampaignRecipientStatus from django.forms.models import model_to_dict -from .forms import NotificationTypeForm -from osf.email import _render_email_html -import json -from collections import defaultdict +from .forms import NotificationTypeForm, NotificationCampaignCreateForm from mako.lexer import Lexer from mako.parsetree import ControlLine -import re +from string import Formatter +from osf.email import _render_email_html +from osf.email.notification_campaign import FILTER_PRESETS, counter_subquery, build_query +from website import settings +from urllib.parse import urlencode + def delete_selected_notifications(selected_ids): NotificationSubscription.objects.filter(id__in=selected_ids).delete() @@ -82,7 +94,7 @@ def generate_mock_json(structure, list_name=None): return result -def build_safe_context(template: str) -> dict: +def build_safe_context(template: str, subject: str) -> dict: templatenode = Lexer(text=template).parse() identifiers_location = [] for node in templatenode.get_children(): @@ -103,6 +115,9 @@ def build_safe_context(template: str) -> dict: mock_json = generate_mock_json(identifier_structure) context = {identifier: f'mock_{identifier}' for identifier in flatten_identifiers if identifier not in TEMPLATE_IDENTIFIER_BLACKLIST} context.update(mock_json) + + # subject + context.update({key: key for _, key, _, _ in Formatter().parse(subject) if key}) return context class NotificationsList(PermissionRequiredMixin, ListView): @@ -282,12 +297,12 @@ def get_context_data(self, *args, **kwargs): return kwargs else: if notification_type.is_digest_type: - inner_context = build_safe_context(notification_type.template) + inner_context = build_safe_context(notification_type.template, notification_type.subject) inner_template = _render_email_html(notification_type, ctx=inner_context, return_original_error=True) safe_context = {'notifications': [inner_template]} return_context = inner_context else: - safe_context = build_safe_context(notification_type.template) + safe_context = build_safe_context(notification_type.template, notification_type.subject) return_context = safe_context if notification_type.is_digest_type: @@ -300,6 +315,7 @@ def get_context_data(self, *args, **kwargs): except Exception as e: kwargs['rendered_template'] = f"Error rendering template: {str(e)}" + kwargs['rendered_subject'] = notification_type.subject.format(**safe_context) kwargs['context'] = json.dumps(return_context, indent=4) return kwargs @@ -327,3 +343,408 @@ class NotificationTypeChangeForm(PermissionRequiredMixin, UpdateView): def get_success_url(self, *args, **kwargs): return reverse_lazy('notifications:type_display', kwargs={'pk': self.kwargs.get('pk')}) + + +class NotificationCampaignsList(PermissionRequiredMixin, ListView): + paginate_by = 25 + template_name = 'notifications/notification_campaigns_list.html' + ordering = 'name' + permission_required = 'osf.view_notificationcampaign' + raise_exception = True + model = NotificationCampaign + + def get_queryset(self): + qs = NotificationCampaign.objects.all().order_by(self.ordering) + q = self.request.GET.get('q') + if q: + qs = qs.filter( + Q(name__icontains=q) | + Q(status__icontains=q) | + Q(notification_type__name__icontains=q) + ) + return qs + + def get_context_data(self, **kwargs): + context = super().get_context_data(**kwargs) + q = self.request.GET.get('q', '') + context['q'] = q + # append search param to pagination links + if q: + context['extra_query_params'] = f"&q={q}" + else: + context['extra_query_params'] = '' + + context['notification_campaigns'] = context['object_list'] + context['page'] = context['page_obj'] + context['active_campaign'] = NotificationCampaign.objects.filter(status=NotificationCampaignStatus.RUNNING).first() + return context + + +class NotificationCampaignDetail(PermissionRequiredMixin, DetailView): + model = NotificationCampaign + template_name = 'notifications/notification_campaigns_detail.html' + permission_required = 'osf.change_notificationcampaign' + raise_exception = True + + def get_object(self, queryset=None): + return NotificationCampaign.objects.get(id=self.kwargs.get('pk')) + + def get_context_data(self, *args, **kwargs): + notification_campaign = self.get_object() + metadata = notification_campaign.metadata or {} + + context = { + 'notification_campaign': notification_campaign, + 'display_fields': [ + ('Name', notification_campaign.name), + ('Notification Type', notification_campaign.notification_type), + ('Created By', notification_campaign.created_by), + ('Developer reminder ', 'Sent' if notification_campaign.developer_reminder_sent else ''), + ('Status', notification_campaign.get_status_display()), + ('Recipients', notification_campaign.recipient_count), + ('Sent', notification_campaign.sent_count), + ('Failed', notification_campaign.failed_count), + ('Retries', notification_campaign.retries), + ('Created', notification_campaign.created_at), + ('Started', notification_campaign.started_at), + ('Completed', notification_campaign.completed_at), + ('Updated at', notification_campaign.updated_at), + ], + 'template': notification_campaign.notification_type.template, + 'metadata': metadata, + 'filters_json': json.dumps(notification_campaign.metadata['filters']), + 'sent_filters_json': json.dumps({ + 'manual': [ + {'field': 'notificationcampaignrecipient', 'value': notification_campaign.id, 'lookup': 'campaign'}, + {'field': 'notificationcampaignrecipient', 'value': 'sent', 'lookup': 'status'} + ] + }), + 'failed_filters_json': json.dumps({ + 'manual': [ + {'field': 'notificationcampaignrecipient', 'value': notification_campaign.id, 'lookup': 'campaign'}, + {'field': 'notificationcampaignrecipient', 'value': 'failed', 'lookup': 'status'} + ] + }), + 'other_metadata': { + k: v + for k, v in metadata.items() + if k not in {'filters', 'context', 'execution', 'template'} + }, + 'allow_restart_stuck': True if timezone.now() - notification_campaign.updated_at > timedelta(minutes=15) else False, + } + + if notification_campaign.status != NotificationCampaignStatus.CREATED: + processed = notification_campaign.sent_count + notification_campaign.failed_count + pending = max(notification_campaign.recipient_count - processed, 0) + + sent_percent = ( + notification_campaign.sent_count / notification_campaign.recipient_count * 100 + if notification_campaign.recipient_count else 0 + ) + failed_percent = ( + notification_campaign.failed_count / notification_campaign.recipient_count * 100 + if notification_campaign.recipient_count else 0 + ) + + pending_percent = max(100 - sent_percent - failed_percent, 0) + elapsed = None + speed = None + eta = None + estimated_finish = None + last_activity_ago = None + failure_rate = None + + if notification_campaign.started_at: + end_time = notification_campaign.completed_at or timezone.now() + elapsed = end_time - notification_campaign.started_at + elapsed_seconds = elapsed.total_seconds() + if processed > 0 and elapsed_seconds > 0: + speed = processed / elapsed_seconds + if pending: + eta = timedelta(seconds=int(pending / speed)) + estimated_finish = timezone.now() + eta + + if notification_campaign.updated_at: + last_activity_ago = timezone.now() - notification_campaign.updated_at + if processed: + failure_rate = notification_campaign.failed_count / processed * 100 + else: + failure_rate = 0 + + context.update({ + 'processed': processed, + 'pending': pending, + 'sent_percent': sent_percent, + 'failed_percent': failed_percent, + 'pending_percent': pending_percent, + 'elapsed': elapsed, + 'speed': speed, + 'eta': eta, + 'estimated_finish': estimated_finish, + 'last_activity_ago': last_activity_ago, + 'failure_rate': failure_rate, + }) + + return context + + +LOOKUPS = { + models.CharField: { + 'exact': 'Equals', + 'iexact': 'Equals (case insensitive)', + 'contains': 'Contains', + 'icontains': 'Contains (case insensitive)', + 'not_contains': 'Does not contain', + 'not_icontains': 'Does not contain (case insensitive)', + 'startswith': 'Starts with', + 'istartswith': 'Starts with (case insensitive)', + 'endswith': 'Ends with', + 'iendswith': 'Ends with (case insensitive)', + 'regex': 'Matches regex', + 'iregex': 'Matches regex (case insensitive)', + 'in': 'In', + 'isnull': 'Is empty', + }, + models.TextField: { + 'exact': 'Equals', + 'iexact': 'Equals (case insensitive)', + 'contains': 'Contains', + 'icontains': 'Contains (case insensitive)', + 'not_contains': 'Does not contain', + 'not_icontains': 'Does not contain (case insensitive)', + 'startswith': 'Starts with', + 'istartswith': 'Starts with (case insensitive)', + 'endswith': 'Ends with', + 'iendswith': 'Ends with (case insensitive)', + 'regex': 'Matches regex', + 'iregex': 'Matches regex (case insensitive)', + 'isnull': 'Is empty', + }, + models.IntegerField: { + 'exact': 'Equals', + 'gt': 'Greater than', + 'gte': 'Greater than or equal to', + 'lt': 'Less than', + 'lte': 'Less than or equal to', + 'in': 'In', + 'isnull': 'Is empty', + }, + models.DateField: { + 'exact': 'On', + 'gt': 'After', + 'gte': 'On or after', + 'lt': 'Before', + 'lte': 'On or before', + 'isnull': 'Is empty', + }, + models.DateTimeField: { + 'exact': 'On', + 'gt': 'After', + 'gte': 'On or after', + 'lt': 'Before', + 'lte': 'On or before', + 'isnull': 'Is empty', + }, + models.BooleanField: { + 'exact': 'Is', + }, +} + + +class NotificationCampaignCreateView(CreateView): + model = NotificationCampaign + form_class = NotificationCampaignCreateForm + template_name = 'notifications/notification_campaing_create.html' + allowed_filters = [ + 'is_active', + 'is_staff', + 'username', + 'last_login', + ] + + def form_valid(self, form): + form.instance.created_by = self.request.user + + if 'manual' in form.cleaned_data['filters']: + if not form.cleaned_data['filters']['manual']['children']: + form.add_error( + 'filters', + 'Manual filters cannot be empty.' + ) + return self.form_invalid(form) + + form.instance.metadata = { + 'filters': form.cleaned_data['filters'], + 'context': form.cleaned_data['context'], + 'execution': { + 'batch_size': form.cleaned_data['batch_size'], + 'max_retries': form.cleaned_data['max_retries'], + 'activity_threshold': form.cleaned_data['activity_threshold'], + 'time_window': form.cleaned_data['time_window'], + }, + 'sendgrid_bulk': form.cleaned_data.get('sendgrid_bulk', False), + } + try: + _render_email_html(form.instance.notification_type, form.cleaned_data['context']) + except Exception as e: + form.add_error( + 'context', + f"Failed to render template: {e}", + ) + return self.form_invalid(form) + + response = super().form_valid(form) + + messages.success( + self.request, + 'Notification campaign created successfully.', + ) + + return response + + def get_success_url(self): + return reverse_lazy( + 'notifications:notification_campaigns_detail', + kwargs={'pk': self.object.pk}, + ) + + def get_context_data(self, **kwargs): + context = super().get_context_data(**kwargs) + context['notification_types'] = NotificationType.objects.order_by('name') + + filter_fields = {} + for field in [f for f in OSFUser._meta.get_fields() if f.name in self.allowed_filters]: + if not field.concrete: + continue + if type(field) not in LOOKUPS.keys(): + continue + filter_fields[field.name] = { + 'label': field.verbose_name, + 'type': field.get_internal_type().lower(), + 'lookups': LOOKUPS.get(type(field), {}) + } + context['filter_fields'] = filter_fields + context['filters'] = [] + context['predefined_filters'] = FILTER_PRESETS.keys() + context['default_context'] = json.dumps({'domain': settings.DOMAIN, 'osf_contact_email': settings.OSF_CONTACT_EMAIL}, indent=4) + return context + + +class NotificationCampaignsRecipientsPreview(PermissionRequiredMixin, ListView): + template_name = 'notifications/notification_campaing_recipients_preview.html' + permission_required = 'osf.view_osfuser' + raise_exception = True + paginate_by = 25 + + def get_queryset(self): + query = Q() + raw_filters = self.request.GET.get('filters', None) + if raw_filters: + json_filters = json.loads(raw_filters) + if predefined := json_filters.get('predefined'): + query = Q(**FILTER_PRESETS.get(predefined, {})) + else: + query = build_query(json_filters.get('manual')) + + qs = OSFUser.objects.filter(query) + qs = qs.annotate( + guid=F('guids___id'), + activity_score=Coalesce(Subquery(counter_subquery), 0) + ).order_by('-activity_score') + + return qs + + def get_context_data(self, **kwargs): + users = self.get_queryset() + + page_size = self.get_paginate_by(users) + paginator, page, query_set, is_paginated = self.paginate_queryset( + users, + page_size, + ) + + filters = self.request.GET.get('filters') + # append search param to pagination links + kwargs.update({'extra_query_params': f'&{urlencode({'filters': filters})}'}) + return super().get_context_data( + **kwargs, + page=page, + users=query_set, + paginator=paginator, + is_paginated=is_paginated, + ) + +class NotificationCampaignsRecipientsView(PermissionRequiredMixin, ListView): + template_name = 'notifications/notification_campaing_recipients_list.html' + permission_required = 'osf.view_osfuser' + raise_exception = True + paginate_by = 25 + + def get_queryset(self): + status = self.request.GET.get('notification_status', None) + campaign_id = self.request.GET.get('campaign_id', None) + if not campaign_id: + return NotificationCampaignRecipient.objects.none() + query = {'campaign_id': campaign_id} + if status == NotificationCampaignRecipientStatus.SENT: + query['status'] = NotificationCampaignRecipientStatus.SENT + elif status == NotificationCampaignRecipientStatus.FAILED: + query['status__in'] = [NotificationCampaignRecipientStatus.FAILED, NotificationCampaignRecipientStatus.SKIPPED] + + qs = NotificationCampaignRecipient.objects.filter(**query) + + return qs.annotate( + guid=F('user__guids___id') + ) + + def get_context_data(self, **kwargs): + users = self.get_queryset() + + page_size = self.get_paginate_by(users) + paginator, page, query_set, is_paginated = self.paginate_queryset( + users, + page_size, + ) + # append search param to pagination links + kwargs.update({'extra_query_params': f'¬ification_status={self.request.GET.get("notification_status")}&campaign_id={self.request.GET.get('campaign_id')}'}) + return super().get_context_data( + **kwargs, + page=page, + query_set=query_set, + paginator=paginator, + is_paginated=is_paginated, + ) + +class StartNotificationCampaign(PermissionRequiredMixin, View): + permission_required = 'osf.change_notificationcampaign' + + def post(self, request, *args, **kwargs): + notification_campaign = get_object_or_404( + NotificationCampaign, + pk=kwargs['pk'], + ) + + if NotificationCampaign.objects.filter(status=NotificationCampaignStatus.RUNNING).exclude(id=notification_campaign.id).exists(): + messages.error(request, 'Another campaign already running') + return redirect( + 'notifications:notification_campaigns_detail', + pk=notification_campaign.pk, + ) + + cancel_campaign = request.GET.get('cancel') == 'true' + if cancel_campaign: + notification_campaign.cancel() + return redirect( + 'notifications:notification_campaigns_detail', + pk=notification_campaign.pk, + ) + + restart_failed = request.GET.get('restart_failed') == 'true' + restart_stuck = request.GET.get('restart_stuck') == 'true' + + notification_campaign.start(restart_failed=restart_failed, restart_stuck=restart_stuck) + + return redirect( + 'notifications:notification_campaigns_detail', + pk=notification_campaign.pk, + ) diff --git a/admin/templates/base.html b/admin/templates/base.html index 5f645ffe267..4554c0ff1b5 100644 --- a/admin/templates/base.html +++ b/admin/templates/base.html @@ -289,7 +289,7 @@ {% endif %} {% endif %} - {% if perms.osf.view_notification or perms.osf.view_notificationtype or perms.osf.view_notificationsubscription %} + {% if perms.osf.view_notification or perms.osf.view_notificationtype or perms.osf.view_notificationsubscription or perms.osf.view_notificationcampaign %}
  • Notifications
  • @@ -307,6 +307,9 @@ {% if perms.osf.view_emailtask %}
  • Email Tasks
  • {% endif %} + {% if perms.osf.view_notificationcampaign %} +
  • Notification Campaigns
  • + {% endif %} diff --git a/admin/templates/notifications/campaign_filter_group.html b/admin/templates/notifications/campaign_filter_group.html new file mode 100644 index 00000000000..b756b1467d8 --- /dev/null +++ b/admin/templates/notifications/campaign_filter_group.html @@ -0,0 +1,22 @@ +
    +
    + Match + {{ group.operator }} +
    + +
    + {% for child in group.children %} + {% if child.children %} + {% include "notifications/campaign_filter_group.html" with group=child is_root=False %} + {% else %} +
    + {{ child.field }} + {{ child.lookup }} + {{ child.value }} +
    + {% endif %} + {% empty %} +
    No filters configured.
    + {% endfor %} +
    +
    diff --git a/admin/templates/notifications/notification_campaigns_detail.html b/admin/templates/notifications/notification_campaigns_detail.html new file mode 100644 index 00000000000..4a64edbd351 --- /dev/null +++ b/admin/templates/notifications/notification_campaigns_detail.html @@ -0,0 +1,455 @@ +{% extends "base.html" %} +{% load static %} +{% load render_bundle from webpack_loader %} + +{% block title %} + {{ notification_campaign.name }} +{% endblock title %} + +{% block content %} + + +
    + {% if messages %} +
      + {% for message in messages %} + {{ message }} + + {% endfor %} +
    + {% endif %} +
    +
    +
    +
    +

    {{ notification_campaign.name }}

    +

    + Campaign #{{ notification_campaign.id }} +

    +
    +
    + {% if notification_campaign.status != 'created' %} +
    + + Status: + {{ notification_campaign.get_status_display }} + + {% if elapsed %} +
    + Running for: {{ elapsed }} + {% endif %} + + {% if last_activity_ago %} +
    + Last activity: {{ last_activity_ago }} ago + {% endif %} +
    +
    + +
    + {{ notification_campaign.sent_count }} +
    + +
    + {{ notification_campaign.failed_count }} +
    + +
    + {{ pending }} +
    + +
    +

    Progress

    + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
    Processed{{ processed }} / {{ notification_campaign.recipient_count }}
    Pending{{ pending }}
    Success rate{{ sent_percent|floatformat:2 }}%
    Failure rate{{ failure_rate|floatformat:2 }}%
    Average speed + {% if speed %} + {{ speed|floatformat:2 }} recipients/sec + {% else %} + — + {% endif %} +
    Elapsed + {% if elapsed %} + {{ elapsed }} + {% else %} + — + {% endif %} +
    ETA + {% if eta %} + {{ eta }} + {% else %} + — + {% endif %} +
    Estimated completion + {% if estimated_finish %} + {{ estimated_finish }} + {% else %} + — + {% endif %} +
    +{% if allow_restart_stuck and notification_campaign.status == "running" %} +
    + Warning! + No activity has been detected for more than 15 minutes. + The campaign may be stuck and can be restarted. +
    +{% endif %} +{% if notification_campaign.developer_reminder_sent and notification_campaign.status == "running" %} +
    + Warning! + The campaign exceeded the expected timeframe ({{ metadata.execution.time_window }}s). A reminder was sent. +
    +{% endif %} + {% endif %} +
    + {% csrf_token %} + +
    +
    + {% csrf_token %} + +
    +
    + {% csrf_token %} + +
    + + +
    +
    +

    General

    + + {% for field, value in display_fields %} + + + + {% if field == 'Recipients' and value != 0 %} + + + {% elif field == 'Sent' and value != 0 %} + + + {% elif field == 'Failed' and value != 0 %} + + + {% else %} + + + + {% endif %} + + {% endfor %} +
    {{ field }}{{ value|safe }} + + View Recipients + + + + View Recipients + + + + View Recipients + + + +
    + {% csrf_token %} + +
    +
    +
    +
    + + + {% if metadata.filters %} +
    +
    +

    Recipient Filters

    + + {% if not "predefined" in metadata.filters %} + + {% if not "predefined" in metadata.filters %} + {% if metadata.filters.manual %} + {% include "notifications/campaign_filter_group.html" with group=metadata.filters.manual is_root=True %} + {% else %} +

    No filters configured.

    + {% endif %} + {% endif %} + + {% elif "predefined" in metadata.filters %} + + + + + + +
    Predefined Filter{{ metadata.filters.predefined }}
    + + {% endif %} + + + Preview Recipients + + +
    +
    + {% endif %} + + + {% if metadata.context %} +
    +
    +

    Context

    + + {% for key, value in metadata.context.items %} + + + + + {% endfor %} +
    {{ key }}
    {{ value }}
    +
    +
    + {% endif %} + + + {% if metadata.execution %} +
    +
    +

    Execution

    + + {% for key, value in metadata.execution.items %} + + + + + {% endfor %} +
    {{ key }}{{ value }}
    +
    +
    + {% endif %} + + {% if template or metadata.template %} +
    +
    +

    Email Template

    + + + +
    + + {% if template %} +
    +
    {{ template }}
    +
    + {% endif %} + + {% if metadata.template %} +
    +
    {{ metadata.template }}
    +
    + {% endif %} + +
    +
    +
    + {% endif %} + + + {% if other_metadata %} +
    +
    +

    Additional Metadata

    + + {% for key, value in other_metadata.items %} + + + + + {% endfor %} +
    {{ key }}
    {{ value }}
    +
    +
    + {% endif %} + +
    +{% endblock content %} + +{% block bottom_js %} + +{% endblock %} diff --git a/admin/templates/notifications/notification_campaigns_list.html b/admin/templates/notifications/notification_campaigns_list.html new file mode 100644 index 00000000000..03e3d05e457 --- /dev/null +++ b/admin/templates/notifications/notification_campaigns_list.html @@ -0,0 +1,68 @@ +{% extends "base.html" %} +{% load render_bundle from webpack_loader %} +{% load static %} +{% block title %} + List of Notification Types +{% endblock title %} +{% block content %} +

    List of Notification Campaigns

    + {% if active_campaign %} +

    Active Campaign

    + + + + + + + + + + + + + + + + + + +
    NameNotification TypeStatusStarted atCompleted at
    {{ active_campaign.name }}{{ active_campaign.notification_type.name }}{{ active_campaign.status }}{{ active_campaign.started_at }}{{ active_campaign.completed_at }}
    + {% endif %} +
    +
    + +
    +
    + + {% include "util/pagination.html" with items=page status=status %} +
    + + +
    + + + + + + + + + + + + {% for notification_capmaign in notification_campaigns %} + + + + + + + + + {% endfor %} + +
    NameNotification TypeStatusStarted atCompleted at
    {{ notification_capmaign.name }}{{ notification_capmaign.notification_type.name }}{{ notification_capmaign.status }}{{ notification_capmaign.started_at }}{{ notification_capmaign.completed_at }}
    + +{% endblock content %} diff --git a/admin/templates/notifications/notification_campaing_create.html b/admin/templates/notifications/notification_campaing_create.html new file mode 100644 index 00000000000..873fafb0635 --- /dev/null +++ b/admin/templates/notifications/notification_campaing_create.html @@ -0,0 +1,609 @@ +{% extends "base.html" %} +{% load static %} +{% load render_bundle from webpack_loader %} + +{% block title %} + Create Notification Campaign +{% endblock title %} + +{% block content %} + + +
    + {% if messages %} +
      + {% for message in messages %} + {{ message }} + + {% endfor %} +
    + {% endif %} +
    +
    + {% if form.errors %} +
    + {{ form.errors }} +
    + {% endif %} +
    +
    + +
    +
    +

    Create Notification Campaign

    +
    +
    + +
    + {% csrf_token %} + + +
    +
    +

    General

    + + + + + + + + + + +
    Campaign Name + +
    Notification Type + +
    +
    +
    + + +
    +
    +

    Recipient Filters

    + +
    + + + +
    + +
    +
    +
    + + + + + + {{ filter_fields|json_script:"filter-fields" }} + {{ filters|json_script:"initial-filters" }} + +
    +
    + +
    +
    +

    Template Context (JSON)

    + + +
    +
    + + +
    +
    +

    Execution

    + + + + + + + + + + + + + + + + + + + + + + + + + +
    Batch Size + +
    Max Retries + +
    Activity Threshold + +

    + Non-spam users at or above this activity total are sent first +

    +
    Time window + +

    + The time in seconds before the developer reminder is sent. +

    +
    Sendgrid Bulk + +
    +
    +
    + + +
    +
    + +
    +
    + +
    + + + +
    +
    + +
    +{% endblock content %} + +{% block bottom_js %} + +{% endblock %} diff --git a/admin/templates/notifications/notification_campaing_recipients_list.html b/admin/templates/notifications/notification_campaing_recipients_list.html new file mode 100644 index 00000000000..9602e8c90e8 --- /dev/null +++ b/admin/templates/notifications/notification_campaing_recipients_list.html @@ -0,0 +1,56 @@ +{% extends 'base.html' %} +{% load static %} +{% block title %} +User Search Results +{% endblock title %} +{% block content %} + {% load node_extras %} + {% include "util/pagination.html" with items=page status=status %} + {% if perms.osf.mark_spam %} +
    + {% csrf_token %} + {% endif %} + + + + + + + + + + + + + {% for record in query_set %} + + + + + + + + + {% endfor %} + +
    GUIDUsernameStatusErrorUpdated atActivity score
    + + {{ record.guid }} + + + {{record.user.username}} + + {{ record.status }} + + {{ record.error_message }} + + {{ record.updated_at }} + + {{ record.activity_score }} +
    +
    + + {% if not query_set|length %} +

    No results found

    + {% endif %} +{% endblock content %} diff --git a/admin/templates/notifications/notification_campaing_recipients_preview.html b/admin/templates/notifications/notification_campaing_recipients_preview.html new file mode 100644 index 00000000000..7d3cd7e111b --- /dev/null +++ b/admin/templates/notifications/notification_campaing_recipients_preview.html @@ -0,0 +1,56 @@ +{% extends 'base.html' %} +{% load static %} +{% block title %} +User Search Results +{% endblock title %} +{% block content %} + {% load node_extras %} + {% include "util/pagination.html" with items=page status=status %} + {% if perms.osf.mark_spam %} +
    + {% csrf_token %} + {% endif %} + + + + + + + + + + + + + {% for user in users %} + + + + + + + + + {% endfor %} + +
    GUIDUsernameFullnameDate confirmedDate disabledActivity score
    + + {{ user.guid }} + + + {{user.username}} + + {{ user.fullname }} + + {{ user.is_confirmed }} + + {{ user.is_disabled }} + + {{ user.activity_score }} +
    +
    + + {% if not users|length %} +

    No results found

    + {% endif %} +{% endblock content %} diff --git a/admin/templates/notifications/notification_type_preview.html b/admin/templates/notifications/notification_type_preview.html index e9aefbe3284..14c9d507c3a 100644 --- a/admin/templates/notifications/notification_type_preview.html +++ b/admin/templates/notifications/notification_type_preview.html @@ -1,5 +1,6 @@

    Notification Template Preview

    Rendered Template

    +

    Subject: {{ rendered_subject|safe }}

    diff --git a/admin_tests/notifications/test_campaigns.py b/admin_tests/notifications/test_campaigns.py new file mode 100644 index 00000000000..b1a04d42bfe --- /dev/null +++ b/admin_tests/notifications/test_campaigns.py @@ -0,0 +1,352 @@ +import json +import pytest +from unittest import mock + +from django.contrib.auth.models import Permission +from django.contrib.messages.storage.fallback import FallbackStorage +from django.core.exceptions import PermissionDenied +from django.test import RequestFactory +from django.urls import reverse +from django.utils import timezone +from datetime import timedelta + +from admin.notifications.forms import NotificationCampaignCreateForm +from admin.notifications.views import ( + NotificationCampaignCreateView, + NotificationCampaignDetail, + NotificationCampaignsList, + StartNotificationCampaign, +) +from admin_tests.utilities import setup_form_view +from osf.models import NotificationType +from osf.models.notification_campaign import ( + NotificationCampaign, +) +from osf_tests.factories import AuthUserFactory +from tests.base import AdminTestCase +from website import settings + + +def patch_messages(request): + setattr(request, 'session', 'session') + messages = FallbackStorage(request) + setattr(request, '_messages', messages) + + +def grant_permission(user, codename): + user.user_permissions.add(Permission.objects.get(codename=codename)) + for attr in ('_perm_cache', '_user_perm_cache', '_group_perm_cache'): + if hasattr(user, attr): + delattr(user, attr) + + +@pytest.fixture +def notification_type(): + notification_type, _ = NotificationType.objects.get_or_create( + name='blank', + defaults={'subject': 'Test', 'template': 'Hello {{ name }}'}, + ) + return notification_type + + +def _valid_form_data(notification_type, **overrides): + data = { + 'name': 'My Campaign', + 'notification_type': notification_type.id, + 'context': '{"greeting": "hi"}', + 'filters': json.dumps({'predefined': 'active'}), + 'batch_size': settings.DEFAULT_CAMPAIGN_BATCH_SIZE, + 'max_retries': settings.DEFAULT_CAMPAIGN_MAX_RETRIES, + 'activity_threshold': settings.DEFAULT_CAMPAIGN_ACTIVITY_THRESHOLD, + 'sendgrid_bulk': False, + 'time_window': 8 * 60 * 60, + } + data.update(overrides) + return data + + +class TestNotificationCampaignCreateForm: + + def test_valid_form_parses_context_and_filters(self, notification_type): + form = NotificationCampaignCreateForm(data=_valid_form_data(notification_type)) + assert form.is_valid() + assert form.cleaned_data['context'] == {'greeting': 'hi'} + assert form.cleaned_data['filters'] == {'predefined': 'active'} + assert form.cleaned_data['batch_size'] == settings.DEFAULT_CAMPAIGN_BATCH_SIZE + assert form.cleaned_data['max_retries'] == settings.DEFAULT_CAMPAIGN_MAX_RETRIES + assert form.cleaned_data['activity_threshold'] == settings.DEFAULT_CAMPAIGN_ACTIVITY_THRESHOLD + assert form.cleaned_data['time_window'] == 8 * 60 * 60 + assert form.cleaned_data['sendgrid_bulk'] is False + + def test_defaults_come_from_settings(self): + form = NotificationCampaignCreateForm() + assert form.fields['batch_size'].initial == settings.DEFAULT_CAMPAIGN_BATCH_SIZE + assert form.fields['max_retries'].initial == settings.DEFAULT_CAMPAIGN_MAX_RETRIES + assert form.fields['activity_threshold'].initial == settings.DEFAULT_CAMPAIGN_ACTIVITY_THRESHOLD + assert form.fields['time_window'].initial == 8 * 60 * 60 + assert form.fields['sendgrid_bulk'].initial is False + + def test_invalid_context_json(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, context='{not-json') + ) + assert not form.is_valid() + assert 'context' in form.errors + + def test_invalid_filters_json(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, filters='[1, 2,') + ) + assert not form.is_valid() + assert 'filters' in form.errors + + def test_empty_context_and_filters_default_to_empty_dict(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, context='', filters='') + ) + assert form.is_valid() + assert form.cleaned_data['context'] == {} + assert form.cleaned_data['filters'] == {} + + def test_batch_size_must_be_at_least_one(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, batch_size=0) + ) + assert not form.is_valid() + assert 'batch_size' in form.errors + + def test_max_retries_cannot_be_negative(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, max_retries=-1) + ) + assert not form.is_valid() + assert 'max_retries' in form.errors + + def test_activity_threshold_cannot_be_negative(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, activity_threshold=-1) + ) + assert not form.is_valid() + assert 'activity_threshold' in form.errors + + def test_time_window_must_be_at_least_one(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, time_window=0) + ) + assert not form.is_valid() + assert 'time_window' in form.errors + + def test_time_window_is_accepted(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, time_window=3600) + ) + assert form.is_valid() + assert form.cleaned_data['time_window'] == 3600 + + def test_name_is_required(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, name='') + ) + assert not form.is_valid() + assert 'name' in form.errors + + def test_notification_type_is_required(self): + form = NotificationCampaignCreateForm( + data={ + 'name': 'My Campaign', + 'context': '{}', + 'filters': '{}', + 'batch_size': 10, + 'max_retries': 1, + 'activity_threshold': 5, + } + ) + assert not form.is_valid() + assert 'notification_type' in form.errors + + def test_sendgrid_bulk_can_be_enabled(self, notification_type): + form = NotificationCampaignCreateForm( + data=_valid_form_data(notification_type, sendgrid_bulk=True) + ) + assert form.is_valid() + assert form.cleaned_data['sendgrid_bulk'] is True + + +class TestNotificationCampaignCreateView(AdminTestCase): + + def setUp(self): + super().setUp() + self.user = AuthUserFactory() + self.notification_type, _ = NotificationType.objects.get_or_create( + name='blank', + defaults={'subject': 'Test', 'template': 'Hello'}, + ) + + def test_form_valid_persists_execution_metadata(self): + request = RequestFactory().post( + reverse('notifications:notification_campaigns_create'), + data=_valid_form_data( + self.notification_type, + batch_size=25, + max_retries=4, + activity_threshold=77, + sendgrid_bulk=True, + time_window=28800 + ), + ) + request.user = self.user + patch_messages(request) + + form = NotificationCampaignCreateForm(data=request.POST) + assert form.is_valid() + + view = setup_form_view( + NotificationCampaignCreateView(), + request, + form, + ) + view.form_valid(form) + + campaign = NotificationCampaign.objects.get(name='My Campaign') + assert campaign.created_by == self.user + assert campaign.metadata['execution'] == { + 'batch_size': 25, + 'max_retries': 4, + 'activity_threshold': 77, + 'time_window': 28800 + } + assert campaign.metadata['sendgrid_bulk'] is True + assert campaign.metadata['filters'] == {'predefined': 'active'} + assert campaign.metadata['context'] == {'greeting': 'hi'} + + @mock.patch('admin.notifications.views._render_email_html', side_effect=Exception('bad template')) + def test_form_valid_rejects_unrenderable_context(self, mock_render): + request = RequestFactory().post( + reverse('notifications:notification_campaigns_create'), + data=_valid_form_data(self.notification_type), + ) + request.user = self.user + patch_messages(request) + + form = NotificationCampaignCreateForm(data=request.POST) + assert form.is_valid() + + view = setup_form_view( + NotificationCampaignCreateView(), + request, + form, + ) + with mock.patch.object(view, 'form_invalid', return_value=mock.Mock(status_code=200)) as mock_invalid: + view.form_valid(form) + + mock_invalid.assert_called_once_with(form) + assert 'context' in form.errors + assert 'Failed to render template' in form.errors['context'][0] + assert not NotificationCampaign.objects.filter(name='My Campaign').exists() + + +class TestNotificationCampaignAdminPermissions(AdminTestCase): + + def setUp(self): + super().setUp() + self.user = AuthUserFactory() + self.notification_type, _ = NotificationType.objects.get_or_create( + name='blank', + defaults={'subject': 'Test', 'template': 'Hello'}, + ) + self.campaign = NotificationCampaign.objects.create( + name='Campaign', + notification_type=self.notification_type, + metadata={'filters': {}, 'context': {}, 'execution': {}}, + ) + + def test_list_requires_view_permission(self): + request = RequestFactory().get(reverse('notifications:notification_campaigns_list')) + request.user = self.user + + with self.assertRaises(PermissionDenied): + NotificationCampaignsList.as_view()(request) + + grant_permission(self.user, 'view_notificationcampaign') + response = NotificationCampaignsList.as_view()(request) + assert response.status_code == 200 + + def test_detail_requires_change_permission(self): + request = RequestFactory().get( + reverse('notifications:notification_campaigns_detail', kwargs={'pk': self.campaign.pk}) + ) + request.user = self.user + + with self.assertRaises(PermissionDenied): + NotificationCampaignDetail.as_view()(request, pk=self.campaign.pk) + + grant_permission(self.user, 'change_notificationcampaign') + response = NotificationCampaignDetail.as_view()(request, pk=self.campaign.pk) + assert response.status_code == 200 + + def test_detail_allow_restart_stuck_false_when_recently_updated(self): + grant_permission(self.user, 'change_notificationcampaign') + request = RequestFactory().get( + reverse('notifications:notification_campaigns_detail', kwargs={'pk': self.campaign.pk}) + ) + request.user = self.user + + response = NotificationCampaignDetail.as_view()(request, pk=self.campaign.pk) + + assert response.status_code == 200 + assert response.context_data['allow_restart_stuck'] is False + + def test_detail_allow_restart_stuck_true_when_updated_long_time_ago(self): + grant_permission(self.user, 'change_notificationcampaign') + NotificationCampaign.objects.filter(pk=self.campaign.pk).update( + updated_at=timezone.now() - timedelta(minutes=16), + ) + request = RequestFactory().get( + reverse('notifications:notification_campaigns_detail', kwargs={'pk': self.campaign.pk}) + ) + request.user = self.user + + response = NotificationCampaignDetail.as_view()(request, pk=self.campaign.pk) + + assert response.status_code == 200 + assert response.context_data['allow_restart_stuck'] is True + + def test_start_requires_change_notificationcampaign_permission(self): + + request = RequestFactory().post( + reverse('notifications:notification_campaigns_start', kwargs={'pk': self.campaign.pk}) + ) + request.user = self.user + + with self.assertRaises(PermissionDenied): + StartNotificationCampaign.as_view()(request, pk=self.campaign.pk) + + grant_permission(self.user, 'change_notificationcampaign') + with mock.patch.object(NotificationCampaign, 'start') as mock_start: + response = StartNotificationCampaign.as_view()(request, pk=self.campaign.pk) + assert response.status_code == 302 + mock_start.assert_called_once_with(restart_failed=False, restart_stuck=False) + + def test_start_rejects_when_another_campaign_is_running(self): + from osf.models.notification_campaign import NotificationCampaignStatus + + grant_permission(self.user, 'change_notificationcampaign') + NotificationCampaign.objects.create( + name='Already Running', + notification_type=self.notification_type, + status=NotificationCampaignStatus.RUNNING, + metadata={'filters': {}, 'context': {}, 'execution': {}}, + ) + request = RequestFactory().post( + reverse('notifications:notification_campaigns_start', kwargs={'pk': self.campaign.pk}) + ) + request.user = self.user + patch_messages(request) + + with mock.patch.object(NotificationCampaign, 'start') as mock_start: + response = StartNotificationCampaign.as_view()(request, pk=self.campaign.pk) + + assert response.status_code == 302 + mock_start.assert_not_called() + self.campaign.refresh_from_db() + assert self.campaign.status == NotificationCampaignStatus.CREATED diff --git a/notifications.yaml b/notifications.yaml index e3b7b286a30..6161abd9fc6 100644 --- a/notifications.yaml +++ b/notifications.yaml @@ -801,3 +801,10 @@ notification_types: object_content_type_model_name: abstractnode template: 'website/templates/empty.html.mako' tests: [] + + - name: blank + subject: 'PLACEHOLDER FOR EMAIL SUBJECT' + __docs__: ... + object_content_type_model_name: osfuser + template: 'website/templates/blank.html.mako' + tests: [] diff --git a/osf/email/notification_campaign.py b/osf/email/notification_campaign.py new file mode 100644 index 00000000000..b33a3963af2 --- /dev/null +++ b/osf/email/notification_campaign.py @@ -0,0 +1,399 @@ +import logging + +from osf.models import NotificationType, NotificationTypeEnum, OSFUser, UserActivityCounter, Email +from osf.models.spam import SpamStatus +from django.db import transaction +from django.db.models import OuterRef, Subquery, Case, When, CharField, Count, Q +from django.db.models.functions import Coalesce +from framework.celery_tasks import app as celery_app +from celery import 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 +from itertools import batched + +logger = logging.getLogger(__name__) + +BULK_CREATE_SIZE = 5000 +FILTER_PRESETS = { + 'all': {}, + 'active': {'is_active': True}, + 'internal': {'is_active': True, 'is_staff': True, 'username__endswith': '@cos.io'}, +} + +first_email_subquery = ( + Email.objects + .filter(user=OuterRef('user_id')) + .values('address')[:1] +) + + +counter_subquery = ( + UserActivityCounter.objects + .filter(_id=OuterRef('guids___id')) + .values('total')[:1] +) + + +def build_query(node): + """ + Convert a filter tree into a Django Q object. + """ + + if 'field' in node: + value = node['value'] + lookup = node['lookup'] + negated_lookups = { # not native Django field lookups + 'not_contains': 'contains', + 'not_icontains': 'icontains', + } + + if lookup == 'in': + value = [v.strip() for v in value.split(',')] + + if lookup in negated_lookups: + return ~Q(**{ + f'{node["field"]}__{negated_lookups[lookup]}': value + }) + + return Q(**{ + f'{node["field"]}__{lookup}': value + }) + + operator = node.get('operator', 'AND') + children = node.get('children', []) + + if not children: + return Q() + + query = build_query(children[0]) + + for child in children[1:]: + if operator == 'AND': + query &= build_query(child) + else: + query |= build_query(child) + + return query + +def create_campaign_recipients(filters, campaign_id): + qs = ( + OSFUser.objects + .filter(filters) + .annotate(activity_score=Coalesce(Subquery(counter_subquery), 0)) + .values_list( + 'id', + 'activity_score', + ) + ) + + for rows in batched(qs.iterator(chunk_size=BULK_CREATE_SIZE), BULK_CREATE_SIZE): + NotificationCampaignRecipient.objects.bulk_create( + [ + NotificationCampaignRecipient( + campaign_id=campaign_id, + user_id=user_id, + activity_score=activity_score, + ) + for user_id, activity_score in rows + ], + ignore_conflicts=True + ) + + +def get_campaign_recipient_batches( + campaign_id, + batch_size, + restart_failed=False, + min_activity=None, + max_activity=None, + spam=None, +): + qs = NotificationCampaignRecipient.objects.filter( + campaign_id=campaign_id, + ) + + if restart_failed: + qs = qs.filter(status=NotificationCampaignRecipientStatus.FAILED) + else: + qs = qs.filter(status=NotificationCampaignRecipientStatus.PENDING) + + # Minimum and maximum activity are mutually exclusive and use the same threshold. + if min_activity is not None: + qs = qs.filter(activity_score__gte=min_activity) + + if max_activity is not None: + qs = qs.filter(activity_score__lt=max_activity) + + if spam is True: + qs = qs.filter(user__spam_status=SpamStatus.SPAM) + elif spam is False: + qs = qs.exclude(user__spam_status=SpamStatus.SPAM) + + yield from batched( + qs.values_list('id', flat=True).iterator(chunk_size=batch_size), + batch_size, + ) + +def build_campaign_group( + campaign_id, + batch_size, + restart_failed=False, + min_activity=None, + max_activity=None, + spam=None, + **send_kwargs, +): + tasks = [] + + for batch in get_campaign_recipient_batches( + campaign_id=campaign_id, + batch_size=batch_size, + restart_failed=restart_failed, + min_activity=min_activity, + max_activity=max_activity, + spam=spam, + ): + tasks.append( + send_campaign_batch.si( + recipients_ids=batch, + campaign_id=campaign_id, + **send_kwargs, + ) + ) + + return group(tasks) + + +def get_campaign_recipient_stats(campaign_id): + return NotificationCampaignRecipient.objects.filter( + campaign_id=campaign_id + ).aggregate( + recipient_count=Count('id'), + sent_count=Count( + 'id', + filter=Q(status=NotificationCampaignRecipientStatus.SENT), + ), + failed_count=Count( + 'id', + filter=Q( + status__in=[ + NotificationCampaignRecipientStatus.FAILED, + NotificationCampaignRecipientStatus.SKIPPED, + ] + ), + ), + ) + +@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) + if kwargs.get('run_id') != campaign.run_id: + return + + final_status = NotificationCampaignStatus.COMPLETED + + if campaign.status != NotificationCampaignStatus.CANCELLED: + failed_recipients = NotificationCampaignRecipient.objects.filter(campaign=campaign, status=NotificationCampaignRecipientStatus.FAILED) + 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 failed_recipients_count: + if campaign.retries < max_retries: + message = ( + f"[Notification Campaign] Retrying " + f"{failed_recipients_count} failed recipients for campaign {campaign_id}" + ) + logger.info(message) + sentry.log_message(message) + campaign.retries += 1 + campaign.save(update_fields=['retries']) + retry_group = build_campaign_group( + batch_size=batch_size, + campaign_id=campaign_id, + restart_failed=True, + notification_type_name=campaign.notification_type.name, + context=campaign.metadata.get('context', {}), + run_id=campaign.run_id, + ) + chain( + retry_group, + process_campaign_retry.si(campaign_id=campaign_id, run_id=campaign.run_id), + ).apply_async() + return + + final_status = NotificationCampaignStatus.PARTIALLY_COMPLETED + else: + message = f'[Notification Campaign] Campaign {campaign_id} {campaign.name} was cancelled.' + logger.info(message) + sentry.log_message(message) + + # Refresh in case the campaign was cancelled while we were running. + campaign.refresh_from_db(fields=['status', 'completed_at']) + + # Sync statistics regardless of status. + stats = get_campaign_recipient_stats(campaign_id) + campaign.recipient_count = stats['recipient_count'] + campaign.sent_count = stats['sent_count'] + campaign.failed_count = stats['failed_count'] + + if campaign.completed_at is None: + campaign.completed_at = timezone.now() + + # Don't overwrite CANCELLED. + if campaign.status != NotificationCampaignStatus.CANCELLED: + campaign.status = final_status + + campaign.save() + + +@celery_app.task(name='email.start_notification_campaign') +def start_notification_campaign(campaign_id, restart_failed=False, restart_stuck=False): + campaign = NotificationCampaign.objects.get(id=campaign_id) + filters = campaign.metadata.get('filters', {}) + context = campaign.metadata.get('context', {}) + notification_type_name = campaign.notification_type.name + + if hasattr(NotificationTypeEnum, notification_type_name): + del getattr(NotificationTypeEnum, notification_type_name).instance + + if predefined_filter_name := filters.get('predefined'): + filters = Q(**FILTER_PRESETS.get(predefined_filter_name, {})) + else: + filters = build_query(filters.get('manual', [])) + + if not restart_failed and not restart_stuck: + create_campaign_recipients(filters=filters, campaign_id=campaign_id) + campaign.recipient_count = NotificationCampaignRecipient.objects.filter(campaign_id=campaign_id).count() + campaign.save() + + 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, + run_id=campaign.run_id + ) + + workflow = [] + high_activity_tasks = build_campaign_group( + min_activity=activity_threshold, + spam=False, + **batch_task_kwargs + ) + if high_activity_tasks: + workflow.append(high_activity_tasks) + + low_activity_tasks = build_campaign_group( + max_activity=activity_threshold, + spam=False, + **batch_task_kwargs + ) + if low_activity_tasks: + workflow.append(low_activity_tasks) + + spam_users_tasks = build_campaign_group( + spam=True, + **batch_task_kwargs + ) + if spam_users_tasks: + workflow.append(spam_users_tasks) + + chain(*workflow, process_campaign_retry.si(campaign_id=campaign_id, run_id=campaign.run_id)).apply_async() + + +@celery_app.task(name='email.send_campaign_batch', ignore_result=False) +def send_campaign_batch(context, recipients_ids, notification_type_name='blank', campaign_id=None, run_id=None): + campaign = NotificationCampaign.objects.get(id=campaign_id) + if campaign.run_id != run_id: + return + if campaign.status == NotificationCampaignStatus.CANCELLED: + logger.warning(f"Campaign {campaign_id} was cancelled") + return + if hasattr(NotificationTypeEnum, notification_type_name): + notification_type = getattr(NotificationTypeEnum, notification_type_name).instance + else: + notification_type = NotificationType.objects.filter( + name=notification_type_name + ).first() # TODO cache + + if notification_type is None: + if campaign.status != NotificationCampaignStatus.FAILED: + campaign.status = NotificationCampaignStatus.FAILED + campaign.save() + return + + execution_time_window = campaign.metadata.get('execution', {}).get('time_window', 8 * 60 * 60) + if campaign.started_at < timezone.now() - timedelta(seconds=execution_time_window): + if not campaign.developer_reminder_sent: + message = f'[Notification Campaign] Campaign {campaign_id} exceeded its execution time window ({execution_time_window}s).' + logger.warning(message) + sentry.log_message(message) + campaign.developer_reminder_sent = True + campaign.save() + + recipients_qs = NotificationCampaignRecipient.objects.filter(id__in=recipients_ids).select_related('user') + recipient_records = [] + recipients_qs_annotated = recipients_qs.annotate( + recipient_address=Case( + When(user__username__contains='@', then='user__username'), + default=Subquery(first_email_subquery), + output_field=CharField(), + ) + ) + valid_emails_qs = recipients_qs_annotated.exclude(recipient_address__isnull=True) + invalid_emails_qs = recipients_qs_annotated.filter(recipient_address__isnull=True) + invalid_emails_qs.update(status=NotificationCampaignRecipientStatus.SKIPPED, error_message='Invalid email address') + + if campaign.metadata.get('sendgrid_bulk', False): + recipient_emails = list(valid_emails_qs.values_list('recipient_address', flat=True)) + try: + send_email_with_send_grid(to_addr=recipient_emails, notification_type=notification_type, context=context) + valid_emails_qs.update(status=NotificationCampaignRecipientStatus.SENT, error_message=None) + except Exception as exc: + message = f'[Notification Campaign] Campaign {campaign_id} sendgrid bulk request failed. {str(exc)}' + logger.error(message) + sentry.log_exception(message) + valid_emails_qs.update(status=NotificationCampaignRecipientStatus.FAILED, error_message=str(exc)) + + else: + for recipient in valid_emails_qs: + try: + notification_type.emit( + user=recipient.user, + event_context=context, + save=False, # Too many write operations + ) + + recipient.status = NotificationCampaignRecipientStatus.SENT + recipient.error_message = None + recipient_records.append(recipient) + except Exception as exc: + message = f'[Notification Campaign] Campaign {campaign_id} sendgrid request failed for user {recipient.user.username}. {str(exc)}' + logger.error(message) + sentry.log_exception(message) + + recipient.status = NotificationCampaignRecipientStatus.FAILED + recipient.error_message = str(exc) + recipient_records.append(recipient) + + NotificationCampaignRecipient.objects.bulk_update(recipient_records, ['status', 'error_message']) + + # Lock the campaign row so concurrent batches cannot + # overwrite counters with a stale aggregate snapshot + with transaction.atomic(): + notification_campaign = NotificationCampaign.objects.select_for_update().get(pk=campaign_id) + stats = get_campaign_recipient_stats(campaign_id) + notification_campaign.sent_count = stats['sent_count'] + notification_campaign.failed_count = stats['failed_count'] + notification_campaign.save(update_fields=['sent_count', 'failed_count']) + + logger.info('Batch finished') # TODO: add/update logs diff --git a/osf/management/commands/create_confirmed_test_users.py b/osf/management/commands/create_confirmed_test_users.py new file mode 100644 index 00000000000..99e1861fcd8 --- /dev/null +++ b/osf/management/commands/create_confirmed_test_users.py @@ -0,0 +1,182 @@ +import logging +import random + +from django.core.management.base import BaseCommand +from django.db import transaction +from django.utils import timezone + +from osf.models import OSFUser, SpamStatus, UserActivityCounter +from website.app import setup_django +from website.security import random_string + +setup_django() + +logger = logging.getLogger(__name__) + + +def generate_users(prefix, suffix, domain, total, start=1): + """Generate usernames paired with full names: prefix+NNNN+suffix@domain + """ + if start + total > 9999: + raise Exception(f'The max numbered user cannot be greater than 9999: start({start}) + total({total}) = {start + total}!') + return { + f'{prefix}+{str(i).zfill(4)}+{suffix}@{domain}': f'{prefix}{str(i).zfill(4)} {suffix}{str(i).zfill(4)}' + for i in range(start, total + start) + } + + +def create_confirmed_test_users( + prefix, + suffix='enter', + domain='cos.io', + total=100, + start=1, + password=None, + set_activity=True, + no_email=False, + flagged=False, + dry_run=False +): + """Create a given number of confirmed users, with generated usernames and full names. For each created user, + optionally creates a matching UserActivityCounter entry with a random total between 1 and 100. + """ + created_user_ids = [] + username_to_fullname = generate_users(prefix, suffix, domain, total, start) + + for raw_username, fullname in username_to_fullname.items(): + username = raw_username.lower().strip() + user_password = password or random_string(16) + if dry_run: + logger.info(f'Dry run: would create confirmed user "{username}" ({fullname}).') + continue + try: + user = OSFUser.create_confirmed( + username=username, + password=user_password, + fullname=fullname, + ) + except Exception as e: + logger.error(f'Failed to create confirmed user "{username}" ({fullname}): error={e}.') + continue + user.accepted_terms_of_service = timezone.now() + if no_email: + user.emails.all().delete() + user.username = user._id + if flagged: + user.spam_status = SpamStatus.FLAGGED + user.save() + logger.info(f'Created confirmed user "{username}" ({fullname})') + created_user_ids.append(user._id) + if set_activity: + UserActivityCounter.objects.bulk_create( + [ + UserActivityCounter( + _id=_id, + action={}, + date={}, + total=random.randint(1, 100) + ) + for _id in created_user_ids + ], + ignore_conflicts=True, + ) + logger.info(f'Done. Created {len(created_user_ids)} user(s); skipped {len(username_to_fullname) - len(created_user_ids)}.') + + +class Command(BaseCommand): + help = '''Create a given number of confirmed users, with generated usernames and full names. + + python3 manage.py create_confirmed_test_users --prefix longze --suffix enter --domain cos.io --total 100 + ''' + + def add_arguments(self, parser): + super().add_arguments(parser) + parser.add_argument( + '--prefix', + type=str, + required=True, + help='Prefix for generated usernames and full names', + ) + parser.add_argument( + '--suffix', + type=str, + default='enter', + help='Suffix for generated usernames and full names', + ) + parser.add_argument( + '--domain', + type=str, + default='cos.io', + help='Domain for generated usernames', + ) + parser.add_argument( + '--total', + type=int, + default=100, + help='Total number of users to create', + ) + parser.add_argument( + '--start', + type=int, + default=1, + help='The starting number of the first user to create', + ) + parser.add_argument( + '--password', + type=str, + dest='password', + help='Password to set for every created user.', + ) + parser.add_argument( + '--no-activity', + action='store_true', + dest='no_activity', + help='Skip setting a random activity total for the created users', + ) + parser.add_argument( + '--no-email', + action='store_true', + dest='no_email', + help='Set username to guid and remove all emails', + ) + parser.add_argument( + '--flagged', + action='store_true', + dest='flagged', + help='Flag user (flagged spam but not confirmed)', + ) + parser.add_argument( + '--dry', + action='store_true', + dest='dry_run', + help='Dry run; log what would happen without creating any users', + ) + + def handle(self, *args, **options): + prefix = options.get('prefix') + suffix = options.get('suffix') + domain = options.get('domain') + total = options.get('total') + start = options.get('start') + password = options.get('password') + set_activity = not options.get('no_activity', False) + no_email = options.get('no_email', False) + flagged = options.get('flagged', False) + dry_run = options.get('dry_run', False) + + if dry_run: + logger.info('This is a dry run; no users will be created.') + + with transaction.atomic(): + create_confirmed_test_users( + prefix, + suffix=suffix, + domain=domain, + total=total, + start=start, + password=password, + set_activity=set_activity, + no_email=no_email, + flagged=flagged, + dry_run=dry_run, + ) diff --git a/osf/migrations/0045_project_enter.py b/osf/migrations/0045_project_enter.py new file mode 100644 index 00000000000..77784fa79d7 --- /dev/null +++ b/osf/migrations/0045_project_enter.py @@ -0,0 +1,68 @@ +# Generated by Django 4.2.26 on 2026-07-23 09:30 + +from django.conf import settings +from django.db import migrations, models +import django.db.models.deletion + + +class Migration(migrations.Migration): + + dependencies = [ + ('osf', '0044_notification_scheduled'), + ] + + operations = [ + migrations.CreateModel( + name='NotificationCampaign', + fields=[ + ('id', models.AutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), + ('created_at', models.DateTimeField(auto_now_add=True)), + ('updated_at', models.DateTimeField(auto_now=True)), + ('run_id', models.UUIDField(blank=True, null=True, unique=True)), + ('name', models.CharField(max_length=255)), + ('status', models.CharField(choices=[('created', 'Created'), ('running', 'Running'), ('completed', 'Completed'), ('partially_completed', 'Partially Completed'), ('failed', 'Failed'), ('cancelled', 'Cancelled'), ('ended', 'Ended')], default='created', max_length=20)), + ('started_at', models.DateTimeField(blank=True, null=True)), + ('completed_at', models.DateTimeField(blank=True, null=True)), + ('metadata', models.JSONField(blank=True, default=dict)), + ('recipient_count', models.PositiveIntegerField(default=0)), + ('sent_count', models.PositiveIntegerField(default=0)), + ('failed_count', models.PositiveIntegerField(default=0)), + ('retries', models.PositiveIntegerField(default=0)), + ('developer_reminder_sent', models.BooleanField(default=False)), + ('created_by', models.ForeignKey(null=True, on_delete=django.db.models.deletion.SET_NULL, to=settings.AUTH_USER_MODEL)), + ('notification_type', models.ForeignKey(on_delete=django.db.models.deletion.PROTECT, to='osf.notificationtype')), + ], + ), + migrations.CreateModel( + name='NotificationCampaignRecipient', + fields=[ + ('id', models.AutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), + ('updated_at', models.DateTimeField(auto_now=True)), + ('status', models.CharField(choices=[('pending', 'Pending'), ('sent', 'Sent'), ('failed', 'Failed'), ('skipped', 'Skipped'), ('postponed', 'Postponed')], db_index=True, default='pending', max_length=20)), + ('error_message', models.TextField(blank=True, null=True)), + ('activity_score', models.IntegerField(default=0)), + ('campaign', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='recipients', to='osf.notificationcampaign')), + ('user', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, to=settings.AUTH_USER_MODEL)), + ], + options={ + 'ordering': ['-activity_score', '-user__date_registered', 'user_id'], + }, + ), + migrations.AddField( + model_name='osfuser', + name='received_notification_campaigns', + field=models.ManyToManyField(through='osf.NotificationCampaignRecipient', to='osf.notificationcampaign'), + ), + migrations.AddIndex( + model_name='notificationcampaignrecipient', + index=models.Index(fields=['campaign', '-activity_score', 'user'], name='campaign_order_idx'), + ), + migrations.AddIndex( + model_name='notificationcampaignrecipient', + index=models.Index(fields=['campaign', 'status', '-activity_score', 'user'], name='campaign_status_order_idx'), + ), + migrations.AlterUniqueTogether( + name='notificationcampaignrecipient', + unique_together={('campaign', 'user')}, + ), + ] diff --git a/osf/migrations/__init__.py b/osf/migrations/__init__.py index 95bfd49a76d..344cf381f5a 100644 --- a/osf/migrations/__init__.py +++ b/osf/migrations/__init__.py @@ -65,6 +65,8 @@ def get_admin_read_permissions(): 'view_notificationtype', 'view_notificationsubscription', 'view_emailtask', + 'view_notificationcampaign', + 'view_notificationcampaignrecipient', ]) @@ -116,6 +118,7 @@ def get_admin_write_permissions(): 'delete_notificationsubscription', 'change_emailtask', 'delete_emailtask', + 'change_notificationcampaign', ]) diff --git a/osf/models/__init__.py b/osf/models/__init__.py index 918ca9aa009..90284dd0d2e 100644 --- a/osf/models/__init__.py +++ b/osf/models/__init__.py @@ -67,6 +67,7 @@ from .notification_subscription import NotificationSubscription from .notification_type import NotificationType, NotificationTypeEnum from .notification import Notification +from .notification_campaign import NotificationCampaign, NotificationCampaignRecipient from .oauth import ( ApiOAuth2Application, diff --git a/osf/models/notification_campaign.py b/osf/models/notification_campaign.py new file mode 100644 index 00000000000..4b03337f25e --- /dev/null +++ b/osf/models/notification_campaign.py @@ -0,0 +1,150 @@ +from django.db import models, transaction +from django.utils import timezone +import uuid + + +class NotificationCampaignStatus(models.TextChoices): + CREATED = 'created', 'Created' + RUNNING = 'running', 'Running' + COMPLETED = 'completed', 'Completed' + PARTIALLY_COMPLETED = 'partially_completed', 'Partially Completed' + FAILED = 'failed', 'Failed' + CANCELLED = 'cancelled', 'Cancelled' + ENDED = 'ended', 'Ended' + +class NotificationCampaignRecipientStatus(models.TextChoices): + PENDING = 'pending', 'Pending' + SENT = 'sent', 'Sent' + FAILED = 'failed', 'Failed' + SKIPPED = 'skipped', 'Skipped' + POSTPONED = 'postponed', 'Postponed' + + +class NotificationCampaign(models.Model): + created_at = models.DateTimeField(auto_now_add=True) + updated_at = models.DateTimeField(auto_now=True) + run_id = models.UUIDField(null=True, blank=True, unique=True) + + name = models.CharField(max_length=255) + + notification_type = models.ForeignKey( + 'NotificationType', + on_delete=models.PROTECT, + ) + created_by = models.ForeignKey( + 'OSFUser', + null=True, + on_delete=models.SET_NULL, + ) + + status = models.CharField( + max_length=20, + choices=NotificationCampaignStatus.choices, + default=NotificationCampaignStatus.CREATED, + ) + + started_at = models.DateTimeField(null=True, blank=True) + completed_at = models.DateTimeField(null=True, blank=True) + + metadata = models.JSONField(default=dict, blank=True) + + # metadata structure: + # { + # "filters": { + # ... + # }, + # "context": { + # ... + # }, + # "execution": { + # "batch_size": , + # "max_retries": , + # "activity_threshold": , + # }, + # "template": , + # } + + recipient_count = models.PositiveIntegerField(default=0) + sent_count = models.PositiveIntegerField(default=0) + failed_count = models.PositiveIntegerField(default=0) + retries = models.PositiveIntegerField(default=0) + + developer_reminder_sent = models.BooleanField(default=False) + + def cancel(self): + self.status = NotificationCampaignStatus.CANCELLED + if self.completed_at is None: + self.completed_at = timezone.now() + self.save() + + def start(self, restart_failed=False, restart_stuck=False): + from osf.email.notification_campaign import start_notification_campaign + self.status = NotificationCampaignStatus.RUNNING + self.started_at = timezone.now() + self.run_id = uuid.uuid4() + if not restart_failed and not restart_stuck: + self.recipient_count = 0 + self.sent_count = 0 + self.failed_count = 0 + self.retries = 0 + self.metadata.update({'template': self.notification_type.template}) + self.save() + + # run start_notification_campaign after transaction + transaction.on_commit( + lambda: start_notification_campaign.delay( + campaign_id=self.id, + restart_failed=restart_failed, + restart_stuck=restart_stuck, + ) + ) + + +class NotificationCampaignRecipient(models.Model): + campaign = models.ForeignKey( + NotificationCampaign, + on_delete=models.CASCADE, + related_name='recipients', + ) + user = models.ForeignKey( + 'OSFUser', + on_delete=models.CASCADE, + ) + updated_at = models.DateTimeField(auto_now=True) + status = models.CharField( + max_length=20, + choices=NotificationCampaignRecipientStatus.choices, + default=NotificationCampaignRecipientStatus.PENDING, + db_index=True, + ) + error_message = models.TextField(null=True, blank=True) + + activity_score = models.IntegerField(default=0) + + class Meta: + unique_together = ('campaign', 'user') + ordering = [ + '-activity_score', + '-user__date_registered', + 'user_id', + ] + + indexes = [ + models.Index( + fields=[ + 'campaign', + '-activity_score', + 'user', + ], + name='campaign_order_idx', + ), + models.Index( + fields=[ + 'campaign', + 'status', + '-activity_score', + 'user', + ], + name='campaign_status_order_idx', + ), + ] diff --git a/osf/models/notification_type.py b/osf/models/notification_type.py index f8162a08bce..9a6629e22e2 100644 --- a/osf/models/notification_type.py +++ b/osf/models/notification_type.py @@ -14,6 +14,7 @@ def get_default_frequency_choices(): class NotificationTypeEnum(str, Enum): EMPTY = 'empty' + BLANK = 'blank' # Desk notifications REVIEWS_SUBMISSION_STATUS = 'reviews_submission_status' ADDONS_BOA_JOB_FAILURE = 'addon_boa_job_failure' diff --git a/osf/models/user.py b/osf/models/user.py index 0218af34692..98461a3cf0a 100644 --- a/osf/models/user.py +++ b/osf/models/user.py @@ -58,6 +58,7 @@ from .session import UserSessionMap from .tag import Tag from .validators import validate_email, validate_social, validate_history_item +from .notification_campaign import NotificationCampaign, NotificationCampaignRecipient from osf.utils.datetime_aware_jsonfield import DateTimeAwareJSONField from osf.utils.fields import NonNaiveDateTimeField, LowercaseEmailField, ensure_str from osf.utils.names import impute_names @@ -397,6 +398,11 @@ class OSFUser(DirtyFieldsMixin, GuidMixin, BaseModel, AbstractBaseUser, Permissi notifications_configured = DateTimeAwareJSONField(default=dict, blank=True) + received_notification_campaigns = models.ManyToManyField( + NotificationCampaign, + through=NotificationCampaignRecipient, + ) + # The time at which the user agreed to our updated ToS and Privacy Policy (GDPR, 25 May 2018) accepted_terms_of_service = NonNaiveDateTimeField(null=True, blank=True) diff --git a/osf_tests/test_notification_campaign.py b/osf_tests/test_notification_campaign.py new file mode 100644 index 00000000000..a54e1f16fed --- /dev/null +++ b/osf_tests/test_notification_campaign.py @@ -0,0 +1,838 @@ +import pytest +import uuid +from datetime import timedelta +from unittest import mock + +from django.utils import timezone +from django.db.models import Q + +from osf.email.notification_campaign import ( + create_campaign_recipients, + get_campaign_recipient_batches, + get_campaign_recipient_stats, + process_campaign_retry, + send_campaign_batch, + start_notification_campaign, + build_query, +) +from osf.models import UserActivityCounter, OSFUser +from osf.models.notification_campaign import ( + NotificationCampaign, + NotificationCampaignRecipient, + NotificationCampaignRecipientStatus, + NotificationCampaignStatus, +) +from osf.models.notification_type import NotificationType +from osf.models.spam import SpamStatus +from osf_tests.factories import UserFactory + +pytestmark = pytest.mark.django_db + + +@pytest.fixture +def notification_type(): + notification_type, _ = NotificationType.objects.get_or_create(name='blank') + return notification_type + + +@pytest.fixture +def campaign(notification_type): + return NotificationCampaign.objects.create( + name='Test campaign', + notification_type=notification_type, + metadata={ + 'filters': {'manual': {'operator': 'AND', 'children': []}}, + 'context': {}, + 'execution': { + 'activity_threshold': 100, + 'batch_size': 2, + 'max_retries': 2, + }, + }, + ) + + +def _set_activity(user, total): + UserActivityCounter.objects.update_or_create( + _id=user._id, + defaults={'total': total, 'action': {}, 'date': {}}, + ) + + +def _recipient_user_ids(campaign_id, **batch_kwargs): + """Flatten all batches into an list of user ids and preserve order""" + user_ids = [] + for batch in get_campaign_recipient_batches(campaign_id=campaign_id, batch_size=1000, **batch_kwargs): + recipients = NotificationCampaignRecipient.objects.filter(id__in=batch) + by_id = {r.id: r.user_id for r in recipients} + user_ids.extend(by_id[recipient_id] for recipient_id in batch) + return user_ids + + +def _recipient_scores(campaign_id, **batch_kwargs): + """Flatten all batches into an list of activity scores and preserve order""" + scores = [] + for batch in get_campaign_recipient_batches(campaign_id=campaign_id, batch_size=1000, **batch_kwargs): + by_id = { + r.id: r.activity_score + for r in NotificationCampaignRecipient.objects.filter(id__in=batch) + } + scores.extend(by_id[recipient_id] for recipient_id in batch) + return scores + + +class TestBuildQuery: + + def test_not_contains_excludes_matching_usernames(self): + email_user = UserFactory(username='user@example.com') + plain_user = UserFactory() + plain_user.username = 'deleted user' + plain_user.save(update_fields=['username']) + + query = build_query({ + 'field': 'username', + 'lookup': 'not_contains', + 'value': '@', + }) + user_ids = set(OSFUser.objects.filter(query).values_list('id', flat=True)) + + assert plain_user.id in user_ids + assert email_user.id not in user_ids + + def test_regex_matches_usernames_without_at(self): + email_user = UserFactory(username='user@example.com') + plain_user = UserFactory() + plain_user.username = 'gdpr-deleted-id' + plain_user.save(update_fields=['username']) + + query = build_query({ + 'field': 'username', + 'lookup': 'regex', + 'value': r'^[^@]+$', + }) + user_ids = set(OSFUser.objects.filter(query).values_list('id', flat=True)) + + assert plain_user.id in user_ids + assert email_user.id not in user_ids + + def test_or_combines_usernames(self): + staging_user = UserFactory(username='tester@staging.example') + plain_user = UserFactory() + plain_user.username = 'uuid-style-name' + plain_user.save(update_fields=['username']) + other_email = UserFactory(username='other@elsewhere.example') + + query = build_query({ + 'operator': 'OR', + 'children': [ + {'field': 'username', 'lookup': 'endswith', 'value': '@staging.example'}, + {'field': 'username', 'lookup': 'not_contains', 'value': '@'}, + ], + }) + user_ids = set(OSFUser.objects.filter(query).values_list('id', flat=True)) + + assert staging_user.id in user_ids + assert plain_user.id in user_ids + assert other_email.id not in user_ids + + +class TestCreateCampaignRecipients: + + def test_creates_recipients_with_activity_scores(self, campaign): + high = UserFactory() + low = UserFactory() + zero = UserFactory() + _set_activity(high, 500) + _set_activity(low, 50) + + create_campaign_recipients( + Q(**{'id__in': [high.id, low.id, zero.id]}), + campaign_id=campaign.id, + ) + + recipients = { + r.user_id: r.activity_score + for r in NotificationCampaignRecipient.objects.filter(campaign=campaign) + } + assert recipients[high.id] == 500 + assert recipients[low.id] == 50 + assert recipients[zero.id] == 0 + assert NotificationCampaignRecipient.objects.filter(campaign=campaign).count() == 3 + + def test_respects_user_filters(self, campaign): + included = UserFactory(is_staff=True) + excluded = UserFactory(is_staff=False) + _set_activity(included, 10) + _set_activity(excluded, 10) + + create_campaign_recipients( + build_query({'operator': 'AND', 'children': [{'field': 'id', 'lookup': 'in', 'value': f'{included.id}, {excluded.id}'}, {'field': 'is_staff', 'lookup': 'exact', 'value': True}]}), + campaign_id=campaign.id, + ) + + user_ids = set( + NotificationCampaignRecipient.objects.filter(campaign=campaign).values_list('user_id', flat=True) + ) + assert user_ids == {included.id} + + def test_ordered_by_activity_score_descending(self, campaign): + older = UserFactory() + newer = UserFactory() + mid = UserFactory() + older.date_registered = timezone.now() - timedelta(days=3) + mid.date_registered = timezone.now() - timedelta(days=2) + newer.date_registered = timezone.now() - timedelta(days=1) + older.save(update_fields=['date_registered']) + mid.save(update_fields=['date_registered']) + newer.save(update_fields=['date_registered']) + + _set_activity(older, 10) + _set_activity(mid, 50) + _set_activity(newer, 50) + + create_campaign_recipients( + Q(**{'id__in': [older.id, newer.id, mid.id]}), + campaign_id=campaign.id, + ) + + scores = _recipient_scores(campaign.id) + assert scores == sorted(scores, reverse=True) + user_ids = _recipient_user_ids(campaign.id) + assert user_ids == [newer.id, mid.id, older.id] + + +class TestGetCampaignRecipientStats: + + def test_empty_campaign(self, campaign): + assert get_campaign_recipient_stats(campaign.id) == { + 'recipient_count': 0, + 'sent_count': 0, + 'failed_count': 0, + } + + def test_counts_sent_failed_and_skipped(self, campaign, notification_type): + sent = UserFactory() + failed = UserFactory() + skipped = UserFactory() + pending = UserFactory() + create_campaign_recipients( + Q(**{'id__in': [sent.id, failed.id, skipped.id, pending.id]}), + campaign_id=campaign.id, + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=sent).update( + status=NotificationCampaignRecipientStatus.SENT + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=failed).update( + status=NotificationCampaignRecipientStatus.FAILED + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=skipped).update( + status=NotificationCampaignRecipientStatus.SKIPPED + ) + + assert get_campaign_recipient_stats(campaign.id) == { + 'recipient_count': 4, + 'sent_count': 1, + 'failed_count': 2, # FAILED + SKIPPED + } + + def test_scopes_to_requested_campaign(self, campaign, notification_type): + other = NotificationCampaign.objects.create( + name='Other campaign', + notification_type=notification_type, + metadata={'filters': {}, 'context': {}, 'execution': {}}, + ) + user = UserFactory() + other_user = UserFactory() + create_campaign_recipients(build_query({'operator': 'AND', 'children': [{'field': 'id', 'lookup': 'in', 'value': f'{user.id}'}]}), campaign_id=campaign.id) + create_campaign_recipients(build_query({'operator': 'AND', 'children': [{'field': 'id', 'lookup': 'in', 'value': f'{other_user.id}'}]}), campaign_id=other.id) + NotificationCampaignRecipient.objects.filter(campaign=campaign).update( + status=NotificationCampaignRecipientStatus.SENT + ) + NotificationCampaignRecipient.objects.filter(campaign=other).update( + status=NotificationCampaignRecipientStatus.FAILED + ) + + assert get_campaign_recipient_stats(campaign.id) == { + 'recipient_count': 1, + 'sent_count': 1, + 'failed_count': 0, + } + assert get_campaign_recipient_stats(other.id) == { + 'recipient_count': 1, + 'sent_count': 0, + 'failed_count': 1, + } + + +class TestGetCampaignRecipientBatches: + + @pytest.fixture + def users_and_recipients(self, campaign): + threshold = 100 + high = UserFactory() + low = UserFactory() + zero = UserFactory() + flagged = UserFactory() + spam = UserFactory() + spam.spam_status = SpamStatus.SPAM + spam.save() + flagged.spam_status = SpamStatus.FLAGGED + flagged.save() + + _set_activity(high, 250) + _set_activity(low, 40) + _set_activity(flagged, 300) + _set_activity(spam, 900) + + create_campaign_recipients( + Q(**{'id__in': [high.id, low.id, zero.id, flagged.id, spam.id]}), + campaign_id=campaign.id, + ) + return { + 'threshold': threshold, + 'high': high, + 'low': low, + 'zero': zero, + 'flagged': flagged, + 'spam': spam, + } + + def test_high_activity_non_spam(self, campaign, users_and_recipients): + data = users_and_recipients + user_ids = _recipient_user_ids( + campaign.id, + min_activity=data['threshold'], + spam=False, + ) + assert set(user_ids) == {data['high'].id, data['flagged'].id} + + def test_low_activity_non_spam_includes_zero(self, campaign, users_and_recipients): + data = users_and_recipients + user_ids = _recipient_user_ids( + campaign.id, + max_activity=data['threshold'], + spam=False, + ) + assert set(user_ids) == {data['low'].id, data['zero'].id} + + def test_spam_only(self, campaign, users_and_recipients): + data = users_and_recipients + user_ids = _recipient_user_ids(campaign.id, spam=True) + assert user_ids == [data['spam'].id] + + def test_phases_cover_all_pending_recipients(self, campaign, users_and_recipients): + data = users_and_recipients + threshold = data['threshold'] + all_ids = ( + set(_recipient_user_ids(campaign.id, min_activity=threshold, spam=False)) + | set(_recipient_user_ids(campaign.id, max_activity=threshold, spam=False)) + | set(_recipient_user_ids(campaign.id, spam=True)) + ) + expected = {data['high'].id, data['low'].id, data['zero'].id, data['flagged'].id, data['spam'].id} + assert all_ids == expected + + def test_batches_respect_batch_size(self, campaign): + users = [UserFactory() for _ in range(5)] + for i, user in enumerate(users): + _set_activity(user, (i + 1) * 10) + + create_campaign_recipients( + Q(**{'id__in': [u.id for u in users]}), + campaign_id=campaign.id, + ) + + batches = list(get_campaign_recipient_batches(campaign_id=campaign.id, batch_size=2)) + assert [len(batch) for batch in batches] == [2, 2, 1] + flat = {recipient_id for batch in batches for recipient_id in batch} + assert len(flat) == 5 + + def test_restart_failed_only_returns_failed(self, campaign, users_and_recipients): + data = users_and_recipients + failed = NotificationCampaignRecipient.objects.get(campaign=campaign, user=data['high']) + failed.status = NotificationCampaignRecipientStatus.FAILED + failed.save(update_fields=['status']) + + pending_ids = _recipient_user_ids(campaign.id, spam=False, min_activity=data['threshold']) + assert data['high'].id not in pending_ids + + failed_ids = _recipient_user_ids(campaign.id, restart_failed=True) + assert failed_ids == [data['high'].id] + + def test_ignore_conflicts_on_duplicate_create(self, campaign): + user = UserFactory() + _set_activity(user, 10) + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + assert NotificationCampaignRecipient.objects.filter(campaign=campaign).count() == 1 + + +class TestNotificationCampaignStart: + + @mock.patch('osf.email.notification_campaign.start_notification_campaign.delay') + def test_campaign_start_sets_running_state(self, mock_delay, campaign): + with mock.patch( + 'osf.models.notification_campaign.transaction.on_commit', + side_effect=lambda callback: callback(), + ): + campaign.start() + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.RUNNING + assert campaign.run_id is not None + assert campaign.started_at is not None + assert campaign.retries == 0 + mock_delay.assert_called_once_with( + campaign_id=campaign.id, + restart_failed=False, + restart_stuck=False, + ) + + @mock.patch('osf.email.notification_campaign.start_notification_campaign.delay') + def test_campaign_start_restart_stuck_counts(self, mock_delay, campaign): + campaign.recipient_count = 10 + campaign.sent_count = 4 + campaign.failed_count = 4 + campaign.retries = 1 + campaign.save() + + with mock.patch( + 'osf.models.notification_campaign.transaction.on_commit', + side_effect=lambda callback: callback(), + ): + campaign.start(restart_stuck=True) + + campaign.refresh_from_db() + assert campaign.recipient_count == 10 + assert campaign.sent_count == 4 + assert campaign.failed_count == 0 + assert campaign.retries == 0 + mock_delay.assert_called_once_with( + campaign_id=campaign.id, + restart_failed=False, + restart_stuck=True, + ) + + @mock.patch('osf.email.notification_campaign.chain') + def test_start_creates_recipients_and_schedules_workflow(self, mock_chain, campaign): + high = UserFactory() + low = UserFactory() + spam = UserFactory() + spam.spam_status = SpamStatus.SPAM + spam.save() + _set_activity(high, 250) + _set_activity(low, 10) + + campaign.metadata['filters'] = { + 'manual': {'operator': 'AND', 'children': [{'field': 'id', 'lookup': 'in', 'value': f'{high.id},{low.id},{spam.id}'}]}, + } + campaign.run_id = uuid.uuid4() + campaign.save() + + mock_chain.return_value.apply_async = mock.Mock() + + start_notification_campaign(campaign.id) + + campaign.refresh_from_db() + recipients = { + r.user_id: r.activity_score + for r in NotificationCampaignRecipient.objects.filter(campaign=campaign) + } + assert recipients == {high.id: 250, low.id: 10, spam.id: 0} + assert campaign.recipient_count == 3 + mock_chain.assert_called_once() + mock_chain.return_value.apply_async.assert_called_once() + + @mock.patch('osf.email.notification_campaign.chain') + def test_start_restart_failed_does_not_recreate_recipients(self, mock_chain, campaign): + user = UserFactory() + _set_activity(user, 50) + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=campaign, user=user) + recipient.status = NotificationCampaignRecipientStatus.FAILED + recipient.save(update_fields=['status']) + + campaign.run_id = uuid.uuid4() + campaign.metadata['filters'] = { + 'manual': {'operator': 'AND', 'children': [{'field': 'id', 'lookup': 'in', 'value': str(user.id)}]}, + } + campaign.save() + mock_chain.return_value.apply_async = mock.Mock() + + start_notification_campaign(campaign.id, restart_failed=True) + + assert NotificationCampaignRecipient.objects.filter(campaign=campaign).count() == 1 + assert NotificationCampaignRecipient.objects.get(pk=recipient.pk).status == NotificationCampaignRecipientStatus.FAILED + + @mock.patch('osf.email.notification_campaign.chain') + def test_start_restart_stuck_does_not_recreate_recipients(self, mock_chain, campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + campaign.run_id = uuid.uuid4() + campaign.recipient_count = 1 + campaign.save() + mock_chain.return_value.apply_async = mock.Mock() + + start_notification_campaign(campaign.id, restart_stuck=True) + + assert NotificationCampaignRecipient.objects.filter(campaign=campaign).count() == 1 + campaign.refresh_from_db() + assert campaign.recipient_count == 1 + + +class TestSendCampaignBatch: + + @pytest.fixture + def running_campaign(self, campaign): + campaign.run_id = uuid.uuid4() + campaign.started_at = timezone.now() + campaign.status = NotificationCampaignStatus.RUNNING + campaign.save() + return campaign + + @mock.patch.object(NotificationType, 'emit') + def test_send_campaign_batch_marks_recipients_sent(self, mock_emit, running_campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + running_campaign.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.SENT + assert running_campaign.sent_count == 1 + mock_emit.assert_called_once() + + @mock.patch.object(NotificationType, 'emit', side_effect=Exception('send failed')) + @mock.patch('osf.email.notification_campaign.sentry.log_exception') + def test_send_campaign_batch_marks_recipients_failed(self, mock_sentry, mock_emit, running_campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + running_campaign.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.FAILED + assert 'send failed' in recipient.error_message + assert running_campaign.failed_count == 1 + + def test_send_campaign_batch_skips_invalid_email_addresses(self, running_campaign): + user = UserFactory() + user.username = 'asd' + user.save(update_fields=['username']) + user.emails.all().delete() + + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + running_campaign.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.SKIPPED + assert recipient.error_message == 'Invalid email address' + assert running_campaign.failed_count == 1 + assert running_campaign.sent_count == 0 + + @mock.patch('osf.email.notification_campaign.send_email_with_send_grid') + def test_send_campaign_batch_sendgrid_bulk_success(self, mock_sendgrid, running_campaign): + user = UserFactory() + running_campaign.metadata['sendgrid_bulk'] = True + running_campaign.save(update_fields=['metadata']) + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + running_campaign.refresh_from_db() + mock_sendgrid.assert_called_once() + assert recipient.status == NotificationCampaignRecipientStatus.SENT + assert running_campaign.sent_count == 1 + assert running_campaign.failed_count == 0 + + @mock.patch( + 'osf.email.notification_campaign.send_email_with_send_grid', + side_effect=Exception('bulk failed'), + ) + @mock.patch('osf.email.notification_campaign.sentry.log_exception') + def test_send_campaign_batch_sendgrid_bulk_failure(self, mock_sentry, mock_sendgrid, running_campaign): + user = UserFactory() + running_campaign.metadata['sendgrid_bulk'] = True + running_campaign.save(update_fields=['metadata']) + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + running_campaign.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.FAILED + assert running_campaign.failed_count == 1 + assert running_campaign.sent_count == 0 + + @mock.patch('osf.email.notification_campaign.sentry.log_message') + def test_send_campaign_batch_logs_when_time_window_exceeded(self, mock_sentry, running_campaign): + user = UserFactory() + running_campaign.started_at = timezone.now() - timedelta(seconds=9) + running_campaign.metadata['execution']['time_window'] = 8 + running_campaign.save() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + with mock.patch.object(NotificationType, 'emit'): + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + running_campaign.refresh_from_db() + assert running_campaign.developer_reminder_sent is True + mock_sentry.assert_called_once() + + def test_send_campaign_batch_marks_failed_when_notification_type_missing(self, running_campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='does-not-exist', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + running_campaign.refresh_from_db() + recipient.refresh_from_db() + assert running_campaign.status == NotificationCampaignStatus.FAILED + assert recipient.status == NotificationCampaignRecipientStatus.PENDING + + def test_send_campaign_batch_skips_stale_run_id(self, running_campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=uuid.uuid4(), + ) + + recipient.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.PENDING + + def test_send_campaign_batch_skips_cancelled_campaign(self, running_campaign): + user = UserFactory() + running_campaign.status = NotificationCampaignStatus.CANCELLED + running_campaign.save(update_fields=['status']) + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.PENDING + + @mock.patch.object(NotificationType, 'emit') + def test_send_campaign_batch_uses_fallback_email_when_username_has_no_at(self, mock_emit, running_campaign): + user = UserFactory() + user.username = 'invalid' + user.save(update_fields=['username']) + user.emails.all().delete() + user.emails.create(address='fallback@example.com') + + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=running_campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=running_campaign, user=user) + + send_campaign_batch( + context={}, + recipients_ids=[recipient.id], + notification_type_name='blank', + campaign_id=running_campaign.id, + run_id=running_campaign.run_id, + ) + + recipient.refresh_from_db() + running_campaign.refresh_from_db() + assert recipient.status == NotificationCampaignRecipientStatus.SENT + assert running_campaign.sent_count == 1 + assert running_campaign.failed_count == 0 + mock_emit.assert_called_once() + + +class TestNotificationCampaignCancel: + + def test_cancel_sets_cancelled_status_and_completed_at(self, campaign): + assert campaign.completed_at is None + + campaign.cancel() + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.CANCELLED + assert campaign.completed_at is not None + + def test_cancel_does_not_overwrite_existing_completed_at(self, campaign): + completed_at = timezone.now() - timedelta(hours=1) + campaign.completed_at = completed_at + campaign.save(update_fields=['completed_at']) + + campaign.cancel() + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.CANCELLED + assert campaign.completed_at == completed_at + + +class TestProcessCampaignRetry: + + def test_process_campaign_retry_marks_completed_and_aggregates_stats(self, campaign): + sent_user = UserFactory() + skipped_user = UserFactory() + create_campaign_recipients( + Q(**{'id__in': [sent_user.id, skipped_user.id]}), + campaign_id=campaign.id, + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=sent_user).update( + status=NotificationCampaignRecipientStatus.SENT + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=skipped_user).update( + status=NotificationCampaignRecipientStatus.SKIPPED + ) + campaign.run_id = uuid.uuid4() + campaign.save(update_fields=['run_id']) + + process_campaign_retry(campaign_id=campaign.id, run_id=campaign.run_id) + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.COMPLETED + assert campaign.recipient_count == 2 + assert campaign.sent_count == 1 + assert campaign.failed_count == 1 + assert campaign.completed_at is not None + + def test_process_campaign_retry_skips_stale_run_id(self, campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + NotificationCampaignRecipient.objects.filter(campaign=campaign).update( + status=NotificationCampaignRecipientStatus.SENT + ) + campaign.run_id = uuid.uuid4() + campaign.status = NotificationCampaignStatus.RUNNING + campaign.save() + + process_campaign_retry(campaign_id=campaign.id, run_id=uuid.uuid4()) + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.RUNNING + assert campaign.completed_at is None + assert campaign.sent_count == 0 + + @mock.patch('osf.email.notification_campaign.sentry.log_message') + def test_process_campaign_retry_keeps_cancelled_status_and_syncs_stats(self, mock_sentry, campaign): + sent_user = UserFactory() + failed_user = UserFactory() + create_campaign_recipients( + Q(**{'id__in': [sent_user.id, failed_user.id]}), + campaign_id=campaign.id, + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=sent_user).update( + status=NotificationCampaignRecipientStatus.SENT + ) + NotificationCampaignRecipient.objects.filter(campaign=campaign, user=failed_user).update( + status=NotificationCampaignRecipientStatus.FAILED + ) + campaign.run_id = uuid.uuid4() + campaign.status = NotificationCampaignStatus.CANCELLED + campaign.save() + + process_campaign_retry(campaign_id=campaign.id, run_id=campaign.run_id) + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.CANCELLED + assert campaign.recipient_count == 2 + assert campaign.sent_count == 1 + assert campaign.failed_count == 1 + assert campaign.completed_at is not None + mock_sentry.assert_called_once() + + @mock.patch('osf.email.notification_campaign.chain') + def test_process_campaign_retry_retries_failed_recipients(self, mock_chain, campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=campaign, user=user) + recipient.status = NotificationCampaignRecipientStatus.FAILED + recipient.save(update_fields=['status']) + campaign.run_id = uuid.uuid4() + campaign.retries = 0 + campaign.status = NotificationCampaignStatus.RUNNING + campaign.save() + mock_chain.return_value.apply_async = mock.Mock() + + process_campaign_retry(campaign_id=campaign.id, run_id=campaign.run_id) + + campaign.refresh_from_db() + assert campaign.retries == 1 + assert campaign.status == NotificationCampaignStatus.RUNNING + mock_chain.assert_called_once() + + def test_process_campaign_retry_marks_partially_completed_after_max_retries(self, campaign): + user = UserFactory() + create_campaign_recipients(Q(**{'id__in': [user.id]}), campaign_id=campaign.id) + recipient = NotificationCampaignRecipient.objects.get(campaign=campaign, user=user) + recipient.status = NotificationCampaignRecipientStatus.FAILED + recipient.save(update_fields=['status']) + campaign.run_id = uuid.uuid4() + campaign.retries = 2 + campaign.save() + + process_campaign_retry(campaign_id=campaign.id, run_id=campaign.run_id) + + campaign.refresh_from_db() + assert campaign.status == NotificationCampaignStatus.PARTIALLY_COMPLETED + assert campaign.failed_count == 1 + assert campaign.recipient_count == 1 + assert campaign.completed_at is not None diff --git a/website/settings/defaults.py b/website/settings/defaults.py index 4b88b16bf1f..583b9ae6cbc 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 @@ -490,6 +495,7 @@ class CeleryConfig: 'api.share.utils', 'scripts.check_manual_restart_approval', 'scripts.enhanced_stuck_registration_audit', + 'osf.email.notification_campaign', } background_migration_modules = { @@ -611,6 +617,7 @@ class CeleryConfig: 'scripts.disable_removed_beat_tasks', 'osf.management.commands.delete_withdrawn_or_failed_registration_files', 'osf.management.commands.migrate_osfmetrics_fix_6to8', + 'osf.email.notification_campaign', ) # Modules that need metrics and release requirements diff --git a/website/templates/blank.html.mako b/website/templates/blank.html.mako new file mode 100644 index 00000000000..ced2e33c787 --- /dev/null +++ b/website/templates/blank.html.mako @@ -0,0 +1,62 @@ + + + + + + + + +<%page args=" + domain='', + osf_contact_email='support@osf.io' +"/> + + + + + + + + + + + + + + + + +
    + OSF +
    + + + + +
    +
    + + + + +
    +

    + Copyright © 2026 + Center For Open Science, All rights reserved. | + + Privacy Policy + +

    +

    + Questions? + ${osf_contact_email} +

    +
    +
    + + diff --git a/website/templates/notify_base.mako b/website/templates/notify_base.mako index 816ae3fbcf0..fde50148660 100644 --- a/website/templates/notify_base.mako +++ b/website/templates/notify_base.mako @@ -23,7 +23,7 @@ node_absolute_url=None, notification_settings_url=None, osf_contact_email='support@osf.io', - year=2025 + year=2026 "/>

    - Copyright © 2025 + Copyright © 2026 Center For Open Science, All rights reserved. | Privacy Policy