From ac960adbc812bdbefd4bd7aaa9f081320f1626e9 Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Fri, 5 Sep 2025 14:32:21 -0400 Subject: [PATCH 1/7] write gene id file if missing --- seqr/utils/search/add_data_utils.py | 25 ++++++++++++++++++++++--- seqr/views/apis/data_manager_api.py | 2 +- seqr/views/utils/airflow_utils.py | 2 +- seqr/views/utils/export_utils.py | 9 +++++++-- 4 files changed, 31 insertions(+), 7 deletions(-) diff --git a/seqr/utils/search/add_data_utils.py b/seqr/utils/search/add_data_utils.py index 848ba1e5a7..ce5b9f3491 100644 --- a/seqr/utils/search/add_data_utils.py +++ b/seqr/utils/search/add_data_utils.py @@ -3,9 +3,10 @@ from django.db.models import F from typing import Callable -from reference_data.models import GENOME_VERSION_LOOKUP +from reference_data.models import GeneInfo, GENOME_VERSION_LOOKUP from seqr.models import Sample, Individual, Project from seqr.utils.communication_utils import send_project_notification +from seqr.utils.file_utils import does_file_exist from seqr.utils.logging_utils import SeqrLogger from seqr.utils.search.utils import backend_specific_call from seqr.utils.search.elasticsearch.es_utils import validate_es_index_metadata_and_get_samples @@ -117,7 +118,7 @@ def _format_loading_pipeline_variables( return variables def prepare_data_loading_request(projects: list[Project], individual_ids: list[int], sample_type: str, dataset_type: str, genome_version: str, - data_path: str, user: User, pedigree_dir: str, raise_pedigree_error: bool = False, + data_path: str, user: User, load_data_dir: str, raise_pedigree_error: bool = False, skip_validation: bool = False, skip_check_sex_and_relatedness: bool = False, vcf_sample_id_map=None): variables = _format_loading_pipeline_variables( projects, @@ -130,8 +131,9 @@ def prepare_data_loading_request(projects: list[Project], individual_ids: list[i variables['skip_validation'] = True if skip_check_sex_and_relatedness: variables['skip_check_sex_and_relatedness'] = True - file_path = _get_pedigree_path(pedigree_dir, genome_version, sample_type, dataset_type) + file_path = _get_pedigree_path(load_data_dir, genome_version, sample_type, dataset_type) _upload_data_loading_files(individual_ids, vcf_sample_id_map or {}, user, file_path, raise_pedigree_error) + backend_specific_call(lambda *args: None, lambda *args: None, _write_gene_id_file)(load_data_dir, user) return variables, file_path @@ -172,6 +174,23 @@ def _upload_data_loading_files(individual_ids: list[int], vcf_sample_id_map: dic raise e +def _write_gene_id_file(load_data_dir, user): + file_name = 'db_id_to_gene_id' + if does_file_exist(f'{load_data_dir}/{file_name}.csv.gz'): + return + + gene_data_loaded = (GeneInfo.objects.filter(gencode_release=GeneInfo.ALL_GENCODE_VERSIONS[0]).exists() and + GeneInfo.objects.filter(gencode_release=GeneInfo.ALL_GENCODE_VERSIONS[-1]).exists()) + if not gene_data_loaded: + raise ValueError( + 'Gene reference data is not yet loaded. If this is a new seqr installation, wait for the initial data load ' + 'to complete. If this is an existing installation, see the documentation for updating data in seqr.' + ) + gene_data = GeneInfo.objects.all().values('gene_id', db_id=F('id')) + file_config = (file_name, ['db_id', 'gene_id'], gene_data) + write_multiple_files([file_config], load_data_dir, user, file_format='csv', gzip_file=True) + + def _get_pedigree_path(pedigree_dir: str, genome_version: str, sample_type: str, dataset_type: str): dag_dataset_type = _dag_dataset_type(sample_type, dataset_type) return f'{pedigree_dir}/{GENOME_VERSION_LOOKUP[genome_version]}/{dag_dataset_type}/pedigrees/{sample_type}' diff --git a/seqr/views/apis/data_manager_api.py b/seqr/views/apis/data_manager_api.py index 82583479fa..5910d97003 100644 --- a/seqr/views/apis/data_manager_api.py +++ b/seqr/views/apis/data_manager_api.py @@ -354,7 +354,7 @@ def load_data(request): ) else: request_json, _ = prepare_data_loading_request( - *loading_args, **loading_kwargs, pedigree_dir=LOADING_DATASETS_DIR, raise_pedigree_error=True, + *loading_args, **loading_kwargs, load_data_dir=LOADING_DATASETS_DIR, raise_pedigree_error=True, ) response = requests.post(f'{PIPELINE_RUNNER_SERVER}/loading_pipeline_enqueue', json=request_json, timeout=60) if response.status_code == 409: diff --git a/seqr/views/utils/airflow_utils.py b/seqr/views/utils/airflow_utils.py index 2bc1d089bf..d6f0dbb04a 100644 --- a/seqr/views/utils/airflow_utils.py +++ b/seqr/views/utils/airflow_utils.py @@ -24,7 +24,7 @@ def trigger_airflow_data_loading(*args, user: User, success_message: str, succes error_message: str, is_internal: bool = False, **kwargs): success = True updated_variables, gs_path = prepare_data_loading_request( - *args, user, pedigree_dir=SEQR_V3_PEDIGREE_GS_PATH, **kwargs, + *args, user, load_data_dir=SEQR_V3_PEDIGREE_GS_PATH, **kwargs, ) updated_variables['sample_source'] = 'Broad_Internal' if is_internal else 'AnVIL' upload_info = [f'Pedigree files have been uploaded to {gs_path}'] diff --git a/seqr/views/utils/export_utils.py b/seqr/views/utils/export_utils.py index 71261e921c..98a1af1704 100644 --- a/seqr/views/utils/export_utils.py +++ b/seqr/views/utils/export_utils.py @@ -1,3 +1,4 @@ +import gzip import openpyxl as xl import os from tempfile import NamedTemporaryFile, TemporaryDirectory @@ -89,14 +90,18 @@ def export_multiple_files(files, zip_filename, **kwargs): return response -def write_multiple_files(files, file_path, user, **kwargs): +def write_multiple_files(files, file_path, user, gzip_file=False, **kwargs): is_gs_path = is_google_bucket_file_path(file_path) if not is_gs_path: os.makedirs(file_path, exist_ok=True) with TemporaryDirectory() as temp_dir_name: dir_name = temp_dir_name if is_gs_path else file_path + open_func = gzip.open if gzip_file else open for filename, content in _format_files_content(files, **kwargs): - with open(f'{dir_name}/{filename}', 'w') as f: + file_path = f'{dir_name}/{filename}' + if gzip_file: + file_path += '.gz' + with open_func(file_path, 'w') as f: f.write(content) if is_gs_path: mv_file_to_gs(f'{temp_dir_name}/*', f'{file_path}/', user) From a71b86e69266243c354f5666fdf2d0204c10acdb Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Fri, 5 Sep 2025 14:34:35 -0400 Subject: [PATCH 2/7] no backend check for file write --- seqr/utils/search/add_data_utils.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/seqr/utils/search/add_data_utils.py b/seqr/utils/search/add_data_utils.py index ce5b9f3491..896e856f47 100644 --- a/seqr/utils/search/add_data_utils.py +++ b/seqr/utils/search/add_data_utils.py @@ -133,7 +133,7 @@ def prepare_data_loading_request(projects: list[Project], individual_ids: list[i variables['skip_check_sex_and_relatedness'] = True file_path = _get_pedigree_path(load_data_dir, genome_version, sample_type, dataset_type) _upload_data_loading_files(individual_ids, vcf_sample_id_map or {}, user, file_path, raise_pedigree_error) - backend_specific_call(lambda *args: None, lambda *args: None, _write_gene_id_file)(load_data_dir, user) + _write_gene_id_file(load_data_dir, user) return variables, file_path From 1cc065a9f0378b4d825f28b3a1f3fc8da9133018 Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Fri, 5 Sep 2025 15:53:56 -0400 Subject: [PATCH 3/7] test write gene ids --- seqr/utils/search/add_data_utils.py | 6 +- seqr/views/apis/data_manager_api_tests.py | 105 ++++++++++++++++------ seqr/views/utils/export_utils.py | 6 +- 3 files changed, 85 insertions(+), 32 deletions(-) diff --git a/seqr/utils/search/add_data_utils.py b/seqr/utils/search/add_data_utils.py index 896e856f47..8bd9157432 100644 --- a/seqr/utils/search/add_data_utils.py +++ b/seqr/utils/search/add_data_utils.py @@ -179,14 +179,14 @@ def _write_gene_id_file(load_data_dir, user): if does_file_exist(f'{load_data_dir}/{file_name}.csv.gz'): return - gene_data_loaded = (GeneInfo.objects.filter(gencode_release=GeneInfo.ALL_GENCODE_VERSIONS[0]).exists() and - GeneInfo.objects.filter(gencode_release=GeneInfo.ALL_GENCODE_VERSIONS[-1]).exists()) + gene_data_loaded = (GeneInfo.objects.filter(gencode_release=int(GeneInfo.CURRENT_VERSION)).exists() and + GeneInfo.objects.filter(gencode_release=int(GeneInfo.ALL_GENCODE_VERSIONS[-1])).exists()) if not gene_data_loaded: raise ValueError( 'Gene reference data is not yet loaded. If this is a new seqr installation, wait for the initial data load ' 'to complete. If this is an existing installation, see the documentation for updating data in seqr.' ) - gene_data = GeneInfo.objects.all().values('gene_id', db_id=F('id')) + gene_data = GeneInfo.objects.all().values('gene_id', db_id=F('id')).order_by('id') file_config = (file_name, ['db_id', 'gene_id'], gene_data) write_multiple_files([file_config], load_data_dir, user, file_format='csv', gzip_file=True) diff --git a/seqr/views/apis/data_manager_api_tests.py b/seqr/views/apis/data_manager_api_tests.py index b692d37347..b2496fe706 100644 --- a/seqr/views/apis/data_manager_api_tests.py +++ b/seqr/views/apis/data_manager_api_tests.py @@ -1376,7 +1376,7 @@ def test_get_loaded_projects(self): self._assert_expected_airtable_errors(url) - def _assert_expected_pm_access(self, get_response): + def _assert_expected_pm_access(self, get_response, mock_current_gene_version=None): response = get_response() self.assertEqual(response.status_code, 200) self.login_data_manager_user() @@ -1387,10 +1387,13 @@ def _assert_expected_airtable_errors(self, url): @responses.activate @mock.patch('seqr.views.utils.airtable_utils.BASE_URL', 'https://seqr.broadinstitute.org/') + @mock.patch('reference_data.models.GeneInfo.CURRENT_VERSION') @mock.patch('seqr.views.utils.export_utils.os.makedirs') + @mock.patch('seqr.views.utils.export_utils.gzip.open') @mock.patch('seqr.views.utils.export_utils.open') @mock.patch('seqr.views.utils.export_utils.TemporaryDirectory') - def test_load_data(self, mock_temp_dir, mock_open, mock_mkdir): + def test_load_data(self, mock_temp_dir, mock_open, mock_gzip_open, mock_mkdir, mock_current_gene_version): + mock_current_gene_version.__int__.return_value = 27 url = reverse(load_data) self.check_pm_login(url) @@ -1409,13 +1412,14 @@ def test_load_data(self, mock_temp_dir, mock_open, mock_mkdir): self.reset_logs() responses.calls.reset() + self._set_file_not_found(has_mv_commands=True) response = self._assert_expected_pm_access( - lambda: self.client.post(url, content_type='application/json', data=json.dumps(body)) + lambda: self.client.post(url, content_type='application/json', data=json.dumps(body)), ) self.assertDictEqual(response.json(), {'success': True}) self._assert_expected_load_data_requests(sample_type='WES', skip_validation=True) - self._has_expected_ped_files(mock_open, mock_mkdir, 'SNV_INDEL', sample_type='WES', has_remap=bool(self.MOCK_AIRTABLE_KEY)) + self._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, 'SNV_INDEL', sample_type='WES', has_remap=bool(self.MOCK_AIRTABLE_KEY)) dag_json = { 'projects_to_run': [ @@ -1432,7 +1436,9 @@ def test_load_data(self, mock_temp_dir, mock_open, mock_mkdir): # Test loading trigger error self._set_loading_trigger_error() + self._set_file_not_found(has_mv_commands=True) mock_open.reset_mock() + mock_gzip_open.reset_mock() mock_mkdir.reset_mock() responses.calls.reset() self.reset_logs() @@ -1440,18 +1446,21 @@ def test_load_data(self, mock_temp_dir, mock_open, mock_mkdir): del body['skipValidation'] del dag_json['skip_validation'] body.update({'datasetType': 'SV', 'filePath': f'{self.CALLSET_DIR}/sv_callset.vcf'}) - self._trigger_error(url, body, dag_json, mock_open, mock_mkdir) + self._trigger_error(url, body, dag_json, mock_open, mock_gzip_open, mock_mkdir) + self._set_file_found() responses.calls.reset() mock_open.reset_mock() + mock_gzip_open.reset_mock() mock_mkdir.reset_mock() body.update({'sampleType': 'WGS', 'projects': [json.dumps(self.PROJECT_OPTION)], 'vcfSamples': VCF_SAMPLES}) del body['datasetType'] response = self.client.post(url, content_type='application/json', data=json.dumps(body)) - self._test_load_single_project(mock_open, mock_mkdir, response, url=url, body=body) + self._test_load_single_project(mock_open, mock_gzip_open, mock_mkdir, response, url=url, body=body) # Test write pedigree error self.reset_logs() + self._set_file_found() responses.calls.reset() mock_mkdir.reset_mock() mock_open.reset_mock() @@ -1466,13 +1475,25 @@ def test_load_data(self, mock_temp_dir, mock_open, mock_mkdir): }), ]) - def _trigger_error(self, url, body, dag_json, mock_open, mock_mkdir): + # Test when gene data is not fully loaded + responses.calls.reset() + self._set_file_not_found(has_mv_commands=True) + mock_open.side_effect = None + mock_current_gene_version.__int__.return_value = 39 + response = self.client.post(url, content_type='application/json', data=json.dumps(body)) + self.assertEqual(response.status_code, 500) + self.assertDictEqual(response.json(), { + 'error': 'Gene reference data is not yet loaded. If this is a new seqr installation, wait for the initial data load to complete. If this is an existing installation, see the documentation for updating data in seqr.', + }) + self.assertEqual(len(responses.calls), 0) + + def _trigger_error(self, url, body, dag_json, mock_open, mock_gzip_open, mock_mkdir): response = self.client.post(url, content_type='application/json', data=json.dumps(body)) self._assert_expected_load_data_requests(trigger_error=True, dataset_type='GCNV', sample_type='WES') self._assert_trigger_error(response, body, dag_json) - self._has_expected_ped_files(mock_open, mock_mkdir, 'GCNV', sample_type='WES') + self._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, 'GCNV', sample_type='WES') - def _has_expected_ped_files(self, mock_open, mock_mkdir, dataset_type, sample_type='WGS', single_project=False, has_remap=False): + def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset_type, sample_type='WGS', single_project=False, has_remap=False, has_gene_id_file=False): mock_open.assert_has_calls([ mock.call(f'{self._local_pedigree_path(dataset_type, sample_type)}/{project}_pedigree.tsv', 'w') for project in self.PROJECTS[(1 if single_project else 0):] @@ -1483,6 +1504,16 @@ def _has_expected_ped_files(self, mock_open, mock_mkdir, dataset_type, sample_ty ] self.assertEqual(len(files), 1 if single_project else 2) + if has_gene_id_file: + mock_gzip_open.assert_not_called() + else: + mock_gzip_open.assert_called_once_with(f'{self.LOCAL_WRITE_DIR}/db_id_to_gene_id.csv.gz', 'w') + file = [ + row.split(',') for row in mock_gzip_open.return_value.__enter__.return_value.write.call_args.args[0].split('\n') + ] + self.assertEqual(len(file), 59) + self.assertListEqual(file[:3], [['db_id', 'gene_id'], ['1', 'ENSG00000223972'], ['2', 'ENSG00000227232']]) + num_rows = 7 if self.MOCK_AIRTABLE_KEY else 8 pedigree_header = PEDIGREE_HEADER + ['VCF_ID'] if has_remap else PEDIGREE_HEADER if not single_project: @@ -1499,10 +1530,10 @@ def _has_expected_ped_files(self, mock_open, mock_mkdir, dataset_type, sample_ty ['R0004_non_analyst_project', 'F000014_14', 'fam14', 'NA21987', '', '', 'M'] + ([''] if has_remap else []), ]) - def _test_load_single_project(self, mock_open, mock_mkdir, response, *args, **kwargs): + def _test_load_single_project(self, mock_open, mock_gzip_open, mock_mkdir, response, *args, **kwargs): self.assertEqual(response.status_code, 200) self.assertDictEqual(response.json(), {'success': True}) - self._has_expected_ped_files(mock_open, mock_mkdir, 'SNV_INDEL', single_project=True) + self._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, 'SNV_INDEL', single_project=True, has_gene_id_file=True) # Only a DAG trigger, no airtable calls as there is no previously loaded WGS SNV_INDEL data for these samples self.assertEqual(len(responses.calls), 1) @@ -1550,6 +1581,7 @@ class LocalDataManagerAPITest(AuthenticationTestCase, DataManagerAPITest): fixtures = ['users', '1kg_project', 'reference_data'] TRIGGER_CALLSET_DIR = '/local_datasets' + LOCAL_WRITE_DIR = TRIGGER_CALLSET_DIR CALLSET_DIR = '' PROJECT_OPTION = PROJECT_OPTION WGS_PROJECT_OPTIONS = [EMPTY_PROJECT_OPTION] @@ -1620,12 +1652,13 @@ def _assert_expected_load_data_requests(self, dataset_type='SNV_INDEL', sample_t body['skip_validation'] = True self.assertDictEqual(json.loads(responses.calls[0].request.body), body) + @staticmethod def _local_pedigree_path(dataset_type, sample_type): return f'/local_datasets/GRCh38/{dataset_type}/pedigrees/{sample_type}' - def _has_expected_ped_files(self, mock_open, mock_mkdir, dataset_type, *args, sample_type='WGS', **kwargs): - super()._has_expected_ped_files(mock_open, mock_mkdir, dataset_type, *args, sample_type, **kwargs) + def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type='WGS', **kwargs): + super()._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type, **kwargs) mock_mkdir.assert_called_once_with(self._local_pedigree_path(dataset_type, sample_type), exist_ok=True) def _assert_success_notification(self, dag_json): @@ -1635,8 +1668,8 @@ def _assert_success_notification(self, dag_json): def _set_loading_trigger_error(self): responses.add(responses.POST, PIPELINE_RUNNER_URL, status=400) - def _trigger_error(self, url, body, dag_json, mock_open, mock_mkdir): - super()._trigger_error(url, body, dag_json, mock_open, mock_mkdir) + def _trigger_error(self, url, body, dag_json, mock_open, mock_gzip_open, mock_mkdir): + super()._trigger_error(url, body, dag_json, mock_open, mock_gzip_open, mock_mkdir) responses.add(responses.POST, PIPELINE_RUNNER_URL, status=409) response = self.client.post(url, content_type='application/json', data=json.dumps(body)) @@ -1711,7 +1744,7 @@ def setUp(self): self.addCleanup(patcher.stop) super().setUp() - def _set_file_not_found(self, file_name=None, sample_guid=None, list_files=False): + def _set_file_not_found(self, file_name=None, sample_guid=None, list_files=False, has_mv_commands=False): self.mock_file_iter.stdout = [] self.mock_does_file_exist.wait.return_value = 1 error = b'CommandException: One or more URLs matched no objects' @@ -1719,12 +1752,23 @@ def _set_file_not_found(self, file_name=None, sample_guid=None, list_files=False self.mock_does_file_exist.communicate.return_value = (b'', error) else: self.mock_does_file_exist.stdout = [error] - self.mock_subprocess.side_effect = [self.mock_does_file_exist] + subprocess_side_effect = [self.mock_does_file_exist] + if has_mv_commands: + mock_mv = mock.MagicMock() + mock_mv.wait.return_value = 0 + subprocess_side_effect = [mock_mv, self.mock_does_file_exist, mock_mv] + self.mock_subprocess.side_effect = subprocess_side_effect return [ (f'==> gsutil ls gs://seqr-scratch-temp/{file_name}/{sample_guid}.json.gz', None), ('CommandException: One or more URLs matched no objects', None), ] + def _set_file_found(self): + self.mock_does_file_exist.wait.return_value = 0 + self.mock_subprocess.side_effect = [ + self.mock_does_file_exist, self.mock_does_file_exist, self.mock_does_file_exist, self.mock_does_file_exist, + ] + def _add_file_iter(self, stdout, is_gz=True): self.mock_does_file_exist.wait.return_value = 0 if not is_gz: @@ -1862,9 +1906,9 @@ def _assert_trigger_error(self, response, body, dag_json, **kwargs): """ self.mock_slack.assert_called_once_with(SEQR_SLACK_LOADING_NOTIFICATION_CHANNEL, error_message) - def _trigger_error(self, url, body, dag_json, mock_open, mock_mkdir): + def _trigger_error(self, url, body, dag_json, mock_open, mock_gzip_open, mock_mkdir): body['vcfSamples'] = None - super()._trigger_error(url, body, dag_json, mock_open, mock_mkdir) + super()._trigger_error(url, body, dag_json, mock_open, mock_gzip_open, mock_mkdir) responses.calls.reset() body['vcfSamples'] = ['ABC123', 'NA19675_1'] @@ -1899,20 +1943,21 @@ def _trigger_error(self, url, body, dag_json, mock_open, mock_mkdir): self._assert_expected_airtable_call(required_sample_field='SV_CallsetPath', project_guid='R0004_non_analyst_project') self.mock_authorized_session.reset_mock() - def _test_load_single_project(self, mock_open, mock_mkdir, response, *args, url=None, body=None, **kwargs): - super()._test_load_single_project(mock_open, mock_mkdir, response, url, body) + def _test_load_single_project(self, mock_open, mock_gzip_open, mock_mkdir, response, *args, url=None, body=None, **kwargs): + super()._test_load_single_project(mock_open, mock_gzip_open, mock_mkdir, response, url, body) self.ADDITIONAL_REQUEST_COUNT = 0 self.assert_airflow_loading_calls(offset=0, dataset_type='SNV_INDEL', trigger_error=True) responses.calls.reset() mock_open.reset_mock() + mock_gzip_open.reset_mock() mock_mkdir.reset_mock() body['projects'] = [json.dumps(option) for option in self.PROJECT_OPTIONS] body['sampleType'] = 'WES' response = self.client.post(url, content_type='application/json', data=json.dumps(body)) self.assertEqual(response.status_code, 200) self.assertDictEqual(response.json(), {'success': True}) - self._has_expected_ped_files(mock_open, mock_mkdir, 'SNV_INDEL', sample_type='WES') + self._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, 'SNV_INDEL', sample_type='WES', has_gene_id_file=True) self.assertEqual(len(responses.calls), 2) self.assert_expected_airtable_call( call_index=0, @@ -1925,14 +1970,22 @@ def _test_load_single_project(self, mock_open, mock_mkdir, response, *args, url= def _local_pedigree_path(*args): return '/mock/tmp' - def _has_expected_ped_files(self, mock_open, mock_mkdir, dataset_type, *args, sample_type='WGS', **kwargs): - super()._has_expected_ped_files(mock_open, mock_mkdir, dataset_type, sample_type, **kwargs) + def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type='WGS', has_gene_id_file=False, **kwargs): + super()._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, dataset_type, sample_type, has_gene_id_file=has_gene_id_file, **kwargs) mock_mkdir.assert_not_called() - self.mock_subprocess.assert_called_once_with( + expected_calls = [mock.call( f'gsutil mv /mock/tmp/* gs://seqr-loading-temp/v3.1/GRCh38/{dataset_type}/pedigrees/{sample_type}/', stdout=-1, stderr=-2, shell=True, # nosec - ) + ), mock.call( + 'gsutil ls gs://seqr-loading-temp/v3.1/db_id_to_gene_id.csv.gz', stdout=-1, stderr=-2, shell=True, + )] + if not has_gene_id_file: + expected_calls.append(mock.call( + 'gsutil mv /mock/tmp/* gs://seqr-loading-temp/v3.1/', stdout=-1, stderr=-2, shell=True, + )) + self.assertEqual(self.mock_subprocess.call_count, len(expected_calls)) + self.mock_subprocess.assert_has_calls(expected_calls) self.mock_subprocess.reset_mock() def _assert_write_pedigree_error(self, response): diff --git a/seqr/views/utils/export_utils.py b/seqr/views/utils/export_utils.py index 98a1af1704..0bc3c8e944 100644 --- a/seqr/views/utils/export_utils.py +++ b/seqr/views/utils/export_utils.py @@ -98,10 +98,10 @@ def write_multiple_files(files, file_path, user, gzip_file=False, **kwargs): dir_name = temp_dir_name if is_gs_path else file_path open_func = gzip.open if gzip_file else open for filename, content in _format_files_content(files, **kwargs): - file_path = f'{dir_name}/{filename}' + current_file = f'{dir_name}/{filename}' if gzip_file: - file_path += '.gz' - with open_func(file_path, 'w') as f: + current_file += '.gz' + with open_func(current_file, 'w') as f: f.write(content) if is_gs_path: mv_file_to_gs(f'{temp_dir_name}/*', f'{file_path}/', user) From 0c5304d9e2ac5d149c049aa1c47c3fd2907a6921 Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Fri, 5 Sep 2025 16:02:03 -0400 Subject: [PATCH 4/7] test write gene ids local --- seqr/views/apis/data_manager_api_tests.py | 19 ++++++++++++++----- 1 file changed, 14 insertions(+), 5 deletions(-) diff --git a/seqr/views/apis/data_manager_api_tests.py b/seqr/views/apis/data_manager_api_tests.py index b2496fe706..cf4ecd2e78 100644 --- a/seqr/views/apis/data_manager_api_tests.py +++ b/seqr/views/apis/data_manager_api_tests.py @@ -1511,7 +1511,7 @@ def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset file = [ row.split(',') for row in mock_gzip_open.return_value.__enter__.return_value.write.call_args.args[0].split('\n') ] - self.assertEqual(len(file), 59) + self.assertEqual(len(file), self.NUM_FIXTURE_GENES) self.assertListEqual(file[:3], [['db_id', 'gene_id'], ['1', 'ENSG00000223972'], ['2', 'ENSG00000227232']]) num_rows = 7 if self.MOCK_AIRTABLE_KEY else 8 @@ -1580,6 +1580,7 @@ def _assert_expected_search_data_update(self, response): class LocalDataManagerAPITest(AuthenticationTestCase, DataManagerAPITest): fixtures = ['users', '1kg_project', 'reference_data'] + NUM_FIXTURE_GENES = 52 TRIGGER_CALLSET_DIR = '/local_datasets' LOCAL_WRITE_DIR = TRIGGER_CALLSET_DIR CALLSET_DIR = '' @@ -1617,11 +1618,14 @@ def setUp(self): self.addCleanup(patcher.stop) super().setUp() - def _set_file_not_found(self, file_name=None, sample_guid=None, list_files=False): + def _set_file_not_found(self, file_name=None, sample_guid=None, list_files=False, has_mv_commands=False): self.mock_does_file_exist.return_value = False self.mock_file_iter.return_value = [] return [] + def _set_file_found(self): + self.mock_does_file_exist.return_value = True + def _add_file_iter(self, stdout, is_gz=True): self.mock_does_file_exist.return_value = True file_iter = self.mock_file_iter if is_gz else self.mock_unzipped_file_iter @@ -1657,9 +1661,13 @@ def _assert_expected_load_data_requests(self, dataset_type='SNV_INDEL', sample_t def _local_pedigree_path(dataset_type, sample_type): return f'/local_datasets/GRCh38/{dataset_type}/pedigrees/{sample_type}' - def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type='WGS', **kwargs): - super()._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type, **kwargs) - mock_mkdir.assert_called_once_with(self._local_pedigree_path(dataset_type, sample_type), exist_ok=True) + def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type='WGS', has_gene_id_file=False, **kwargs): + super()._has_expected_ped_files(mock_open, mock_gzip_open, mock_mkdir, dataset_type, *args, sample_type, has_gene_id_file=has_gene_id_file, **kwargs) + call_paths = [self._local_pedigree_path(dataset_type, sample_type)] + if not has_gene_id_file: + call_paths.append(self.LOCAL_WRITE_DIR) + self.assertEqual(mock_mkdir.call_count, len(call_paths)) + mock_mkdir.assert_has_calls([mock.call(call_path, exist_ok=True) for call_path in call_paths]) def _assert_success_notification(self, dag_json): self.maxDiff = None @@ -1716,6 +1724,7 @@ def _test_validate_dataset_type(self, url): class AnvilDataManagerAPITest(AirflowTestCase, DataManagerAPITest): fixtures = ['users', 'social_auth', '1kg_project', 'reference_data', 'clickhouse_search'] + NUM_FIXTURE_GENES = 59 LOADING_PROJECT_GUID = NON_ANALYST_PROJECT_GUID CALLSET_DIR = 'gs://test_bucket' TRIGGER_CALLSET_DIR = CALLSET_DIR From 1cc7d5e831d2368839b86d29d4f8db17150e5664 Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Fri, 5 Sep 2025 16:06:37 -0400 Subject: [PATCH 5/7] codacy --- seqr/views/apis/data_manager_api_tests.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/seqr/views/apis/data_manager_api_tests.py b/seqr/views/apis/data_manager_api_tests.py index cf4ecd2e78..4cb5a6f4b1 100644 --- a/seqr/views/apis/data_manager_api_tests.py +++ b/seqr/views/apis/data_manager_api_tests.py @@ -1987,11 +1987,11 @@ def _has_expected_ped_files(self, mock_open, mock_gzip_open, mock_mkdir, dataset f'gsutil mv /mock/tmp/* gs://seqr-loading-temp/v3.1/GRCh38/{dataset_type}/pedigrees/{sample_type}/', stdout=-1, stderr=-2, shell=True, # nosec ), mock.call( - 'gsutil ls gs://seqr-loading-temp/v3.1/db_id_to_gene_id.csv.gz', stdout=-1, stderr=-2, shell=True, + 'gsutil ls gs://seqr-loading-temp/v3.1/db_id_to_gene_id.csv.gz', stdout=-1, stderr=-2, shell=True, # nosec )] if not has_gene_id_file: expected_calls.append(mock.call( - 'gsutil mv /mock/tmp/* gs://seqr-loading-temp/v3.1/', stdout=-1, stderr=-2, shell=True, + 'gsutil mv /mock/tmp/* gs://seqr-loading-temp/v3.1/', stdout=-1, stderr=-2, shell=True, # nosec )) self.assertEqual(self.mock_subprocess.call_count, len(expected_calls)) self.mock_subprocess.assert_has_calls(expected_calls) From 7896499ad8d13d42f58ba50c81b71bd6db8c26ac Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Fri, 5 Sep 2025 16:26:55 -0400 Subject: [PATCH 6/7] fix tests --- seqr/views/apis/anvil_workspace_api_tests.py | 23 ++++++++++++++++---- 1 file changed, 19 insertions(+), 4 deletions(-) diff --git a/seqr/views/apis/anvil_workspace_api_tests.py b/seqr/views/apis/anvil_workspace_api_tests.py index ab6c6d60a3..ae8952fdea 100644 --- a/seqr/views/apis/anvil_workspace_api_tests.py +++ b/seqr/views/apis/anvil_workspace_api_tests.py @@ -505,6 +505,7 @@ def _test_get_workspace_files(self, url, response_key, expected_files, mock_subp ]) +@mock.patch('reference_data.models.GeneInfo.CURRENT_VERSION', '27') class LoadAnvilDataAPITest(AirflowTestCase, AirtableTest): fixtures = ['users', 'social_auth', 'reference_data', '1kg_project'] @@ -561,6 +562,9 @@ def setUp(self): patcher = mock.patch('seqr.views.utils.export_utils.open') self.mock_temp_open = patcher.start() self.addCleanup(patcher.stop) + patcher = mock.patch('seqr.views.utils.export_utils.gzip.open') + self.mock_gzip_temp_open = patcher.start() + self.addCleanup(patcher.stop) patcher = mock.patch('seqr.views.apis.anvil_workspace_api.logger') self.mock_api_logger = patcher.start() self.addCleanup(patcher.stop) @@ -796,10 +800,16 @@ def _assert_valid_operation(self, project, test_add_data=True): '\n'.join(['\t'.join(row) for row in [header] + rows]) ) + self.mock_gzip_temp_open.assert_called_with(f'{TEMP_PATH}/db_id_to_gene_id.csv.gz', 'w') + gene_file = self.mock_gzip_temp_open.return_value.__enter__.return_value.write.call_args.args[0].split('\n') + self.assertEqual(len(gene_file), 52) + self.assertListEqual(gene_file[:3], ['db_id,gene_id', '1,ENSG00000223972', '2,ENSG00000227232']) + gs_path = f'gs://seqr-loading-temp/v3.1/{genome_version}/SNV_INDEL/pedigrees/WES/' - self.mock_mv_file.assert_called_with( - f'{TEMP_PATH}/*', gs_path, self.manager_user - ) + self.mock_mv_file.assert_has_calls([ + mock.call(f'{TEMP_PATH}/*', gs_path, self.manager_user), + mock.call(f'{TEMP_PATH}/*', 'gs://seqr-loading-temp/v3.1/', self.manager_user) + ]) self.assert_airflow_loading_calls(additional_tasks_check=test_add_data) @@ -866,11 +876,16 @@ def _assert_valid_operation(self, project, test_add_data=True): 'father__individual_id': None, 'sex': 'F', 'affected': 'N', 'notes': 'a individual note', 'features': [], }, individual_model_data) + @staticmethod + def _raise_move_file_error(from_path, to_path, *args, **kwargs): + if 'pedigrees' in to_path: + raise Exception('Something wrong while moving the file.') + def _test_mv_file_and_triggering_dag_exception(self, url, workspace, sample_data, genome_version, request_body, num_samples=None, sample_type='WES'): # Test saving ID file exception responses.calls.reset() self.mock_authorized_session.reset_mock() - self.mock_mv_file.side_effect = Exception('Something wrong while moving the file.') + self.mock_mv_file.side_effect = self._raise_move_file_error # Test triggering dag exception self.set_dag_trigger_error_response() From 9406cee16a3b1c0fcb6e60c187568ea5617fd2e2 Mon Sep 17 00:00:00 2001 From: Hana Snow Date: Tue, 9 Sep 2025 12:08:06 -0400 Subject: [PATCH 7/7] update hail dir env variable --- .../commands/check_for_new_samples_from_pipeline.py | 4 ++-- .../tests/check_for_new_samples_from_pipeline_tests.py | 2 +- seqr/management/tests/update_individuals_sample_qc_tests.py | 2 +- settings.py | 2 +- 4 files changed, 5 insertions(+), 5 deletions(-) diff --git a/seqr/management/commands/check_for_new_samples_from_pipeline.py b/seqr/management/commands/check_for_new_samples_from_pipeline.py index 899bb6a475..40033f03c7 100644 --- a/seqr/management/commands/check_for_new_samples_from_pipeline.py +++ b/seqr/management/commands/check_for_new_samples_from_pipeline.py @@ -22,7 +22,7 @@ from seqr.views.utils.permissions_utils import is_internal_anvil_project, project_has_anvil from seqr.views.utils.variant_utils import reset_cached_search_results, update_projects_saved_variant_json, \ get_saved_variants -from settings import SEQR_SLACK_LOADING_NOTIFICATION_CHANNEL, HAIL_SEARCH_DATA_DIR, ANVIL_UI_URL, \ +from settings import SEQR_SLACK_LOADING_NOTIFICATION_CHANNEL, PIPELINE_DATA_DIR, ANVIL_UI_URL, \ SEQR_SLACK_ANVIL_DATA_LOADING_CHANNEL logger = logging.getLogger(__name__) @@ -146,7 +146,7 @@ def _get_runs(cls, **kwargs): @staticmethod def _run_path(get_field_format): return RUN_FILE_PATH_TEMPLATE.format( - data_dir=HAIL_SEARCH_DATA_DIR, + data_dir=PIPELINE_DATA_DIR, **{field: get_field_format(field) for field in RUN_PATH_FIELDS} ) diff --git a/seqr/management/tests/check_for_new_samples_from_pipeline_tests.py b/seqr/management/tests/check_for_new_samples_from_pipeline_tests.py index d46c613e38..0700851413 100644 --- a/seqr/management/tests/check_for_new_samples_from_pipeline_tests.py +++ b/seqr/management/tests/check_for_new_samples_from_pipeline_tests.py @@ -369,7 +369,7 @@ def set_up(self): mock_rand_int = patcher.start() mock_rand_int.side_effect = [GUID_ID, GUID_ID, GUID_ID, GUID_ID, GCNV_GUID_ID, GCNV_GUID_ID, GCNV_GUID_ID, GCNV_GUID_ID, GUID_ID, GUID_ID, GUID_ID, GUID_ID] self.addCleanup(patcher.stop) - patcher = mock.patch('seqr.management.commands.check_for_new_samples_from_pipeline.HAIL_SEARCH_DATA_DIR') + patcher = mock.patch('seqr.management.commands.check_for_new_samples_from_pipeline.PIPELINE_DATA_DIR') mock_data_dir = patcher.start() mock_data_dir.__str__.return_value = self.MOCK_DATA_DIR self.addCleanup(patcher.stop) diff --git a/seqr/management/tests/update_individuals_sample_qc_tests.py b/seqr/management/tests/update_individuals_sample_qc_tests.py index fb406f0f7c..361b1b5698 100644 --- a/seqr/management/tests/update_individuals_sample_qc_tests.py +++ b/seqr/management/tests/update_individuals_sample_qc_tests.py @@ -74,7 +74,7 @@ class UpdateIndividualsSampleQC(TestCase): fixtures = ['users', '1kg_project'] def setUp(self): - patcher = mock.patch('seqr.management.commands.check_for_new_samples_from_pipeline.HAIL_SEARCH_DATA_DIR') + patcher = mock.patch('seqr.management.commands.check_for_new_samples_from_pipeline.PIPELINE_DATA_DIR') mock_data_dir = patcher.start() mock_data_dir.__str__.return_value = 'gs://seqr-hail-search-data/v3.1' self.addCleanup(patcher.stop) diff --git a/settings.py b/settings.py index 3762609701..891a8a030d 100644 --- a/settings.py +++ b/settings.py @@ -159,7 +159,7 @@ MEDIA_URL = '/media/' LOADING_DATASETS_DIR = os.environ.get('LOADING_DATASETS_DIR') -HAIL_SEARCH_DATA_DIR = os.environ.get('HAIL_SEARCH_DATA_DIR') +PIPELINE_DATA_DIR = os.environ.get('PIPELINE_DATA_DIR') LOGGING = { 'version': 1,