Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion beagle/__init__.py
Original file line number Diff line number Diff line change
@@ -1 +1 @@
__version__ = "2.2.1"
__version__ = "2.3.0"
3 changes: 1 addition & 2 deletions beagle/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion beagle_etl/admin.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
43 changes: 42 additions & 1 deletion beagle_etl/jobs/metadb_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
37 changes: 37 additions & 0 deletions beagle_etl/migrations/0047_add_scheduled_retry_count.py
Original file line number Diff line number Diff line change
@@ -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),
]
32 changes: 27 additions & 5 deletions beagle_etl/models.py
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -45,6 +47,7 @@ class SmileMessageStatus(IntEnum):
COMPLETED = 3
NOT_SUPPORTED = 4
FAILED = 5
RETRY = 6


class SMILEMessage(BaseModel):
Expand All @@ -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):
Expand All @@ -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"])
Expand Down
19 changes: 16 additions & 3 deletions beagle_etl/smile_message/objects/sample_object.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
3 changes: 2 additions & 1 deletion beagle_etl/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
43 changes: 42 additions & 1 deletion beagle_etl/tests/jobs/test_metadb.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down
2 changes: 1 addition & 1 deletion runner/operator/argos_operator/v2_3_0/argos_operator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
)
Expand Down
2 changes: 1 addition & 1 deletion runner/operator/argos_operator/v2_3_0/bin/make_sample.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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"
Expand Down
7 changes: 0 additions & 7 deletions runner/serializers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Loading