diff --git a/addons/base/views.py b/addons/base/views.py index 04e620c4c54..3f96d763781 100644 --- a/addons/base/views.py +++ b/addons/base/views.py @@ -52,9 +52,11 @@ DraftRegistration, Guid, FileVersionUserMetadata, - FileVersion, NotificationTypeEnum + FileVersion, NotificationTypeEnum, + DownloadEvent, ) from osf.utils import permissions +from osf.utils.download_telemetry import never_breaks_downloads, record_download from osf.external.gravy_valet import request_helpers from website.profile.utils import get_profile_image_url from website.project import decorators @@ -181,6 +183,74 @@ def _download_is_from_mfr(waterbutler_data): ) +def _download_request_is_from_mfr(query_params): + """Same question as :func:`_download_is_from_mfr`, asked of a browser request. + + Here the render mode is a query param on the request itself rather than something + WaterButler reported to us. + """ + return bool( + request.headers.get('X-Cos-Mfr-Render-Request', None) or + query_params.get('mode') == 'render' + ) + + +@never_breaks_downloads +def _record_file_download(target, file_node, query_params, auth, version=None): + """Record a single-file download from the redirect view. + + Only identifiers are handed over — the size, region and materialized path are looked + up in the celery task so the download itself doesn't pay for them. + """ + if _download_request_is_from_mfr(query_params): + return + + record_download( + download_type=DownloadEvent.FILE, + resource_guid=getattr(target, '_id', '') or '', + file_id=getattr(file_node, '_id', None), + version_identifier=getattr(version, 'identifier', None), + storage_provider=getattr(file_node, 'provider', '') or '', + user_guid=getattr(getattr(auth, 'user', None), '_id', None), + ip=request.remote_addr, + source_area=query_params.get('source', ''), + tz=query_params.get('tz', ''), + ) + + +@never_breaks_downloads +def _record_zip_download(payload): + """Record a folder or project zip from the WaterButler callback. + + Zips are requested straight from WaterButler, so this callback is the only point at + which we hear about them. The user's IP and the ``source``/``tz`` link tags are + forwarded to us in ``action_meta``. + """ + metadata = payload.get('metadata') or {} + action_meta = payload.get('action_meta') or {} + + if action_meta.get('is_mfr_render'): + return + + materialized = metadata.get('materialized') or metadata.get('path') or '' + # The provider root is the whole project; anything below it is one folder. + is_whole_project = not materialized.strip('/') + + record_download( + download_type=DownloadEvent.PROJECT if is_whole_project else DownloadEvent.FOLDER_ZIP, + resource_guid=metadata.get('nid') or '', + path=materialized, + storage_provider=metadata.get('provider') or '', + size_bytes=action_meta.get('bytes_downloaded'), + zip_completed=action_meta.get('completed'), + status_code=action_meta.get('status_code'), + user_guid=(payload.get('auth') or {}).get('id'), + ip=action_meta.get('ip'), + source_area=action_meta.get('source', ''), + tz=action_meta.get('tz', ''), + ) + + def make_auth(user): if user is not None: return { @@ -483,11 +553,11 @@ def create_waterbutler_log(payload, **kwargs): with transaction.atomic(): try: auth = payload['auth'] - # Don't log download actions + # Downloads produce no NodeLog, but zips are recorded for telemetry here — + # they never pass through the redirect view where single files are caught. if payload['action'] in DOWNLOAD_ACTIONS: - guid_id = payload['metadata'].get('nid') - - node, _ = Guid.load_referent(guid_id) + if payload['action'] == 'download_zip': + _record_zip_download(payload) return {'status': 'success'} user = OSFUser.load(auth['id']) @@ -983,6 +1053,7 @@ def addon_view_or_download_file(auth, path, provider, **kwargs): })) if action == 'download': + _record_file_download(target, file_node, extras, auth, version=version) format = extras.get('format') _, extension = os.path.splitext(file_node.name) # avoid rendering files with the same format type. @@ -1045,6 +1116,8 @@ def persistent_file_download(auth, **kwargs): query_params = request.args.to_dict() + _record_file_download(file.target, file, query_params, auth) + return make_response( '', http_status.HTTP_302_FOUND, { 'Location': file.generate_waterbutler_url(**query_params), diff --git a/admin/base/settings/defaults.py b/admin/base/settings/defaults.py index 3d89c8d3d04..c255605fd55 100644 --- a/admin/base/settings/defaults.py +++ b/admin/base/settings/defaults.py @@ -85,6 +85,7 @@ 'guardian', 'waffle', 'elasticsearch_metrics.apps.ElasticsearchMetricsConfig', + 'rangefilter', # OSF 'osf', diff --git a/admin/templates/download_events/download_events.html b/admin/templates/download_events/download_events.html new file mode 100644 index 00000000000..bb1563a39ec --- /dev/null +++ b/admin/templates/download_events/download_events.html @@ -0,0 +1,609 @@ +{% extends "admin/change_list.html" %} + +{% block content %} + {% if download_events_dashboard %} + +
+
+

Download telemetry dashboard

+

Summaries and charts are scoped to the same filters and search terms as the table below.

+
+ +
+
+
Total downloads
+
{{ download_events_dashboard.summary.total_downloads }}
+
+
+
Total GB
+
{{ download_events_dashboard.summary.total_gb|floatformat:2 }}
+
+
+
Unique users
+
{{ download_events_dashboard.summary.unique_users }}
+
+
+
File vs zip (including whole-project zip) (GB)
+
{{ download_events_dashboard.split.file.gb|floatformat:2 }} / {{ download_events_dashboard.split.zip.gb|floatformat:2 }}
+
+
+
Failed zip downloads (server error)
+
{{ download_events_dashboard.summary.failed_zips }}
+
+
+
Zip outcomes (completed / cancelled / failed)
+
{{ download_events_dashboard.zip_outcomes.completed }} / {{ download_events_dashboard.zip_outcomes.cancelled }} / {{ download_events_dashboard.zip_outcomes.failed }}
+
+
+ +
+
GB over time
+ +
+ +
+
+ +
+ +
+
+
GB by storage region
+
+ Requests in range — Files: {{ download_events_dashboard.split.file.count }} + · Zips: {{ download_events_dashboard.split.zip.count }} + · Total: {{ download_events_dashboard.summary.total_downloads }} +
+ +
+ +
+
+ +
+
Downloads / GB by user region
+
+ Requests in range — Files: {{ download_events_dashboard.split.file.count }} + · Zips: {{ download_events_dashboard.split.zip.count }} + · Total: {{ download_events_dashboard.summary.total_downloads }} +
+ +
+ +
+
+ +
+
File vs Zip
+ +
+
+

+ Downloads +

+ +
+ +
+
+ +
+

+ GB +

+ +
+ +
+
+ +
+
+
+ +
+
+
Top projects by volume
+ + + + + + + + + + {% for project in download_events_dashboard.top_projects %} + + + + + + {% empty %} + + {% endfor %} + +
ProjectGBDownloads
{{ project.name }}{{ project.gb|floatformat:2 }}{{ project.downloads }}
No project-level activity in the current range.
+
+
+
Top users by volume
+ + + + + + + + + + {% for user in download_events_dashboard.top_users %} + + + + + + {% empty %} + + {% endfor %} + +
UserGBDownloads
{{ user.name }}{{ user.gb|floatformat:2 }}{{ user.downloads }}
No user activity in the current range.
+
+
+
By storage provider
+ + + + + + + + + + {% for provider in download_events_dashboard.storage_providers %} + + + + + + {% empty %} + + {% endfor %} + +
ProviderGBDownloads
{{ provider.name }}{{ provider.gb|floatformat:2 }}{{ provider.downloads }}
No activity in the current range.
+
+
+
+ {% endif %} +

+ {{ block.super }} + + + +{% endblock %} diff --git a/osf/admin.py b/osf/admin.py index d9fed50b7ff..284b8695f66 100644 --- a/osf/admin.py +++ b/osf/admin.py @@ -1,10 +1,15 @@ +from collections import defaultdict +from datetime import timedelta + from django.contrib import admin, messages from django.urls import re_path, reverse, path from django.template.response import TemplateResponse from django_extensions.admin import ForeignKeyAutocompleteAdmin from django.contrib.auth.models import Group -from django.db.models import Q, Count +from django.db.models import Q, Count, Sum, F, Min, Max +from django.db.models.functions import Trunc from django.http import HttpResponseRedirect, HttpResponse, JsonResponse +from django.utils import timezone from django.utils.html import format_html from django.shortcuts import get_object_or_404 from django import forms @@ -12,12 +17,28 @@ from django.contrib.admin import SimpleListFilter import waffle +from rangefilter.filters import DateTimeRangeFilterBuilder + from osf.external.spam.tasks import reclassify_domain_references -from osf.models import OSFUser, Node, NotableDomain, NodeLicense, NotificationType, NotificationSubscription, EmailTask, Notification +from osf.models import ( + OSFUser, + Node, + NotableDomain, + NodeLicense, + NotificationType, + NotificationSubscription, + EmailTask, + Notification, + DownloadEvent +) +from osf.models import AbstractNode from osf.models.notification_type import get_default_frequency_choices from osf.models.notable_domain import DomainReference +DASHBOARD_GROUP_NAME = 'download_telemetry' + + def list_displayable_fields(cls): return [x.name for x in cls._meta.fields if x.editable and not x.is_relation and not x.primary_key] @@ -386,6 +407,415 @@ def user(self, obj): return '(username)' user.short_description = 'User' + +# A zip that WaterButler ended on a 5xx failed through no fault of the user (the server +# buckled) -- distinct from a cancel, which ends completed=False at 200 because the headers +# already went out. This is the threshold that separates the two. +DOWNLOAD_FAILURE_MIN_STATUS = 500 + + +class DownloadOutcomeFilter(SimpleListFilter): + """Filter zips by how they ended: completed, cancelled mid-stream, or failed.""" + + title = 'download outcome' + parameter_name = 'outcome' + + COMPLETED = 'completed' + CANCELLED = 'cancelled' + FAILED = 'failed' + + def lookups(self, request, model_admin): + return [ + (self.COMPLETED, 'Completed'), + (self.CANCELLED, 'Cancelled mid-download'), + (self.FAILED, 'Failed (server error)'), + ] + + def queryset(self, request, queryset): + if self.value() == self.COMPLETED: + return queryset.filter(zip_completed=True) + if self.value() == self.FAILED: + return queryset.filter( + zip_completed=False, status_code__gte=DOWNLOAD_FAILURE_MIN_STATUS) + if self.value() == self.CANCELLED: + return queryset.filter(zip_completed=False).exclude( + status_code__gte=DOWNLOAD_FAILURE_MIN_STATUS) + return queryset + + +@admin.register(DownloadEvent) +class DownloadEventsView(admin.ModelAdmin): + change_list_template = 'download_events/download_events.html' + list_display = ( + 'resource_guid', + 'user', + 'download_type', + 'outcome', + 'zip_completed', + 'status_code', + 'path', + 'size', + 'storage_provider', + 'user_region', + 'storage_region', + 'ip', + 'source_area', + 'created' + ) + list_filter = ( + ( + 'created', + DateTimeRangeFilterBuilder( + title='date and time (UTC)', + ), + ), + 'download_type', + DownloadOutcomeFilter, + 'zip_completed', + 'storage_provider', + ) + ordering = ('-created',) + search_fields = ( + 'user__username', + 'user__fullname', + 'user__guids___id', + 'resource_guid', + 'ip', + 'path', + 'storage_provider', + 'user_region', + 'storage_region', + 'source_area' + ) + search_help_text = 'Search by username, full name, user or node guid, ip, path, storage provider, user or storage region, source area.' + + @admin.display(description='Outcome') + def outcome(self, obj): + """Human-readable end state. Single files have no outcome — they're recorded at the + redirect before any bytes move, so they never report completion.""" + if obj.zip_completed is None: + return '—' + if obj.zip_completed: + return 'Completed' + if obj.status_code and obj.status_code >= DOWNLOAD_FAILURE_MIN_STATUS: + return 'Failed' + return 'Cancelled' + + @admin.display(description='Size (GB)', ordering=F('size_bytes').desc(nulls_last=True)) + def size(self, obj): + """`size_bytes` is null when we could not determine it — that is not zero. + + Sorting puts those last rather than letting Postgres float them to the top. + """ + if obj.size_bytes is None: + return '—' + return f'{self._to_gb(obj.size_bytes)}' + + def changelist_view(self, request, extra_context=None): + for query_string in request.GET: + # when at least one of the "created" filters is set, don't override the filter values + if query_string.startswith('created__range'): + break + else: + # by default, when the page is initially loaded or "created" filter is reset + # show only events within the last hour + request.GET._mutable = True + last_hour_datetime = timezone.now() - timedelta(hours=1) + request.GET['created__range__gte_0'] = last_hour_datetime.date().strftime('%Y-%m-%d') + request.GET['created__range__gte_1'] = last_hour_datetime.time().strftime('%H:%M:%S') + + if extra_context is None: + extra_context = {} + changelist = self.get_changelist_instance(request) + extra_context['download_events_dashboard'] = self.get_dashboard_data(changelist.get_queryset(request)) + return super().changelist_view(request, extra_context=extra_context) + + def _in_dashboard_group(self, request): + """Membership in the allow-list group is the only key to this page. + + Deliberately not falling back to ``super()``/``has_perm``: ``ModelBackend`` + answers True to every permission check for a superuser, so anything that + consults it would let every superuser in — the opposite of what this + dashboard is for. + """ + user = getattr(request, 'user', None) + if user is None or not user.is_authenticated: + return False + return user.groups.filter(name=DASHBOARD_GROUP_NAME).exists() + + def has_module_permission(self, request, obj=None): + """Keeps the model off the admin index for everyone else.""" + return self._in_dashboard_group(request) + + def has_view_permission(self, request, obj=None): + """What ``changelist_view`` actually enforces — the real gate. + + ``has_module_permission`` alone only hides the link; the page itself stays + reachable by URL without this. + """ + return self._in_dashboard_group(request) + + # Append-only telemetry: nothing is editable through the admin, by anyone. + def has_add_permission(self, request): + return False + + def has_change_permission(self, request, obj=None): + return False + + def has_delete_permission(self, request, obj=None): + return False + + def _sum_bytes(self, queryset): + return queryset.aggregate(total_bytes=Sum('size_bytes'))['total_bytes'] or 0 + + def _to_gb(self, total_bytes): + return round((total_bytes or 0) / (1024**3), 2) + + def _percent(self, part, whole): + """Empty ranges are normal — the default window is the last hour.""" + if not whole: + return 0.0 + return round(part * 100 / whole, 2) + + def get_dashboard_data(self, queryset): + file_queryset = queryset.filter(download_type=DownloadEvent.FILE) + zip_queryset = queryset.exclude(download_type=DownloadEvent.FILE) + total_file_downloads = file_queryset.count() + total_zip_downloads = zip_queryset.count() + total_bytes = self._sum_bytes(queryset) + + total_downloads = queryset.count() + total_file_gb = self._to_gb(self._sum_bytes(file_queryset)) + total_zip_gb = self._to_gb(self._sum_bytes(zip_queryset)) + time_series = self._build_time_series(queryset) + storage_regions = self._build_region_breakdown(queryset, 'storage_region') + user_regions = self._build_region_breakdown(queryset, 'user_region') + # downloads and GB grouped by where the bytes came from (osfstorage vs addons) + storage_providers = self._build_region_breakdown(queryset, 'storage_provider') + + # Zip outcomes. Single files are recorded before any bytes move, so they have no + # outcome and are left out of this breakdown entirely. + completed_zips = zip_queryset.filter(zip_completed=True).count() + failed_zips = zip_queryset.filter( + zip_completed=False, status_code__gte=DOWNLOAD_FAILURE_MIN_STATUS).count() + incomplete_zips = zip_queryset.filter(zip_completed=False).count() + zip_outcomes = { + 'completed': completed_zips, + # everything that didn't complete and wasn't a server failure is a user cancel + 'cancelled': incomplete_zips - failed_zips, + 'failed': failed_zips, + } + + total_gb = self._to_gb(total_bytes) + split = { + 'file': { + 'count': total_file_downloads, + 'gb': total_file_gb, + 'count_percent': self._percent(total_file_downloads, total_downloads), + 'gb_percent': self._percent(total_file_gb, total_gb), + }, + 'zip': { + 'count': total_zip_downloads, + 'gb': total_zip_gb, + 'count_percent': self._percent(total_zip_downloads, total_downloads), + 'gb_percent': self._percent(total_zip_gb, total_gb), + }, + } + + return { + 'summary': { + 'total_downloads': total_downloads, + 'total_gb': total_gb, + 'unique_users': queryset.exclude(user_id__isnull=True).values('user_id').distinct().count(), + 'failed_zips': failed_zips, + }, + 'split': split, + 'zip_outcomes': zip_outcomes, + 'time_series': time_series, + 'storage_regions': storage_regions, + 'storage_providers': storage_providers, + 'user_regions': user_regions, + 'top_projects': self._build_top_resource_breakdown(queryset), + 'top_users': self._build_top_user_breakdown(queryset), + } + + EMPTY_TIME_SERIES = {'labels': [], 'file': [], 'zip': []} + + def _build_time_series(self, queryset): + """Bucketed GB over time, aggregated in the database. + + The range can cover millions of rows once the announcement lands, so the + bucketing is a GROUP BY rather than a pass over every event in Python. + """ + bounds = queryset.aggregate(start=Min('created'), end=Max('created')) + start, end = bounds['start'], bounds['end'] + if start is None: + return dict(self.EMPTY_TIME_SERIES) + + bucket_size = self._get_bucket_size(start, end) + step = self._get_bucket_step(bucket_size) + + rows = queryset.annotate( + bucket=Trunc('created', self._get_trunc_kind(bucket_size)) + ).values('bucket', 'download_type').annotate( + total_bytes=Sum('size_bytes'), + ) + + totals = defaultdict(lambda: {'file': 0.0, 'zip': 0.0}) + for row in rows: + key = self._floor_to_bucket(row['bucket'], bucket_size) + side = 'file' if row['download_type'] == DownloadEvent.FILE else 'zip' + totals[key][side] += (row['total_bytes'] or 0) / (1024**3) + + # walk the whole span so gaps render as zero rather than closing up + buckets = [] + current = self._floor_to_bucket(start, bucket_size) + last = self._floor_to_bucket(end, bucket_size) + while current <= last: + buckets.append(current) + current += step + + return { + 'labels': [self._format_bucket_label(bucket, bucket_size) for bucket in buckets], + 'file': [round(totals[bucket]['file'], 2) for bucket in buckets], + 'zip': [round(totals[bucket]['zip'], 2) for bucket in buckets], + } + + def _get_trunc_kind(self, bucket_size): + """15-minute buckets have no Trunc equivalent, so truncate to the hour and + let `_floor_to_bucket` split it down.""" + return { + '15m': 'minute', + '1h': 'hour', + '1d': 'day', + '1w': 'week', + }[bucket_size] + + def _get_bucket_size(self, start, end): + delta = end - start + if delta <= timedelta(hours=2): + return '15m' + if delta <= timedelta(hours=24): + return '1h' + if delta <= timedelta(days=14): + return '1d' + return '1w' + + def _get_bucket_step(self, bucket_size): + if bucket_size == '15m': + return timedelta(minutes=15) + if bucket_size == '1h': + return timedelta(hours=1) + if bucket_size == '1d': + return timedelta(days=1) + return timedelta(days=7) + + def _floor_to_bucket(self, value, bucket_size): + """Snap to the start of the containing bucket. + + The result is used as a dict key, so it has to land on exactly the same + instant whether it came from an event or from walking the axis — every + branch zeroes everything below its own resolution. + """ + if bucket_size == '15m': + return value.replace(minute=(value.minute // 15) * 15, second=0, microsecond=0) + if bucket_size == '1h': + return value.replace(minute=0, second=0, microsecond=0) + midnight = value.replace(hour=0, minute=0, second=0, microsecond=0) + if bucket_size == '1d': + return midnight + return midnight - timedelta(days=value.weekday()) + + def _format_bucket_label(self, value, bucket_size): + if bucket_size == '15m': + return value.strftime('%H:%M') + if bucket_size == '1h': + return value.strftime('%Y-%m-%d %H:%M') + if bucket_size == '1d': + return value.strftime('%Y-%m-%d') + return value.strftime('%Y-%m-%d') + + def _build_region_breakdown(self, queryset, field_name): + """Grouped in the database — the range can cover millions of rows. + + `downloads` is the total request count; `file_count` and `zip_count` split it by + request type (a zip is either a folder or a whole-project zip), so file + zip always + equals the total. + """ + rows = queryset.values(field_name).annotate( + downloads=Count('id'), + total_bytes=Sum('size_bytes'), + file_count=Count('id', filter=Q(download_type=DownloadEvent.FILE)), + zip_count=Count('id', filter=~Q(download_type=DownloadEvent.FILE)), + ) + + breakdown = defaultdict(lambda: {'downloads': 0, 'gb': 0.0, 'file_count': 0, 'zip_count': 0}) + for row in rows: + # blank and null both mean "we could not tell", so they fold together + region_name = (row[field_name] or 'Unknown').strip() or 'Unknown' + breakdown[region_name]['downloads'] += row['downloads'] + breakdown[region_name]['gb'] += (row['total_bytes'] or 0) / (1024**3) + breakdown[region_name]['file_count'] += row['file_count'] + breakdown[region_name]['zip_count'] += row['zip_count'] + + ordered = sorted(breakdown.items(), key=lambda item: item[1]['gb'], reverse=True)[:10] + max_gb = max((data['gb'] for _, data in ordered), default=0) + max_downloads = max((data['downloads'] for _, data in ordered), default=0) + return [ + { + 'name': name, + 'downloads': data['downloads'], + 'file_count': data['file_count'], + 'zip_count': data['zip_count'], + 'gb': round(data['gb'], 2), + 'gb_percent': self._percent(data['gb'], max_gb), + 'download_percent': self._percent(data['downloads'], max_downloads), + } + for name, data in ordered + ] + + def _build_top_resource_breakdown(self, queryset): + rows = queryset.exclude(resource_guid='').values('resource_guid').annotate( + gb_bytes=Sum('size_bytes'), + downloads=Count('id'), + ).order_by(F('gb_bytes').desc(nulls_last=True), '-downloads')[:10] + rows = list(rows) + + # one query for all ten, rather than one per row + guids = [row['resource_guid'] for row in rows] + titles = dict( + AbstractNode.objects.filter(guids___id__in=guids).values_list('guids___id', 'title') + ) + return [ + { + # a deleted project keeps its title, but fall back to the bare guid so + # the row still says something if it ever fails to resolve + 'name': ( + f'{titles[row["resource_guid"]]} ({row["resource_guid"]})' + if row['resource_guid'] in titles + else row['resource_guid'] + ), + 'downloads': row['downloads'], + 'gb': self._to_gb(row['gb_bytes']), + } + for row in rows + ] + + def _build_top_user_breakdown(self, queryset): + rows = queryset.exclude(user__isnull=True).values('user__username', 'user__fullname').annotate( + gb_bytes=Sum('size_bytes'), + downloads=Count('id'), + ).order_by(F('gb_bytes').desc(nulls_last=True), '-downloads')[:10] + return [ + { + 'name': row['user__fullname'] or row['user__username'] or 'Unknown user', + 'downloads': row['downloads'], + 'gb': self._to_gb(row['gb_bytes']), + } + for row in rows + ] + + admin.site.register(OSFUser, OSFUserAdmin) admin.site.register(Node, NodeAdmin) admin.site.register(NotableDomain, NotableDomainAdmin) diff --git a/osf/migrations/0045_downloadevent.py b/osf/migrations/0045_downloadevent.py new file mode 100644 index 00000000000..2ea410e1ce2 --- /dev/null +++ b/osf/migrations/0045_downloadevent.py @@ -0,0 +1,75 @@ +import django.db.models.deletion +from django.db import migrations, models + + +DASHBOARD_GROUP_NAME = 'download_telemetry' + +DASHBOARD_USERS = [ + 'sheredko.andriy@gmail.com', + 'bodintsov@exoft.net', + 'isokhan@exoft.net', + 'ykopka@exoft.net', + 'bgeiger@cos.io', + 'osmand@cos.io', + 'ramya@cos.io', + 'eric@cos.io', +] + + +def create_dashboard_group(apps, schema_editor): + """Create the allow-list group the dashboard loads against and seed it. + + The group carries no permissions of its own — membership is the only gate. + """ + Group = apps.get_model('auth', 'Group') + OSFUser = apps.get_model('osf', 'OSFUser') + group, _ = Group.objects.get_or_create(name=DASHBOARD_GROUP_NAME) + for username in DASHBOARD_USERS: + user = OSFUser.objects.filter(username=username).first() + if user: + group.user_set.add(user) + + +def remove_dashboard_group(apps, schema_editor): + Group = apps.get_model('auth', 'Group') + Group.objects.filter(name=DASHBOARD_GROUP_NAME).delete() + + +class Migration(migrations.Migration): + + dependencies = [ + ('osf', '0044_notification_scheduled'), + ] + + operations = [ + migrations.CreateModel( + name='DownloadEvent', + fields=[ + ('id', models.AutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')), + ('created', models.DateTimeField(auto_now_add=True, db_index=True)), + ('resource_guid', models.CharField(blank=True, db_index=True, default='', max_length=255)), + ('path', models.TextField(blank=True, default='')), + ('download_type', models.CharField(choices=[('file', 'Single file'), ('folder_zip', 'Folder zip'), ('project', 'Whole-project zip')], max_length=16)), + ('zip_completed', models.BooleanField(blank=True, null=True)), + ('size_bytes', models.BigIntegerField(blank=True, null=True)), + ('storage_region', models.CharField(blank=True, default='', max_length=64)), + ('user_region', models.CharField(blank=True, default='', max_length=64)), + ('ip', models.GenericIPAddressField(blank=True, null=True)), + ('source_area', models.CharField(blank=True, default='', max_length=128)), + ('user', models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='download_events', to='osf.osfuser')), + ], + ), + migrations.AddIndex( + model_name='downloadevent', + index=models.Index(fields=['created', 'download_type'], name='download_event_crt_type'), + ), + migrations.AddIndex( + model_name='downloadevent', + index=models.Index(fields=['created', 'storage_region'], name='download_event_crt_regn'), + ), + migrations.AddIndex( + model_name='downloadevent', + index=models.Index(fields=['created', 'user_region'], name='download_event_crt_user'), + ), + migrations.RunPython(create_dashboard_group, remove_dashboard_group), + ] diff --git a/osf/migrations/0046_dashboard_group_staff_access.py b/osf/migrations/0046_dashboard_group_staff_access.py new file mode 100644 index 00000000000..ca1ed5f6eb0 --- /dev/null +++ b/osf/migrations/0046_dashboard_group_staff_access.py @@ -0,0 +1,44 @@ +from django.db import migrations + + +DASHBOARD_GROUP_NAME = 'download_telemetry' + + +def grant_staff_access(apps, schema_editor): + """Let the allow-listed users through Django's admin door. + + ``AdminSite.has_permission`` rejects anyone without ``is_staff`` before any + per-model check runs, so without this the dashboard is unreachable for exactly + the people it was built for. + + This is not superuser. These accounts carry no Django permissions, so the + download events page is the only thing in the admin they can open — every other + registered model stays hidden and unviewable. Kept separate from the migration + that creates the group so it also applies to databases where that one has + already run. + """ + Group = apps.get_model('auth', 'Group') + group = Group.objects.filter(name=DASHBOARD_GROUP_NAME).first() + if group is None: + return + group.user_set.filter(is_staff=False).update(is_staff=True) + + +def revoke_staff_access(apps, schema_editor): + """Deliberately a no-op. + + Some of these accounts are OSF admins who had staff access long before this + feature; reversing the migration must not strip it from them, and we have no + record of who had it beforehand. + """ + + +class Migration(migrations.Migration): + + dependencies = [ + ('osf', '0045_downloadevent'), + ] + + operations = [ + migrations.RunPython(grant_staff_access, revoke_staff_access), + ] diff --git a/osf/migrations/0047_downloadevent_status_code.py b/osf/migrations/0047_downloadevent_status_code.py new file mode 100644 index 00000000000..ff99a0ea8de --- /dev/null +++ b/osf/migrations/0047_downloadevent_status_code.py @@ -0,0 +1,18 @@ +# Generated by Django 4.2.26 on 2026-07-27 12:38 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('osf', '0046_dashboard_group_staff_access'), + ] + + operations = [ + migrations.AddField( + model_name='downloadevent', + name='status_code', + field=models.PositiveSmallIntegerField(blank=True, null=True), + ), + ] diff --git a/osf/migrations/0048_downloadevent_storage_provider.py b/osf/migrations/0048_downloadevent_storage_provider.py new file mode 100644 index 00000000000..b49f0648ca1 --- /dev/null +++ b/osf/migrations/0048_downloadevent_storage_provider.py @@ -0,0 +1,18 @@ +# Generated by Django 4.2.26 on 2026-07-28 17:05 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('osf', '0047_downloadevent_status_code'), + ] + + operations = [ + migrations.AddField( + model_name='downloadevent', + name='storage_provider', + field=models.CharField(blank=True, default='', max_length=32), + ), + ] diff --git a/osf/models/__init__.py b/osf/models/__init__.py index 918ca9aa009..410f62e6872 100644 --- a/osf/models/__init__.py +++ b/osf/models/__init__.py @@ -12,6 +12,7 @@ from .admin_log_entry import AdminLogEntry from .admin_profile import AdminProfile from .analytics import UserActivityCounter, PageCounter +from .download_event import DownloadEvent from .archive import ArchiveJob, ArchiveTarget from .banner import ScheduledBanner from .base import ( diff --git a/osf/models/download_event.py b/osf/models/download_event.py new file mode 100644 index 00000000000..6484dd47b76 --- /dev/null +++ b/osf/models/download_event.py @@ -0,0 +1,73 @@ +from django.db import models + + +class DownloadEvent(models.Model): + """One metadata row per download — never the file contents. + + Foundation for the download-telemetry capture and dashboard. Append-only: + rows are written from the download flow (single files at the osf.io redirect + view, folder/project zips from the WaterButler callback) and read, always + scoped to a time range, by the dashboard. + """ + + FILE = 'file' + FOLDER_ZIP = 'folder_zip' + PROJECT = 'project' + DOWNLOAD_TYPES = ( + (FILE, 'Single file'), + (FOLDER_ZIP, 'Folder zip'), + (PROJECT, 'Whole-project zip'), + ) + + created = models.DateTimeField(auto_now_add=True, db_index=True) + + # what was downloaded + resource_guid = models.CharField(max_length=255, blank=True, default='', db_index=True) + path = models.TextField(blank=True, default='') + download_type = models.CharField(max_length=16, choices=DOWNLOAD_TYPES) + # null for single files (only zips stream through WB, which reports completion) + zip_completed = models.BooleanField(null=True, blank=True) + # The HTTP status WaterButler ended the download on. Together with zip_completed it tells + # the three zip outcomes apart: completed (zip_completed=True), cancelled mid-stream + # (zip_completed=False, status 200/302 — headers already sent), and failed through no + # fault of the user (zip_completed=False, status >= 400). Null for single files and for + # callbacks from a WaterButler build that predates this field. + status_code = models.PositiveSmallIntegerField(null=True, blank=True) + size_bytes = models.BigIntegerField(null=True, blank=True) + + # which storage the bytes came from -- 'osfstorage' or a connected addon + # ('github', 'dropbox', 's3', 'box', ...). Blank for rows recorded before this field. + storage_provider = models.CharField(max_length=32, blank=True, default='') + # storage_region = where the bytes were served from (capacity); + # user_region = roughly where the user is. Kept separate on purpose. + storage_region = models.CharField(max_length=64, blank=True, default='') + user_region = models.CharField(max_length=64, blank=True, default='') + ip = models.GenericIPAddressField(null=True, blank=True) + source_area = models.CharField(max_length=128, blank=True, default='') + + # nullable: anonymous downloads of public files + user = models.ForeignKey( + 'osf.OSFUser', + null=True, + blank=True, + on_delete=models.SET_NULL, + related_name='download_events', + ) + + class Meta: + # `created` is indexed on the field; these cover the dashboard's + # time-range group-bys. + indexes = [ + models.Index(fields=['created', 'download_type'], name='download_event_crt_type'), + models.Index(fields=['created', 'storage_region'], name='download_event_crt_regn'), + models.Index(fields=['created', 'user_region'], name='download_event_crt_user'), + ] + + def __repr__(self): + return ( + f'' + ) + + def __str__(self): + return self.__repr__() diff --git a/osf/utils/download_telemetry.py b/osf/utils/download_telemetry.py new file mode 100644 index 00000000000..66dc5a3e9a0 --- /dev/null +++ b/osf/utils/download_telemetry.py @@ -0,0 +1,170 @@ +import functools +import logging + +from addons.osfstorage.settings import DEFAULT_REGION_NAME +from framework.celery_tasks import app +from framework.postcommit_tasks.handlers import enqueue_postcommit_task + +logger = logging.getLogger(__name__) + +# Set on every user until they pick something in their profile, so it says nothing +# about where they actually are. +UNSET_USER_TIMEZONE = 'Etc/UTC' + +# Identifiers worth having in the log line to track a failure back to one download. +# Deliberately excludes the IP. +LOGGED_CONTEXT_KEYS = ('download_type', 'resource_guid', 'file_id', 'user_guid') + + +def never_breaks_downloads(fn): + """Swallow and log anything this raises. + + Wraps the whole capture, not just the write — gathering the values is as capable of + raising as storing them is, and neither is a reason for a download to fail. + """ + @functools.wraps(fn) + def wrapped(*args, **kwargs): + try: + return fn(*args, **kwargs) + except Exception as exc: + # exc_info carries the traceback; the rest names the failure and which + # download it was, so a report is actionable without reproducing it. + logger.exception( + 'Failed to record a download event in %s: %s: %s [%s]', + fn.__name__, + type(exc).__name__, + exc, + ', '.join( + f'{key}={kwargs[key]!r}' + for key in LOGGED_CONTEXT_KEYS + if kwargs.get(key) + ) or 'no context', + ) + return wrapped + + +@never_breaks_downloads +def record_download(**kwargs): + """Enqueue a :class:`DownloadEvent` write.""" + enqueue_postcommit_task(write_download_event, (), kwargs, celery=True) + + +@app.task(max_retries=5, default_retry_delay=60) +def write_download_event( + download_type, + resource_guid='', + path='', + file_id=None, + version_identifier=None, + size_bytes=None, + storage_provider='', + storage_region_id=None, + zip_completed=None, + status_code=None, + user_guid=None, + ip=None, + source_area='', + tz='', +): + """Resolve the expensive bits and write one row. + + Callers hand over identifiers rather than loaded objects so that the download request + itself does no extra queries — everything that needs a lookup is resolved here. + """ + from osf.models import BaseFileNode, DownloadEvent, OSFUser + + user = OSFUser.load(user_guid) if user_guid else None + file_node = BaseFileNode.load(file_id) if file_id else None + file_version = _load_file_version(file_node, version_identifier) + + if file_version is not None: + if size_bytes is None: + size_bytes = file_version.size + if storage_region_id is None: + storage_region_id = file_version.region_id + + storage_region = _region_name(storage_region_id) or _resource_region_name(resource_guid) + + if not path and file_node is not None: + path = getattr(file_node, 'materialized_path', '') or '' + + DownloadEvent.objects.create( + download_type=download_type, + resource_guid=_truncate(resource_guid, 255), + path=path or '', + size_bytes=size_bytes if size_bytes is not None and size_bytes >= 0 else None, + zip_completed=zip_completed, + status_code=status_code, + storage_provider=_truncate(storage_provider, 32), + storage_region=_truncate(storage_region, 64), + user_region=_truncate(derive_user_region(tz, user, storage_region), 64), + ip=ip or None, + source_area=_truncate(source_area, 128), + user=user, + ) + + +def derive_user_region(tz, user, storage_region): + """Best available guess at where the user is, most to least trustworthy. + + The live browser timezone is the only real signal; the rest are fallbacks so the + dashboard isn't mostly blank. An empty string means we genuinely don't know, which + is more useful than a wrong guess. + """ + if tz: + return tz + + profile_timezone = getattr(user, 'timezone', '') + if profile_timezone and profile_timezone != UNSET_USER_TIMEZONE: + return profile_timezone + + # Everything defaults to the US region, so it only tells us something when it's been + # deliberately changed. + if storage_region and storage_region != DEFAULT_REGION_NAME: + return storage_region + + return '' + + +def _load_file_version(file_node, version_identifier): + """The version that was served, for its size and region.""" + if file_node is None: + return None + + from osf.models import FileVersion + + versions = FileVersion.objects.filter(basefilenode=file_node) + if version_identifier: + return versions.filter(identifier=version_identifier).first() + return versions.order_by('-created').first() + + +def _region_name(region_id): + if not region_id: + return '' + + from addons.osfstorage.models import Region + + region = Region.objects.filter(id=region_id).first() + return region.name if region else '' + + +def _resource_region_name(resource_guid): + """Where a zip was served from — zips have no single file version to read it off.""" + if not resource_guid: + return '' + + from osf.models import Guid + + resource, _ = Guid.load_referent(resource_guid) + region = getattr(resource, 'osfstorage_region', None) + return getattr(region, 'name', '') or '' + + +def _truncate(value, max_length): + """Keep user-controllable values inside their column. + + ``source`` and ``tz`` arrive off the query string, so they're whatever the caller + put there. + """ + return (value or '')[:max_length] diff --git a/poetry.lock b/poetry.lock index 6a2bc3e2d44..8281b60fb6c 100644 --- a/poetry.lock +++ b/poetry.lock @@ -986,6 +986,18 @@ tzdata = {version = "*", markers = "sys_platform == \"win32\""} argon2 = ["argon2-cffi (>=19.1.0)"] bcrypt = ["bcrypt"] +[[package]] +name = "django-admin-rangefilter" +version = "0.13.5" +description = "django-admin-rangefilter app, add the filter by a custom date range on the admin UI." +optional = false +python-versions = "!=3.0.*,!=3.1.*,!=3.2.*,!=3.3.*,!=3.4.*,>=2.7" +groups = ["main"] +files = [ + {file = "django_admin_rangefilter-0.13.5-py2.py3-none-any.whl", hash = "sha256:7fdcd1ee9007e2b76a5bd35bc017e785bdc7b17a24b8a36aef4ce80741356c47"}, + {file = "django_admin_rangefilter-0.13.5.tar.gz", hash = "sha256:3134e9e877f59ccad5949cd25cd2db9bf20d170fd8070be584875fdca96d3a14"}, +] + [[package]] name = "django-bulk-update" version = "2.2.0" @@ -4757,4 +4769,4 @@ testing = ["coverage (>=5.0.3)", "zope.event", "zope.testing"] [metadata] lock-version = "2.1" python-versions = "^3.12" -content-hash = "5599dfc677ced71d0e9097822fb1d1f6f0e22320eacca87299d1542b89e9f286" +content-hash = "4f88d33107a210745397689a6e82afaa2bda27c74f61d0681bbda7785de928f5" diff --git a/pyproject.toml b/pyproject.toml index 70309e505fa..2f140ec3e9f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -82,6 +82,7 @@ django-guardian = "2.4.0" # Admin requirements django-webpack-loader = {git = "https://github.com/CenterForOpenScience/django-webpack-loader.git", rev = "6b62fef7d6bc9d25d7b7b7a303f4580ad24831a6"} # branch is feature/v1-webpack-stats django-sendgrid-v5 = "1.2.3" # metadata says python 3.10 not supported, but tests pass +django-admin-rangefilter = "0.13.5" # OSF models django-typed-models = "0.14.0" diff --git a/tests/test_download_events_dashboard.py b/tests/test_download_events_dashboard.py new file mode 100644 index 00000000000..5cdfefabcad --- /dev/null +++ b/tests/test_download_events_dashboard.py @@ -0,0 +1,369 @@ +from datetime import timedelta + +from django.apps import apps as global_apps +from django.contrib.admin.sites import AdminSite +from django.contrib.auth.models import Group, Permission +from django.utils import timezone + +from osf.admin import DASHBOARD_GROUP_NAME, DownloadEventsView +from osf.models import DownloadEvent +from osf_tests.factories import AuthUserFactory, ProjectFactory +from tests.base import OsfTestCase + + +class FakeRequest: + def __init__(self, user): + self.user = user + self.GET = {} + + +def make_event(**kwargs): + defaults = { + 'download_type': DownloadEvent.FILE, + 'size_bytes': 1024 ** 3, + 'resource_guid': '', + } + defaults.update(kwargs) + return DownloadEvent.objects.create(**defaults) + + +class TestDashboardAccess(OsfTestCase): + """Membership in the allow-list group is the only key. + + Not a permission check — ModelBackend answers True to every permission for a + superuser, which would defeat the point of the dashboard. + """ + + def setUp(self): + super().setUp() + self.admin = DownloadEventsView(DownloadEvent, AdminSite()) + self.group, _ = Group.objects.get_or_create(name=DASHBOARD_GROUP_NAME) + + def _user(self, in_group=False, superuser=False, with_perm=False): + user = AuthUserFactory() + user.is_staff = True + user.is_superuser = superuser + user.save() + if in_group: + self.group.user_set.add(user) + if with_perm: + user.user_permissions.add(Permission.objects.get(codename='view_downloadevent')) + return type(user).objects.get(pk=user.pk) + + def test_group_member_gets_in(self): + request = FakeRequest(self._user(in_group=True)) + + assert self.admin.has_view_permission(request) is True + assert self.admin.has_module_permission(request) is True + + def test_superuser_outside_the_group_is_locked_out(self): + """The whole point: not even admins, unless they're on the list.""" + request = FakeRequest(self._user(superuser=True)) + + assert self.admin.has_view_permission(request) is False + assert self.admin.has_module_permission(request) is False + + def test_django_view_permission_alone_is_not_enough(self): + request = FakeRequest(self._user(with_perm=True)) + + assert self.admin.has_view_permission(request) is False + + def test_ordinary_staff_is_locked_out(self): + request = FakeRequest(self._user()) + + assert self.admin.has_view_permission(request) is False + + def test_it_is_read_only_even_for_group_members(self): + """Append-only telemetry — nothing is editable through the admin.""" + request = FakeRequest(self._user(in_group=True)) + + assert self.admin.has_add_permission(request) is False + assert self.admin.has_change_permission(request) is False + assert self.admin.has_delete_permission(request) is False + + def test_it_is_hidden_from_the_admin_index(self): + """`has_module_permission` is what keeps it off the list of pages.""" + request = FakeRequest(self._user(superuser=True)) + + assert self.admin.has_module_permission(request) is False + + +class TestDashboardData(OsfTestCase): + + def setUp(self): + super().setUp() + self.admin = DownloadEventsView(DownloadEvent, AdminSite()) + + def test_empty_range_does_not_blow_up(self): + """The default window is the last hour, so empty is the normal case.""" + data = self.admin.get_dashboard_data(DownloadEvent.objects.none()) + + assert data['summary']['total_downloads'] == 0 + assert data['summary']['total_gb'] == 0 + assert data['split']['file']['count_percent'] == 0 + assert data['split']['zip']['gb_percent'] == 0 + assert data['time_series']['labels'] == [] + assert data['top_projects'] == [] + + def test_events_with_unknown_size_do_not_blow_up(self): + """`size_bytes` is null when we could not determine it.""" + make_event(size_bytes=None, storage_region='Germany') + make_event(size_bytes=None, download_type=DownloadEvent.PROJECT) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['summary']['total_downloads'] == 2 + assert data['summary']['total_gb'] == 0 + assert data['storage_regions'][0]['gb'] == 0 + + def test_all_zero_sizes_do_not_blow_up(self): + make_event(size_bytes=0, storage_region='Germany') + make_event(size_bytes=0, storage_region='Germany') + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['storage_regions'][0]['gb_percent'] == 0 + + def test_totals_and_split(self): + make_event(size_bytes=2 * 1024 ** 3) + make_event(size_bytes=2 * 1024 ** 3, download_type=DownloadEvent.PROJECT) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['summary']['total_downloads'] == 2 + assert data['summary']['total_gb'] == 4 + assert data['split']['file']['count_percent'] == 50 + assert data['split']['zip']['count_percent'] == 50 + + def test_zip_outcomes_split_completed_cancelled_failed(self): + # a completed zip, a user cancel (False at 200), and a server failure (False at 5xx) + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=True, status_code=200) + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=False, status_code=200) + make_event(download_type=DownloadEvent.PROJECT, zip_completed=False, status_code=502) + # a single file — has no outcome, must not land in any bucket + make_event(download_type=DownloadEvent.FILE) + + outcomes = self.admin.get_dashboard_data(DownloadEvent.objects.all())['zip_outcomes'] + + assert outcomes == {'completed': 1, 'cancelled': 1, 'failed': 1} + + def test_failed_zips_summary_counts_only_server_errors(self): + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=False, status_code=500) + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=False, status_code=503) + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=False, status_code=200) + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=True, status_code=200) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['summary']['failed_zips'] == 2 + + def test_incomplete_zip_without_a_status_counts_as_cancelled(self): + """Callbacks from a WaterButler build predating status_code have status None — treat + those as cancels, not failures, so we never over-report failures.""" + make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=False, status_code=None) + + outcomes = self.admin.get_dashboard_data(DownloadEvent.objects.all())['zip_outcomes'] + + assert outcomes == {'completed': 0, 'cancelled': 1, 'failed': 0} + + def test_outcome_column_labels(self): + completed = make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=True) + cancelled = make_event(download_type=DownloadEvent.FOLDER_ZIP, zip_completed=False, status_code=200) + failed = make_event(download_type=DownloadEvent.PROJECT, zip_completed=False, status_code=500) + single = make_event(download_type=DownloadEvent.FILE) + + assert self.admin.outcome(completed) == 'Completed' + assert self.admin.outcome(cancelled) == 'Cancelled' + assert self.admin.outcome(failed) == 'Failed' + assert self.admin.outcome(single) == '—' + + def test_storage_provider_breakdown(self): + make_event(storage_provider='osfstorage', size_bytes=3 * 1024 ** 3) + make_event(storage_provider='osfstorage', size_bytes=1 * 1024 ** 3) + make_event(storage_provider='github', size_bytes=2 * 1024 ** 3) + + providers = self.admin.get_dashboard_data(DownloadEvent.objects.all())['storage_providers'] + by_name = {row['name']: row for row in providers} + + assert by_name['osfstorage']['downloads'] == 2 + assert by_name['osfstorage']['gb'] == 4 + assert by_name['github']['downloads'] == 1 + assert by_name['github']['gb'] == 2 + + def test_blank_storage_provider_folds_into_unknown(self): + make_event(storage_provider='') + + providers = self.admin.get_dashboard_data(DownloadEvent.objects.all())['storage_providers'] + + assert [row['name'] for row in providers] == ['Unknown'] + + def test_region_breakdown_splits_requests_by_type(self): + # Germany: 2 files + 1 folder zip + 1 project zip = 4 total, 3 zips + make_event(storage_region='Germany', download_type=DownloadEvent.FILE) + make_event(storage_region='Germany', download_type=DownloadEvent.FILE) + make_event(storage_region='Germany', download_type=DownloadEvent.FOLDER_ZIP) + make_event(storage_region='Germany', download_type=DownloadEvent.PROJECT) + + regions = self.admin.get_dashboard_data(DownloadEvent.objects.all())['storage_regions'] + germany = next(r for r in regions if r['name'] == 'Germany') + + assert germany['file_count'] == 2 + assert germany['zip_count'] == 2 + # file + zip always equals the total — the invariant QA will check + assert germany['file_count'] + germany['zip_count'] == germany['downloads'] == 4 + + def test_region_type_counts_match_the_period_totals(self): + """Summed across regions, file/zip counts equal the header's split totals — so the + subtitle and the bars can never disagree.""" + make_event(storage_region='Germany', download_type=DownloadEvent.FILE) + make_event(storage_region='United States', download_type=DownloadEvent.FILE) + make_event(storage_region='Germany', download_type=DownloadEvent.FOLDER_ZIP) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + regions = data['storage_regions'] + + assert sum(r['file_count'] for r in regions) == data['split']['file']['count'] == 2 + assert sum(r['zip_count'] for r in regions) == data['split']['zip']['count'] == 1 + + def test_region_type_counts_present_even_when_empty(self): + data = self.admin.get_dashboard_data(DownloadEvent.objects.none()) + + assert data['storage_regions'] == [] + assert data['user_regions'] == [] + + def test_blank_and_null_regions_fold_into_unknown(self): + make_event(storage_region='') + make_event(storage_region=' ') + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert [row['name'] for row in data['storage_regions']] == ['Unknown'] + assert data['storage_regions'][0]['downloads'] == 2 + + def test_top_projects_shows_title_and_guid(self): + user = AuthUserFactory() + node = ProjectFactory(creator=user, title='Panic Download Project') + make_event(resource_guid=node._id, size_bytes=5 * 1024 ** 3) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['top_projects'][0]['name'] == f'Panic Download Project ({node._id})' + assert data['top_projects'][0]['gb'] == 5 + + def test_top_projects_falls_back_to_the_bare_guid(self): + """An unresolvable guid still has to say something.""" + make_event(resource_guid='notaguid', size_bytes=1024 ** 3) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['top_projects'][0]['name'] == 'notaguid' + + def test_time_series_buckets_by_type(self): + make_event(size_bytes=1024 ** 3) + make_event(size_bytes=3 * 1024 ** 3, download_type=DownloadEvent.FOLDER_ZIP) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert len(data['time_series']['labels']) >= 1 + assert sum(data['time_series']['file']) == 1 + assert sum(data['time_series']['zip']) == 3 + + def test_time_series_spans_gaps(self): + """Quiet buckets render as zero instead of collapsing the axis.""" + recent = make_event(size_bytes=1024 ** 3) + old = make_event(size_bytes=1024 ** 3) + DownloadEvent.objects.filter(pk=old.pk).update( + created=timezone.now() - timedelta(days=4) + ) + DownloadEvent.objects.filter(pk=recent.pk).update(created=timezone.now()) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert len(data['time_series']['labels']) == 5 + assert data['time_series']['file'][0] == 1 + assert data['time_series']['file'][-1] == 1 + assert data['time_series']['file'][2] == 0 + + def test_every_bucket_size_places_all_the_bytes(self): + """The bucket key has to land on the same instant whether it came from an + event or from walking the axis, at every granularity.""" + spans = { + '15m': [timedelta(minutes=m) for m in (0, 20, 50, 80)], + '1h': [timedelta(hours=h) for h in (0, 3, 9, 20)], + '1d': [timedelta(days=d) for d in (0, 2, 5, 10)], + '1w': [timedelta(days=d) for d in (0, 10, 20, 30)], + } + now = timezone.now() + for bucket_size, offsets in spans.items(): + DownloadEvent.objects.all().delete() + for offset in offsets: + event = make_event(size_bytes=1024 ** 3) + DownloadEvent.objects.filter(pk=event.pk).update(created=now - offset) + + series = self.admin._build_time_series(DownloadEvent.objects.all()) + + assert sum(series['file']) == float(len(offsets)), ( + f'{bucket_size} buckets dropped data' + ) + + def test_time_series_with_a_single_event(self): + """start == end, so the range delta is zero.""" + make_event(size_bytes=2 * 1024 ** 3) + + series = self.admin._build_time_series(DownloadEvent.objects.all()) + + assert sum(series['file']) == 2.0 + + def test_unique_users_ignores_anonymous(self): + user = AuthUserFactory() + make_event(user=user) + make_event(user=user) + make_event(user=None) + + data = self.admin.get_dashboard_data(DownloadEvent.objects.all()) + + assert data['summary']['unique_users'] == 1 + + +class TestStaffAccessMigration(OsfTestCase): + """Django's admin rejects anyone without `is_staff` before our gate runs, so + the allow-listed users need it to reach the page at all.""" + + def setUp(self): + super().setUp() + from importlib import import_module + self.migration = import_module('osf.migrations.0046_dashboard_group_staff_access') + self.group, _ = Group.objects.get_or_create(name=DASHBOARD_GROUP_NAME) + + def test_it_grants_staff_to_group_members_only(self): + member = AuthUserFactory() + outsider = AuthUserFactory() + self.group.user_set.add(member) + + self.migration.grant_staff_access(global_apps, None) + + member.refresh_from_db() + outsider.refresh_from_db() + assert member.is_staff is True + assert outsider.is_staff is False + + def test_it_does_not_grant_superuser(self): + """is_staff opens the admin door; is_superuser would bypass every gate.""" + member = AuthUserFactory() + self.group.user_set.add(member) + + self.migration.grant_staff_access(global_apps, None) + + member.refresh_from_db() + assert member.is_superuser is False + + def test_reversing_does_not_strip_staff_from_existing_admins(self): + admin_user = AuthUserFactory() + admin_user.is_staff = True + admin_user.save() + self.group.user_set.add(admin_user) + + self.migration.revoke_staff_access(global_apps, None) + + admin_user.refresh_from_db() + assert admin_user.is_staff is True diff --git a/tests/test_download_telemetry.py b/tests/test_download_telemetry.py new file mode 100644 index 00000000000..ac71b6afe41 --- /dev/null +++ b/tests/test_download_telemetry.py @@ -0,0 +1,364 @@ +import importlib +import logging +import time + +import pytest +from django.apps import apps as django_apps +from django.contrib.auth.models import Group + +from addons.osfstorage.settings import DEFAULT_REGION_NAME +from api_tests.utils import create_test_file +from framework.auth import signing +from osf.models import DownloadEvent, OSFUser +from osf.utils.download_telemetry import derive_user_region, record_download +from osf_tests.factories import AuthUserFactory, ProjectFactory +from tests.base import OsfTestCase + +download_event_migration = importlib.import_module('osf.migrations.0045_downloadevent') + + +@pytest.mark.django_db +class TestDownloadEventModel: + """The table and indexes the dashboard's time-range group-bys rely on.""" + + def test_expected_fields(self): + field_names = {field.name for field in DownloadEvent._meta.get_fields()} + + assert field_names >= { + 'created', + 'resource_guid', + 'path', + 'download_type', + 'zip_completed', + 'size_bytes', + 'storage_region', + 'user_region', + 'ip', + 'source_area', + 'user', + } + + def test_expected_indexes(self): + index_names = {index.name for index in DownloadEvent._meta.indexes} + + assert index_names == { + 'download_event_crt_type', + 'download_event_crt_regn', + 'download_event_crt_user', + } + + def test_deleting_the_user_keeps_the_row_and_nulls_the_user(self): + user = AuthUserFactory() + download_event = DownloadEvent.objects.create(download_type=DownloadEvent.FILE, user=user) + OSFUser.objects.filter(id=user.id).delete() + download_event.refresh_from_db() + + assert download_event.user_id is None + + +@pytest.mark.django_db +class TestDownloadEventDashboardGroupMigration: + + def test_create_dashboard_group_adds_dashboard_users(self): + seed_user = AuthUserFactory(username=download_event_migration.DASHBOARD_USERS[0]) + download_event_migration.create_dashboard_group(django_apps, None) + group = Group.objects.get(name=download_event_migration.DASHBOARD_GROUP_NAME) + + assert seed_user in group.user_set.all() + + def test_create_dashboard_group_skips_not_dashboard_users(self): + download_event_migration.create_dashboard_group(django_apps, None) + group = Group.objects.get(name=download_event_migration.DASHBOARD_GROUP_NAME) + + assert group.user_set.count() == 0 + + +class TestZipDownloadTelemetry(OsfTestCase): + """Folder and project zips, recorded from the WaterButler callback. + + Zips are requested straight from WaterButler, so this callback is the only place we + hear about them. + """ + + def setUp(self): + super().setUp() + self.user = AuthUserFactory() + self.node = ProjectFactory(creator=self.user) + self.url = self.node.api_url_for('create_waterbutler_log') + + def build_payload(self, materialized='/', action='download_zip', **action_meta): + meta = dict( + bytes_downloaded=2048, + completed=True, + ip='198.51.100.7', + source='files', + tz='Europe/Kyiv', + ) + meta.update(action_meta) + options = { + 'auth': {'id': self.user._id}, + 'action': action, + 'provider': 'osfstorage', + 'time': time.time() + 1000, + 'metadata': { + 'nid': self.node._id, + 'materialized': materialized, + 'path': materialized, + 'kind': 'folder', + 'provider': 'osfstorage', + }, + 'action_meta': meta, + } + message, signature = signing.default_signer.sign_payload(options) + return {'payload': message, 'signature': signature} + + def test_project_zip_is_recorded(self): + res = self.app.put(self.url, json=self.build_payload(materialized='/')) + + assert res.status_code == 200 + event = DownloadEvent.objects.get() + assert event.download_type == DownloadEvent.PROJECT + assert event.resource_guid == self.node._id + assert event.user == self.user + assert event.size_bytes == 2048 + assert event.zip_completed is True + assert event.ip == '198.51.100.7' + assert event.source_area == 'files' + assert event.user_region == 'Europe/Kyiv' + + def test_folder_zip_is_recorded_with_its_path(self): + self.app.put(self.url, json=self.build_payload(materialized='/data/raw/')) + + event = DownloadEvent.objects.get() + assert event.download_type == DownloadEvent.FOLDER_ZIP + assert event.path == '/data/raw/' + + def test_incomplete_zip_is_recorded_as_incomplete(self): + self.app.put(self.url, json=self.build_payload(completed=False)) + + assert DownloadEvent.objects.get().zip_completed is False + + def test_successful_zip_records_its_status_code(self): + self.app.put(self.url, json=self.build_payload(completed=True, status_code=200)) + + event = DownloadEvent.objects.get() + assert event.zip_completed is True + assert event.status_code == 200 + + def test_failed_zip_is_distinguishable_from_a_cancel(self): + """A server failure comes through as completed=False at a 5xx; a user cancel comes + through as completed=False at 200 (headers already sent). The status tells them apart.""" + self.app.put(self.url, json=self.build_payload(completed=False, status_code=500)) + failed = DownloadEvent.objects.get() + assert failed.zip_completed is False + assert failed.status_code == 500 + + DownloadEvent.objects.all().delete() + + self.app.put(self.url, json=self.build_payload(completed=False, status_code=200)) + cancelled = DownloadEvent.objects.get() + assert cancelled.zip_completed is False + assert cancelled.status_code == 200 + + def test_mfr_render_is_not_recorded(self): + self.app.put(self.url, json=self.build_payload(is_mfr_render=True)) + + assert not DownloadEvent.objects.exists() + + def test_single_file_action_is_not_recorded_here(self): + """Single files are caught at the redirect view — recording them here too would + double count every one of them.""" + self.app.put(self.url, json=self.build_payload(action='download_file')) + + assert not DownloadEvent.objects.exists() + + def test_oversized_source_is_truncated_to_the_column(self): + self.app.put(self.url, json=self.build_payload(source='f' * 500)) + + assert len(DownloadEvent.objects.get().source_area) == 128 + + def test_payload_missing_the_new_waterbutler_fields_records(self): + """Guards the rollout window: osf.io must keep reading zip callbacks from a + WaterButler build that predates bytes_downloaded/completed/ip in action_meta.""" + options = { + 'auth': {'id': self.user._id}, + 'action': 'download_zip', + 'provider': 'osfstorage', + 'time': time.time() + 1000, + 'metadata': { + 'nid': self.node._id, + 'materialized': '/', + 'path': '/', + 'kind': 'folder', + 'provider': 'osfstorage', + }, + 'action_meta': {}, + } + message, signature = signing.default_signer.sign_payload(options) + res = self.app.put(self.url, json={'payload': message, 'signature': signature}) + + assert res.status_code == 200 + event = DownloadEvent.objects.get() + assert event.size_bytes is None + assert event.zip_completed is None + assert event.status_code is None + assert event.ip is None + + def test_storage_region_comes_from_the_projects_node(self): + self.app.put(self.url, json=self.build_payload()) + event = DownloadEvent.objects.get() + + assert event.storage_region == self.node.osfstorage_region.name + + def test_storage_provider_comes_from_the_callback(self): + self.app.put(self.url, json=self.build_payload()) + + assert DownloadEvent.objects.get().storage_provider == 'osfstorage' + + def test_callback_still_succeeds_when_recording_fails(self, ): + with pytest.MonkeyPatch.context() as patch: + patch.setattr( + 'addons.base.views.record_download', + lambda **kwargs: (_ for _ in ()).throw(ValueError('boom')), + ) + res = self.app.put(self.url, json=self.build_payload()) + + assert res.status_code == 200 + + +class TestSingleFileDownloadTelemetry(OsfTestCase): + """Single files, recorded at the redirect view before we 302 on to WaterButler.""" + + def setUp(self): + super().setUp() + self.user = AuthUserFactory() + self.node = ProjectFactory(creator=self.user) + self.file = create_test_file(self.node, self.user, size=4096) + self.guid = self.file.get_guid()._id + + def test_download_is_recorded_with_link_tags(self): + res = self.app.get( + f'/download/{self.guid}/?source=file-detail&tz=Europe%2FKyiv', + auth=self.user.auth, + ) + + assert res.status_code == 302 + event = DownloadEvent.objects.get() + assert event.download_type == DownloadEvent.FILE + assert event.resource_guid == self.node._id + assert event.user == self.user + assert event.source_area == 'file-detail' + assert event.user_region == 'Europe/Kyiv' + + def test_size_and_region_come_from_the_file_version(self): + self.app.get(f'/download/{self.guid}/', auth=self.user.auth) + + event = DownloadEvent.objects.get() + assert event.size_bytes == 4096 + assert event.storage_region == self.node.osfstorage_region.name + + def test_storage_provider_comes_from_the_file(self): + self.app.get(f'/download/{self.guid}/', auth=self.user.auth) + + assert DownloadEvent.objects.get().storage_provider == 'osfstorage' + + def test_zip_completed_is_unset_for_single_files(self): + """Only zips stream through WaterButler, so nothing reports completion here.""" + self.app.get(f'/download/{self.guid}/', auth=self.user.auth) + + assert DownloadEvent.objects.get().zip_completed is None + + def test_anonymous_download_is_recorded_without_a_user(self): + self.node.is_public = True + self.node.save() + + self.app.get(f'/download/{self.guid}/') + + event = DownloadEvent.objects.get() + assert event.user is None + + def test_mfr_render_is_not_recorded(self): + self.app.get(f'/download/{self.guid}/?mode=render', auth=self.user.auth) + + assert not DownloadEvent.objects.exists() + + def test_download_still_succeeds_when_recording_fails(self): + with pytest.MonkeyPatch.context() as patch: + patch.setattr( + 'addons.base.views.record_download', + lambda **kwargs: (_ for _ in ()).throw(ValueError('boom')), + ) + res = self.app.get(f'/download/{self.guid}/', auth=self.user.auth) + + assert res.status_code == 302 + + +class TestUserRegionDerivation: + """The fallback chain, most to least trustworthy.""" + + class FakeUser: + def __init__(self, timezone): + self.timezone = timezone + + def test_live_browser_timezone_wins(self): + user = self.FakeUser('America/New_York') + assert derive_user_region('Europe/Kyiv', user, 'Germany') == 'Europe/Kyiv' + + def test_falls_back_to_profile_timezone(self): + user = self.FakeUser('America/New_York') + assert derive_user_region('', user, 'Germany') == 'America/New_York' + + def test_default_profile_timezone_is_not_a_signal(self): + """Every user has Etc/UTC until they change it, so it says nothing.""" + user = self.FakeUser('Etc/UTC') + assert derive_user_region('', user, 'Germany') == 'Germany' + + def test_falls_back_to_storage_region(self): + assert derive_user_region('', None, 'Germany') == 'Germany' + + def test_default_storage_region_is_not_a_signal(self): + assert derive_user_region('', None, DEFAULT_REGION_NAME) == '' + + def test_unknown_is_empty(self): + assert derive_user_region('', None, '') == '' + + +@pytest.mark.django_db +class TestRecordDownloadNeverRaises: + + def test_enqueue_failure_is_swallowed(self, monkeypatch): + def explode(*args, **kwargs): + raise ValueError('boom') + + monkeypatch.setattr('osf.utils.download_telemetry.enqueue_postcommit_task', explode) + + record_download(download_type=DownloadEvent.FILE, resource_guid='abcde') + + assert not DownloadEvent.objects.exists() + + def test_failure_is_logged_with_the_cause_and_the_download(self, monkeypatch, caplog): + """A report has to be actionable without reproducing it.""" + def explode(*args, **kwargs): + raise ValueError('boom') + + monkeypatch.setattr('osf.utils.download_telemetry.enqueue_postcommit_task', explode) + + with caplog.at_level(logging.ERROR, logger='osf.utils.download_telemetry'): + record_download( + download_type=DownloadEvent.FILE, + resource_guid='abcde', + user_guid='zyxwv', + ip='198.51.100.7', + ) + + record = caplog.records[0] + message = record.getMessage() + assert 'ValueError' in message + assert 'boom' in message + assert 'record_download' in message + assert "resource_guid='abcde'" in message + assert "user_guid='zyxwv'" in message + # the IP is not debugging information + assert '198.51.100.7' not in message + # the traceback is still attached + assert record.exc_info is not None diff --git a/website/settings/defaults.py b/website/settings/defaults.py index 4b88b16bf1f..9800831b297 100644 --- a/website/settings/defaults.py +++ b/website/settings/defaults.py @@ -603,6 +603,7 @@ class CeleryConfig: 'osf.management.commands.sync_doi_metadata', 'api.providers.tasks', 'api.users.tasks', + 'osf.utils.download_telemetry', 'osf.management.commands.daily_reporters_go', 'osf.management.commands.monthly_reporters_go', 'osf.external.spam.tasks',