Skip to content
9 changes: 9 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 @@ -31,6 +33,13 @@ class SMILEMessagesAdmin(AdminAdvancedFiltersMixin, ModelAdmin):
advanced_filter_fields = ("request_id", "topic", "status")
ordering = ("-created_date",)
list_display = ("created_date", "request_id", "topic", "status")
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 @@ -180,17 +180,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
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
3 changes: 2 additions & 1 deletion beagle_etl/urls.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
from rest_framework import routers
from django.urls import path, include
from beagle_etl.views import AssayViewSet
from beagle_etl.views import AssayViewSet, ForceImportView


router = routers.DefaultRouter()

urlpatterns = [
path("", include(router.urls)),
path("assay", AssayViewSet.as_view()),
path("import/<str:request_id>/", ForceImportView.as_view()),
]
27 changes: 26 additions & 1 deletion beagle_etl/views.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
from django.conf import settings
from rest_framework import 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.permissions import IsAuthenticated
from beagle_etl.models import ETLConfiguration, SMILEMessage
from beagle_etl.jobs.metadb_jobs import new_request
from drf_yasg.utils import swagger_auto_schema
from .serializers import (
AssaySerializer,
Expand All @@ -10,6 +14,27 @@
)


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 AssayViewSet(GenericAPIView):
serializer_class = AssaySerializer
queryset = ETLConfiguration.objects.all()
Expand Down
Loading