Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
5e340f0
initial force import feature
buehlere Jun 23, 2026
317ed9e
Update request_view.py
buehlere Jun 23, 2026
c38a41e
Update request_view.py
buehlere Jun 23, 2026
8db9b68
add create request to celery
buehlere Jun 23, 2026
6db5b7e
get rid of async on endpoint
buehlere Jun 23, 2026
27a4e12
Update metadb_jobs.py
buehlere Jun 23, 2026
e03ddea
run as background task, change url, add admin panel
buehlere Jul 1, 2026
481883a
Update sample_object.py
buehlere Jul 1, 2026
108175e
fix get on samples
buehlere Jul 1, 2026
162c56c
Bump gitpython from 3.1.37 to 3.1.58
dependabot[bot] Aug 9, 2026
91a9ed1
Bump django from 6.0.4 to 6.0.7
dependabot[bot] Aug 9, 2026
17e55ad
Merge branch 'develop' into feature/force-import
buehlere Aug 18, 2026
3dafc38
Merge pull request #1548 from mskcc/master
sivkovic Aug 18, 2026
b27dedc
trying postfix for failed imports
buehlere Aug 18, 2026
95a5499
add SMILE message lookup
buehlere Aug 19, 2026
b3e0bdc
Merge branch 'develop' into feature/force-import
buehlere Aug 24, 2026
04b5419
Merge pull request #1528 from mskcc/feature/force-import
buehlere Aug 24, 2026
1da827d
Implement default mapping for file_groups different then IGO
sivkovic Aug 24, 2026
fe86e98
Merge pull request #1547 from mskcc/dependabot/pip/django-6.0.7
sivkovic Aug 25, 2026
040307f
Merge pull request #1546 from mskcc/dependabot/pip/gitpython-3.1.58
sivkovic Aug 25, 2026
686d573
Fix directory path
sivkovic Aug 26, 2026
3aee35a
Merge pull request #1549 from mskcc/feature/stage_file_from_different…
sivkovic Aug 26, 2026
5c38acf
Version bump 2.4.0
sivkovic Aug 26, 2026
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.3.0"
__version__ = "2.4.0"
4 changes: 4 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 Expand Up @@ -488,6 +491,7 @@

FASTQ_DEFAULT_LOCATION_PREFIX = os.environ.get("BEAGLE_FASTQ_DEFAULT_LOCATION_PREFIX")
FASTQ_IRIS_LOCATION_PREFIX = os.environ.get("BEAGLE_FASTQ_IRIS_LOCATION_PREFIX")
FASTQ_DEFAULT_STAGING_PATH = os.environ.get("BEAGLE_FASTQ_DEFAULT_STAGING_PATH", "")

DEFAULT_LOG_PREFIX = os.environ.get("BEAGLE_DEFAULT_LOG_PREFIX", "")
DEFAULT_LOG_PATH = os.environ.get("BEAGLE_DEFAULT_LOG_PATH", "/tmp")
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
20 changes: 13 additions & 7 deletions file_manager/copy_service/copy_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,19 +25,25 @@ def copy(path_from, path_to):
os.chmod(path_to, settings.COPY_FILE_PERMISSION)

@staticmethod
def remap(gene_panel, path, mapping=settings.DEFAULT_MAPPING):
prefix, dst = CopyService._get_mapping(gene_panel, path, mapping)
def remap(gene_panel, path, file_group=settings.IMPORT_FILE_GROUP, mapping=settings.DEFAULT_MAPPING):
prefix, dst = CopyService._get_mapping(gene_panel, path, file_group, mapping)
if prefix and dst:
path = path.replace(prefix, dst)
logger.info("New path {path}".format(path=path))
return path

@staticmethod
def _get_mapping(gene_panel, path, mapping=settings.DEFAULT_MAPPING):
recipe_mapping = mapping.get(gene_panel, {})
for prefix, dst in recipe_mapping.items():
if path.startswith(prefix):
return prefix, dst
def _get_mapping(gene_panel, path, file_group=settings.IMPORT_FILE_GROUP, mapping=settings.DEFAULT_MAPPING):
if file_group == settings.IMPORT_FILE_GROUP:
recipe_mapping = mapping.get(gene_panel, {})
for prefix, dst in recipe_mapping.items():
if path.startswith(prefix):
return prefix, dst
else:
return (
settings.FASTQ_IRIS_LOCATION_PREFIX,
os.path.join(settings.FASTQ_DEFAULT_STAGING_PATH, file_group) + "/",
)
return None, None

@staticmethod
Expand Down
4 changes: 2 additions & 2 deletions file_manager/file_manager/file_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ def stage_sample(self, sample_id):
files_to_stage = 0
for f in files:
if not f.file.is_available:
new_path = CopyService.remap(gene_panel, f.file.path)
new_path = CopyService.remap(gene_panel, f.file.path, str(f.file.file_group.id))
if new_path != f.file.path:
files_to_stage += 1

Expand All @@ -55,7 +55,7 @@ def stage_file(self, file_obj, gene_panel, sample_job=None):
Returns: Task signature or None
"""
if not file_obj.is_available:
new_path = CopyService.remap(gene_panel, file_obj.path)
new_path = CopyService.remap(gene_panel, file_obj.path, str(file_obj.file_group.id))
if new_path != file_obj.path:
fp_job, created = FileProviderJob.objects.provide_file(
file_obj, file_obj.path, new_path, sample_job=sample_job
Expand Down
Loading
Loading