diff --git a/beagle/__init__.py b/beagle/__init__.py index b19ee4b7..55e47090 100644 --- a/beagle/__init__.py +++ b/beagle/__init__.py @@ -1 +1 @@ -__version__ = "2.2.1" +__version__ = "2.3.0" diff --git a/beagle/settings.py b/beagle/settings.py index b2233886..61a3be5f 100644 --- a/beagle/settings.py +++ b/beagle/settings.py @@ -471,8 +471,7 @@ "IMPACT341,IMPACT+ (341 genes plus custom content),IMPACT468,HemePACT_v4,HemePACT_v3,IMPACT505,IMPACT410", ).split(",") -PERMISSION_DENIED_CC = json.loads(os.environ.get("BEAGLE_PERMISSION_DENIED_CC", "{}")) -PERMISSION_DENIED_EMAILS = json.loads(os.environ.get("BEAGLE_PERMISSION_DENIED_EMAIL", "{}")) +PERMISSION_DENIED_EMAILS = os.environ.get("BEAGLE_PERMISSION_DENIED_EMAIL", "").split(",") JOB_HANGING_ALERT_EMAILS = os.environ.get("BEAGLE_JOB_HANGING_ALERT_EMAILS", "").split(",") # Tempo diff --git a/beagle_etl/admin.py b/beagle_etl/admin.py index 068b6f8c..da0ae168 100644 --- a/beagle_etl/admin.py +++ b/beagle_etl/admin.py @@ -30,7 +30,8 @@ class SMILEMessagesAdmin(AdminAdvancedFiltersMixin, ModelAdmin): list_filter = ("request_id", "topic", "status") advanced_filter_fields = ("request_id", "topic", "status") ordering = ("-created_date",) - list_display = ("created_date", "request_id", "topic", "status") + list_display = ("created_date", "request_id", "gene_panel", "topic", "status") + search_fields = ("request_id", "gene_panel") class RequestCallbackJobAdmin(ModelAdmin): diff --git a/beagle_etl/jobs/metadb_jobs.py b/beagle_etl/jobs/metadb_jobs.py index 17c8e0d6..570f37de 100644 --- a/beagle_etl/jobs/metadb_jobs.py +++ b/beagle_etl/jobs/metadb_jobs.py @@ -29,10 +29,18 @@ WESJobFailedEvent, VoyagerCantProcessRequestAllNormalsEvent, SMILEUpdateEvent, + ErrorImportingFilesEvent, ) from notifier.tasks import send_notification from notifier.helper import get_emails_to_notify -from beagle_etl.models import Operator, ETLConfiguration, SMILEMessage, RequestCallbackJob, RequestCallbackJobStatus +from beagle_etl.models import ( + Operator, + ETLConfiguration, + SMILEMessage, + RequestCallbackJob, + RequestCallbackJobStatus, + SmileMessageStatus, +) from file_system.serializers import UpdateFileSerializer from file_system.repository.file_repository import FileRepository from file_system.models import File @@ -202,6 +210,7 @@ def new_request(message_id): # Validate samples and fastqs log, status = data.validate_all_samples() + logger.info(f"Request validation log for SMILEMessage id:{message_id}: {log}") message.add_log(log) jgn_id = None @@ -214,6 +223,38 @@ def new_request(message_id): study, _ = Study.objects.get_or_create(study_id=StudyObject.generate_study_id(data.labHeadName)) valid_samples = {k for k, v in status.items() if v.status == "COMPLETED"} + + retry_samples = {k for k, v in status.items() if v.status == "RETRY"} + + if retry_samples: + sample_status = sorted([sample.to_dict() for sample in status.values()], key=lambda d: d["sample"]) + message.set_sample_status(sample_status) + message.add_log(f"Permission Denied error during import for igoRequestId:{message.request_id} id:{message_id}") + logger.error(f"Permission Denied error during import for igoRequestId:{message.request_id} id:{message_id}") + message.retry() + message.refresh_from_db() + if message.status in (SmileMessageStatus.RETRY,): + logger.info( + f"Retrying import for igoRequestId:{message.request_id} id:{message_id}, " + f"attempt {message.retry_count} of 2, next attempt scheduled at {message.scheduled}" + ) + elif message.status in (SmileMessageStatus.FAILED,): + logger.error( + f"Import failed for igoRequestId:{message.request_id} id:{message_id} " + f"after {message.retry_count - 1} retries due to permission errors" + ) + for email in settings.PERMISSION_DENIED_EMAILS: + e = ErrorImportingFilesEvent( + job_notifier=settings.BEAGLE_NOTIFIER_EMAIL_GROUP, + email_to=email, + subject=f"VOYAGER: Permission Denied error during import for igoRequestId:{message.request_id} id:{message_id}", + email_from=settings.BEAGLE_NOTIFIER_EMAIL_FROM, + request_id=message.request_id, + msg=f"Samples {', '.join(sorted(retry_samples))} failed to import because fastqs don't have correct permissions", + ) + send_notification.delay(e.to_dict()) + return + request_metadata = data.request_metadata() import_status = True diff --git a/beagle_etl/migrations/0047_add_scheduled_retry_count.py b/beagle_etl/migrations/0047_add_scheduled_retry_count.py new file mode 100644 index 00000000..27e64f97 --- /dev/null +++ b/beagle_etl/migrations/0047_add_scheduled_retry_count.py @@ -0,0 +1,37 @@ +# Generated by Django 6.0.4 on 2026-07-08 20:56 + +import django.utils.timezone +from django.db import migrations, models + + +def set_scheduled_to_created_date(apps, schema_editor): + SMILEMessage = apps.get_model("beagle_etl", "SMILEMessage") + messages = list(SMILEMessage.objects.all()) + for message in messages: + message.scheduled = message.created_date + SMILEMessage.objects.bulk_update(messages, ["scheduled"], batch_size=1000) + + +def noop(apps, schema_editor): + pass + + +class Migration(migrations.Migration): + + dependencies = [ + ("beagle_etl", "0046_delete_skipproject"), + ] + + operations = [ + migrations.AddField( + model_name="smilemessage", + name="retry_count", + field=models.IntegerField(default=0), + ), + migrations.AddField( + model_name="smilemessage", + name="scheduled", + field=models.DateTimeField(default=django.utils.timezone.now), + ), + migrations.RunPython(set_scheduled_to_created_date, noop), + ] diff --git a/beagle_etl/models.py b/beagle_etl/models.py index 595040c4..60320970 100644 --- a/beagle_etl/models.py +++ b/beagle_etl/models.py @@ -1,8 +1,10 @@ import uuid import logging from enum import IntEnum +from datetime import timedelta from django.db import models from django.db.models import JSONField +from django.utils import timezone from django.contrib.postgres.fields import ArrayField from notifier.tasks import notifier_start from notifier.models import Notifier, JobGroup, JobGroupNotifier @@ -45,6 +47,7 @@ class SmileMessageStatus(IntEnum): COMPLETED = 3 NOT_SUPPORTED = 4 FAILED = 5 + RETRY = 6 class SMILEMessage(BaseModel): @@ -61,14 +64,22 @@ class SMILEMessage(BaseModel): default=SmileMessageStatus.PENDING, db_index=True, ) + scheduled = models.DateTimeField(default=timezone.now, editable=True) + retry_count = models.IntegerField(default=0) def in_progress(self): self.status = SmileMessageStatus.IN_PROGRESS - job_group = JobGroup.objects.create() - self.job_group = job_group - job_group_notifier_id = notifier_start(job_group, self.request_id) - job_group_notifier = JobGroupNotifier.objects.get(id=job_group_notifier_id) if job_group_notifier_id else None - self.job_group_notifier = job_group_notifier + # Retries re-enter in_progress() with a job_group already set from the first + # attempt; skip creating another one so retries update the original ticket + # instead of opening a new one each pass. + if not self.job_group: + job_group = JobGroup.objects.create() + self.job_group = job_group + job_group_notifier_id = notifier_start(job_group, self.request_id) + job_group_notifier = ( + JobGroupNotifier.objects.get(id=job_group_notifier_id) if job_group_notifier_id else None + ) + self.job_group_notifier = job_group_notifier self.save(update_fields=["job_group", "job_group_notifier", "status"]) def complete(self, request_metadata=None): @@ -83,6 +94,17 @@ def failed(self, request_metadata=None): if self.job_group_notifier and request_metadata: self._generate_description(request_metadata) + def retry(self): + if self.retry_count >= 2: + # Fail after 2 retries + self.status = SmileMessageStatus.FAILED + else: + # Retry import after 24 hours + self.status = SmileMessageStatus.RETRY + self.scheduled = self.scheduled + timedelta(hours=24) + self.retry_count += 1 + self.save(update_fields=["scheduled", "status", "retry_count"]) + def not_supported(self): self.status = SmileMessageStatus.NOT_SUPPORTED self.save(update_fields=["status"]) diff --git a/beagle_etl/smile_message/objects/sample_object.py b/beagle_etl/smile_message/objects/sample_object.py index 8b8f9106..6b1dbc5e 100644 --- a/beagle_etl/smile_message/objects/sample_object.py +++ b/beagle_etl/smile_message/objects/sample_object.py @@ -507,9 +507,22 @@ def validate_with_file_checks(self, redelivery: bool = False, log: str = "") -> log = self._process_validation_results(validation_results, log, redelivery) except Exception as e: if isinstance(e, ETLExceptions): - sample_status = SampleStatus( - sample_id=self.primaryId, igocomplete=self.igoComplete, code=e.code, status="FAILED", message=str(e) - ) + if isinstance(e, FailedToCopyFilePermissionDeniedException): + sample_status = SampleStatus( + sample_id=self.primaryId, + igocomplete=self.igoComplete, + code=e.code, + status="RETRY", + message=str(e), + ) + else: + sample_status = SampleStatus( + sample_id=self.primaryId, + igocomplete=self.igoComplete, + code=e.code, + status="FAILED", + message=str(e), + ) else: sample_status = SampleStatus( sample_id=self.primaryId, igocomplete=self.igoComplete, code=None, status="FAILED", message=str(e) diff --git a/beagle_etl/tasks.py b/beagle_etl/tasks.py index 4e38d90f..3f9e46de 100644 --- a/beagle_etl/tasks.py +++ b/beagle_etl/tasks.py @@ -65,12 +65,13 @@ def get_pending_smile_messages(): # Get all pending messages with priority pending_messages = SMILEMessage.objects.filter( - status=SmileMessageStatus.PENDING, + status__in=(SmileMessageStatus.PENDING, SmileMessageStatus.RETRY), topic__in=[ settings.METADB_NATS_NEW_REQUEST, settings.METADB_NATS_SAMPLE_UPDATE, settings.METADB_NATS_REQUEST_UPDATE, ], + scheduled__lte=datetime.datetime.now(tz=pytz.UTC), ).annotate( topic_priority=Case( When(topic=settings.METADB_NATS_NEW_REQUEST, then=0), diff --git a/beagle_etl/tests/jobs/test_metadb.py b/beagle_etl/tests/jobs/test_metadb.py index 7959cec0..a24c045a 100644 --- a/beagle_etl/tests/jobs/test_metadb.py +++ b/beagle_etl/tests/jobs/test_metadb.py @@ -3,9 +3,10 @@ """ import os -from mock import patch import json +from mock import patch from deepdiff import DeepDiff +from datetime import timedelta from django.conf import settings from django.contrib.auth.models import User from django.test import TestCase, override_settings @@ -157,6 +158,46 @@ def test_new_request( msg.refresh_from_db() self.assertEqual(msg.status, SmileMessageStatus.COMPLETED) + @patch("os.access") + @patch("notifier.tasks.notifier_start") + @patch("notifier.tasks.send_notification.delay") + @patch("notifier.models.JobGroupNotifier.objects.get") + @patch("file_system.tasks.populate_job_group_notifier_metadata.delay") + @patch("os.path.exists") + @patch("beagle_etl.jobs.helper_jobs.calculate_checksum.delay") + @override_settings(IMPORT_FILE_GROUP="1a1b29cf-3bc2-4f6c-b376-d4c5d701166a") + def test_retry_permission_denied( + self, + calculate_checksum, + path_exists, + populate_job_group_notifier_metadata, + job_group_notifier_get, + send_notification, + notifier_start, + access, + ): + calculate_checksum.return_value = None + path_exists.return_value = True + populate_job_group_notifier_metadata.return_value = None + job_group_notifier_get.return_value = self.job_group_notifier + notifier_start.return_value = True + send_notification.return_value = True + access.return_value = False + msg = SMILEMessage.objects.create( + topic="new-request", request_id="08944_B", gene_panel="", message=self.new_request_str + ) + scheduled_before_retry = msg.scheduled + msg.in_progress() + new_request(str(msg.id)) + msg.refresh_from_db() + files = FileRepository.filter( + metadata={settings.REQUEST_ID_METADATA_KEY: "08944_B"}, file_group=self.file_group_id + ) + self.assertEqual(files.count(), 0) + self.assertEqual(msg.status, SmileMessageStatus.RETRY) + self.assertEqual(msg.retry_count, 1) + self.assertEqual(msg.scheduled, (scheduled_before_retry + timedelta(hours=24))) + @patch("notifier.models.JobGroupNotifier.objects.get") @patch("notifier.tasks.send_notification.delay") @patch("file_system.tasks.populate_job_group_notifier_metadata.delay") diff --git a/runner/operator/argos_operator/v2_3_0/argos_operator.py b/runner/operator/argos_operator/v2_3_0/argos_operator.py index 87d6be55..e2e162b1 100644 --- a/runner/operator/argos_operator/v2_3_0/argos_operator.py +++ b/runner/operator/argos_operator/v2_3_0/argos_operator.py @@ -250,7 +250,7 @@ def get_files_for_pairs(self, pairing): tumor_current = tumor.first() bait_set = tumor_current.metadata["baitSet"] preservation_types = tumor_current.metadata["preservation"] - sample_origin = tumor_current.metadata["sampleOrigin"] + sample_origin = tumor_current.metadata.get("sampleOrigin") or "" pooled_normal_files, bait_set_reformatted, sample_name = get_pooled_normal_files( run_ids, preservation_types, bait_set, sample_origin ) diff --git a/runner/operator/argos_operator/v2_3_0/bin/make_sample.py b/runner/operator/argos_operator/v2_3_0/bin/make_sample.py index 99979668..50a7f561 100644 --- a/runner/operator/argos_operator/v2_3_0/bin/make_sample.py +++ b/runner/operator/argos_operator/v2_3_0/bin/make_sample.py @@ -144,7 +144,7 @@ def build_sample(data, ignore_sample_formatting=False): bait_set = meta["baitSet"] tumor_type = meta["tumorOrNormal"] specimen_type = meta[settings.SAMPLE_CLASS_METADATA_KEY] - sample_origin = meta["sampleOrigin"] + sample_origin = meta.get("sampleOrigin") or "" species = meta["species"] cmo_sample_name = format_sample_name( meta[settings.CMO_SAMPLE_NAME_METADATA_KEY], specimen_type, ignore_sample_formatting diff --git a/runner/operator/argos_operator/v2_3_0/bin/retrieve_samples_by_query.py b/runner/operator/argos_operator/v2_3_0/bin/retrieve_samples_by_query.py index e073542b..c6a7563e 100644 --- a/runner/operator/argos_operator/v2_3_0/bin/retrieve_samples_by_query.py +++ b/runner/operator/argos_operator/v2_3_0/bin/retrieve_samples_by_query.py @@ -92,7 +92,7 @@ def get_descriptor(bait_set, pooled_normals, preservation_types, run_ids, sample # sample_name is FROZENPOOLEDNORMAL unless FFPE is in any of the preservation types # in preservation_types; plc = preservations lower case - plc = set([x.lower() for x in preservation_types]) + plc = set([x.lower() for x in preservation_types if x]) # filter out None/empty values run_ids_suffix_list = [i for i in run_ids if i] # remove empty or false string values run_ids_suffix = "_".join(set(run_ids_suffix_list)) sample_name = "FROZENPOOLEDNORMAL_" + run_ids_suffix @@ -156,8 +156,8 @@ def get_preservation_type(preservation_types, sample_origin): - if preservation_type is not empty: - return 'ffpe' if is_ffpe(); otherwise 'frozen' """ - plc = set([x.lower() for x in preservation_types]) # preservation lower case - solc = set([x.lower() for x in sample_origin]) # sample origin lower case + plc = set([x.lower() for x in preservation_types if x]) # preservation lower case + solc = set([x.lower() for x in sample_origin if x]) # sample origin lower case preservation = "" if not is_list_empty(plc) and not is_list_empty(solc): @@ -233,7 +233,7 @@ def build_preservation_query(data): Main logic: if FFPE in data, return FFPE query; else, return FROZEN query """ - plc = set([x.lower() for x in data]) + plc = set([x.lower() for x in data if x]) # filter out None/empty values value = "FROZEN" if "ffpe" in plc: value = "FFPE" diff --git a/runner/serializers.py b/runner/serializers.py index c0005a89..b19abfac 100644 --- a/runner/serializers.py +++ b/runner/serializers.py @@ -278,13 +278,6 @@ class RestartRunSerializer(serializers.Serializer): clean = serializers.BooleanField(default=False) -# TODO: Delete this -class RequestIdOperatorSerializer(serializers.Serializer): - request_ids = serializers.ListField(child=serializers.CharField(max_length=30), allow_empty=True) - run_ids = serializers.ListField(child=serializers.UUIDField(), allow_empty=True) - pipeline_name = serializers.CharField(max_length=100) - - class RequestIdsOperatorSerializer(serializers.Serializer): request_ids = serializers.ListField(child=serializers.CharField(max_length=30), allow_empty=True) pipeline = serializers.CharField(max_length=30, allow_null=False, allow_blank=False)