Skip to content
3 changes: 3 additions & 0 deletions beagle/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(",")
Expand Down
10 changes: 10 additions & 0 deletions beagle_etl/admin.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,14 @@
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,
ETLConfiguration,
SMILEMessage,
RequestCallbackJob,
)
from .jobs.metadb_jobs import new_request
from advanced_filters.admin import AdminAdvancedFiltersMixin


Expand All @@ -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):
Expand Down
1 change: 1 addition & 0 deletions beagle_etl/celery.py
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down
6 changes: 3 additions & 3 deletions beagle_etl/jobs/metadb_jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
20 changes: 19 additions & 1 deletion beagle_etl/serializers.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
from rest_framework import serializers
from .models import ETLConfiguration
from .models import ETLConfiguration, SMILEMessage, SmileMessageStatus


def ValidateDict(value):
Expand Down Expand Up @@ -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)
6 changes: 3 additions & 3 deletions beagle_etl/smile_message/objects/request_object.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
23 changes: 16 additions & 7 deletions beagle_etl/smile_message/objects/sample_object.py
Original file line number Diff line number Diff line change
@@ -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 (
Expand Down Expand Up @@ -217,28 +217,34 @@ 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", {})
status = SampleStatusObj(**status_data) if status_data else SampleStatusObj(False, "{}")

# 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(
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 3 additions & 1 deletion beagle_etl/urls.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,13 @@
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()

urlpatterns = [
path("", include(router.urls)),
path("assay", AssayViewSet.as_view()),
path("import/messages/", SMILEMessageViewSet.as_view({"get": "list"})),
path("import/<str:request_id>/", ForceImportView.as_view()),
]
57 changes: 55 additions & 2 deletions beagle_etl/views.py
Original file line number Diff line number Diff line change
@@ -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()
Expand Down
8 changes: 8 additions & 0 deletions compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 5 additions & 2 deletions notifier/email/email_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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__)
Expand All @@ -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
Expand Down
Loading