diff --git a/beagle/settings.py b/beagle/settings.py index 61a3be5fa..fe9f8350b 100644 --- a/beagle/settings.py +++ b/beagle/settings.py @@ -459,6 +459,9 @@ BEAGLE_NOTIFIER_EMAIL_ABOUT_NEW_USERS = os.environ.get("BEAGLE_NOTIFIER_EMAIL_ABOUT_NEW_USERS") BEAGLE_NOTIFIER_EMAIL_FROM = os.environ.get("BEAGLE_NOTIFIER_EMAIL_FROM") +SMTP_HOST = os.environ.get("SMTP_HOST", "localhost") +SMTP_PORT = int(os.environ.get("SMTP_PORT", "25")) + BEAGLE_NOTIFIER_VOYAGER_STATUS_EMAIL_TO = os.environ.get("BEAGLE_NOTIFIER_VOYAGER_STATUS_EMAIL_TO", "").split(",") BEAGLE_NOTIFIER_VOYAGER_STATUS_BLACKLIST = os.environ.get("BEAGLE_NOTIFIER_VOYAGER_STATUS_BLACKLIST", "").split(",") BEAGLE_NOTIFIER_VOYAGER_STATUS_PIPELINES = os.environ.get("BEAGLE_NOTIFIER_VOYAGER_STATUS_PIPELINES", "").split(",") diff --git a/beagle_etl/admin.py b/beagle_etl/admin.py index da0ae1689..dc5a8b1c2 100644 --- a/beagle_etl/admin.py +++ b/beagle_etl/admin.py @@ -1,5 +1,6 @@ from django.contrib import admin from django.contrib.admin import ModelAdmin +from django.conf import settings from lib.admin import link_relation from .models import ( Operator, @@ -7,6 +8,7 @@ SMILEMessage, RequestCallbackJob, ) +from .jobs.metadb_jobs import new_request from advanced_filters.admin import AdminAdvancedFiltersMixin @@ -30,8 +32,16 @@ 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", "gene_panel", "topic", "status") search_fields = ("request_id", "gene_panel") + actions = ["force_import"] + + @admin.action(description="Force import selected SMILE messages (skip validation)") + def force_import(self, request, queryset): + for message in queryset.filter(topic=settings.METADB_NATS_NEW_REQUEST): + new_request.delay(str(message.id), force_import=True) + self.message_user(request, f"Force import triggered for {queryset.count()} message(s).") class RequestCallbackJobAdmin(ModelAdmin): diff --git a/beagle_etl/celery.py b/beagle_etl/celery.py index d6d6145ff..d7b58feb6 100644 --- a/beagle_etl/celery.py +++ b/beagle_etl/celery.py @@ -50,6 +50,7 @@ def setup_task_logger(logger, *args, **kwargs): "beagle_etl.tasks.job_processor": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, "beagle_etl.tasks.process_smile_events": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, "beagle_etl.tasks.process_job_with_lock": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, + "beagle_etl.jobs.metadb_jobs.new_request": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, "beagle_etl.jobs.metadb_jobs.update_job": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, "beagle_etl.jobs.metadb_jobs.not_supported": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, "beagle_etl.jobs.metadb_jobs.request_callback": {"queue": settings.BEAGLE_DEFAULT_QUEUE}, diff --git a/beagle_etl/jobs/metadb_jobs.py b/beagle_etl/jobs/metadb_jobs.py index 570f37dee..d125aedfd 100644 --- a/beagle_etl/jobs/metadb_jobs.py +++ b/beagle_etl/jobs/metadb_jobs.py @@ -188,17 +188,17 @@ def request_update_notification(request_id): @shared_task -def new_request(message_id): +def new_request(message_id, force_import=False): message = SMILEMessage.objects.get(id=message_id) try: - data = RequestMetadata.from_dict(json.loads(message.message)) + data = RequestMetadata.from_dict(json.loads(message.message), force_import=force_import) except Exception as e: message.add_log(str(e)) message.failed() return - if not data.isCmoRequest: + if not data.isCmoRequest and not force_import: # Non CmoRequests not supported logger.info(f"Request {data.igoRequestId} is not CMO Request") message.add_log(f"Request {data.igoRequestId} is not CMO Request") diff --git a/beagle_etl/serializers.py b/beagle_etl/serializers.py index c2325dac6..702e3f2ea 100644 --- a/beagle_etl/serializers.py +++ b/beagle_etl/serializers.py @@ -1,5 +1,5 @@ from rest_framework import serializers -from .models import ETLConfiguration +from .models import ETLConfiguration, SMILEMessage, SmileMessageStatus def ValidateDict(value): @@ -29,3 +29,21 @@ class AssayUpdateSerializer(serializers.Serializer): class RequestIdLimsPullSerializer(serializers.Serializer): request_ids = serializers.ListField(child=serializers.CharField(max_length=30)) redelivery = serializers.BooleanField(default=False) + + +class SMILEMessageSerializer(serializers.ModelSerializer): + status = serializers.SerializerMethodField() + + class Meta: + model = SMILEMessage + fields = "__all__" + + def get_status(self, obj): + return SmileMessageStatus(obj.status).name + + +class SMILEMessageListSerializer(serializers.Serializer): + request_id = serializers.CharField(required=False) + topic = serializers.CharField(required=False) + gene_panel = serializers.CharField(required=False) + status = serializers.ChoiceField([(s.name, s.value) for s in SmileMessageStatus], allow_blank=True, required=False) diff --git a/beagle_etl/smile_message/objects/request_object.py b/beagle_etl/smile_message/objects/request_object.py index 3a6e332a0..9ce7d1e9a 100644 --- a/beagle_etl/smile_message/objects/request_object.py +++ b/beagle_etl/smile_message/objects/request_object.py @@ -41,15 +41,15 @@ class RequestMetadata: pooledNormals: Optional[List[str]] = None @classmethod - def from_dict(cls, data: Dict[str, Any]) -> "RequestMetadata": + def from_dict(cls, data: Dict[str, Any], force_import: bool = False) -> "RequestMetadata": """Deserialize from dictionary.""" # Handle nested status status_data = data.get("status", {}) status = RequestStatus(**status_data) if status_data else RequestStatus(False, "{}") # Handle nested samples - samples_data = data.get("samples", []) - samples = [SampleMetadata.from_dict(sample) for sample in samples_data] + samples_data = data.get("samples") or [] + samples = [SampleMetadata.from_dict(sample, force_import=force_import) for sample in samples_data] # Handle delivery date conversion delivery_date = None diff --git a/beagle_etl/smile_message/objects/sample_object.py b/beagle_etl/smile_message/objects/sample_object.py index 6b1dbc5e3..47e373649 100644 --- a/beagle_etl/smile_message/objects/sample_object.py +++ b/beagle_etl/smile_message/objects/sample_object.py @@ -1,6 +1,6 @@ import re import logging -from dataclasses import dataclass +from dataclasses import dataclass, field from typing import List, Optional, Dict, Any from django.conf import settings from beagle_etl.exceptions import ( @@ -217,9 +217,10 @@ class SampleMetadata: tubeId: Optional[str] = None cfDNA2dBarcode: Optional[str] = None cmoInfoIgoId: Optional[str] = None + skip_validation: bool = field(default=False, compare=False, repr=False) @classmethod - def from_dict(cls, data: Dict[str, Any]) -> "SampleMetadata": + def from_dict(cls, data: Dict[str, Any], force_import: bool = False) -> "SampleMetadata": """Deserialize SampleMetadata from dictionary.""" # Handle nested status status_data = data.get("status", {}) @@ -227,18 +228,23 @@ def from_dict(cls, data: Dict[str, Any]) -> "SampleMetadata": # Handle nested cmoSampleIdFields cmo_fields_data = data.get("cmoSampleIdFields", {}) - cmo_fields = CmoSampleIdFields(**cmo_fields_data) if cmo_fields_data else CmoSampleIdFields("", "", "", "") + cmo_fields = CmoSampleIdFields( + naToExtract=cmo_fields_data.get("naToExtract", ""), + normalizedPatientId=cmo_fields_data.get("normalizedPatientId", ""), + sampleType=cmo_fields_data.get("sampleType", ""), + recipe=cmo_fields_data.get("recipe", ""), + ) # Handle libraries - libraries_data = data.get("libraries", []) + libraries_data = data.get("libraries") or [] libraries = [Library.from_dict(lib) for lib in libraries_data] # Handle sample aliases - sample_aliases_data = data.get("sampleAliases", []) + sample_aliases_data = data.get("sampleAliases") or [] sample_aliases = [SampleAlias(**alias) for alias in sample_aliases_data] # Handle patient aliases - patient_aliases_data = data.get("patientAliases", []) + patient_aliases_data = data.get("patientAliases") or [] patient_aliases = [PatientAlias(**alias) for alias in patient_aliases_data] return cls( @@ -270,14 +276,17 @@ def from_dict(cls, data: Dict[str, Any]) -> "SampleMetadata": igoComplete=data.get("igoComplete"), status=status, cmoSampleIdFields=cmo_fields, - qcReports=data.get("qcReports", []), + qcReports=data.get("qcReports") or [], libraries=libraries, sampleAliases=sample_aliases, patientAliases=patient_aliases, additionalProperties=data.get("additionalProperties", {}), + skip_validation=force_import, ) def __post_init__(self): + if self.skip_validation: + return self._validate_primary_id() required = PANEL_REQUIRED_FIELDS.get(self.genePanel, DEFAULT_REQUIRED_FIELDS) self._validate_required_fields(required) diff --git a/beagle_etl/urls.py b/beagle_etl/urls.py index ca232e417..ee55eca07 100644 --- a/beagle_etl/urls.py +++ b/beagle_etl/urls.py @@ -1,6 +1,6 @@ from rest_framework import routers from django.urls import path, include -from beagle_etl.views import AssayViewSet +from beagle_etl.views import AssayViewSet, ForceImportView, SMILEMessageViewSet router = routers.DefaultRouter() @@ -8,4 +8,6 @@ urlpatterns = [ path("", include(router.urls)), path("assay", AssayViewSet.as_view()), + path("import/messages/", SMILEMessageViewSet.as_view({"get": "list"})), + path("import//", ForceImportView.as_view()), ] diff --git a/beagle_etl/views.py b/beagle_etl/views.py index b16b6c718..865a2b1c9 100644 --- a/beagle_etl/views.py +++ b/beagle_etl/views.py @@ -1,15 +1,68 @@ -from rest_framework import status +from django.conf import settings +from rest_framework import mixins, status from rest_framework.response import Response from rest_framework.generics import GenericAPIView -from beagle_etl.models import ETLConfiguration +from rest_framework.views import APIView +from rest_framework.viewsets import GenericViewSet +from rest_framework.permissions import IsAuthenticated +from beagle_etl.models import ETLConfiguration, SMILEMessage, SmileMessageStatus +from beagle_etl.jobs.metadb_jobs import new_request from drf_yasg.utils import swagger_auto_schema from .serializers import ( AssaySerializer, AssayElementSerializer, AssayUpdateSerializer, + SMILEMessageSerializer, + SMILEMessageListSerializer, ) +class ForceImportView(APIView): + permission_classes = (IsAuthenticated,) + + def post(self, _request, request_id): + message = ( + SMILEMessage.objects.filter(request_id=request_id, topic=settings.METADB_NATS_NEW_REQUEST) + .order_by("-created_date") + .first() + ) + if not message: + return Response( + {"detail": f"No new-request SMILEMessage found for request_id {request_id}."}, + status=status.HTTP_404_NOT_FOUND, + ) + new_request.delay(str(message.id), force_import=True) + return Response( + {"detail": f"Force import triggered for request {request_id} (message {message.id})."}, + status=status.HTTP_202_ACCEPTED, + ) + + +class SMILEMessageViewSet(mixins.ListModelMixin, GenericViewSet): + queryset = SMILEMessage.objects.order_by("-created_date").all() + serializer_class = SMILEMessageListSerializer + permission_classes = (IsAuthenticated,) + + @swagger_auto_schema(query_serializer=SMILEMessageListSerializer) + def list(self, request, *args, **kwargs): + serializer = self.get_serializer(data=request.query_params) + if not serializer.is_valid(): + return Response(serializer.errors, status=status.HTTP_400_BAD_REQUEST) + validated = serializer.validated_data + queryset = self.queryset + if validated.get("request_id"): + queryset = queryset.filter(request_id=validated["request_id"]) + if validated.get("topic"): + queryset = queryset.filter(topic=validated["topic"]) + if validated.get("gene_panel"): + queryset = queryset.filter(gene_panel=validated["gene_panel"]) + if validated.get("status"): + queryset = queryset.filter(status=SmileMessageStatus[validated["status"]].value) + page = self.paginate_queryset(queryset) + serializer = SMILEMessageSerializer(page, many=True) + return self.get_paginated_response(serializer.data) + + class AssayViewSet(GenericAPIView): serializer_class = AssaySerializer queryset = ETLConfiguration.objects.all() diff --git a/compose.yaml b/compose.yaml index 92d7abc17..344f31472 100644 --- a/compose.yaml +++ b/compose.yaml @@ -592,6 +592,14 @@ services: beagle_pgbouncer: condition: service_healthy restart: false + postfix: + image: boky/postfix:5.1.0-debian + container_name: ${BEAGLE_DEPLOYMENT}-postfix + restart: always + networks: + - voyager_net + environment: + - ALLOWED_SENDER_DOMAINS=mskcc.org volumes: postgres_path: driver: local diff --git a/notifier/email/email_client.py b/notifier/email/email_client.py index 167f8d4fc..dd833a213 100644 --- a/notifier/email/email_client.py +++ b/notifier/email/email_client.py @@ -5,6 +5,8 @@ from email.mime.text import MIMEText from email.mime.multipart import MIMEMultipart +from django.conf import settings + class EmailClient(object): logger = logging.getLogger(__name__) @@ -15,12 +17,13 @@ def __init__(self, email_to, email_from, subject, content): self.content = content self.email_from = email_from self.domain = "mskcc.org" - self.SMTP_server = "localhost" + self.SMTP_server = settings.SMTP_HOST + self.SMTP_port = settings.SMTP_PORT def send(self): server = None try: - server = smtplib.SMTP(self.SMTP_server) + server = smtplib.SMTP(self.SMTP_server, self.SMTP_port) msg = MIMEMultipart("alternative") msg["Subject"] = self.subject msg["From"] = self.email_from