From da465d99c2385cd4d141c8529104b21aeb735187 Mon Sep 17 00:00:00 2001 From: vinay Date: Tue, 5 Aug 2025 21:05:26 +0100 Subject: [PATCH 1/5] support new ftp path layout --- .../metadata/api/adaptors/genome.py | 7 +- .../production/metadata/api/adaptors/vep.py | 12 +-- .../metadata/api/factories/utils.py | 37 +++++++++- src/tests/test_utils.py | 74 ++++++++++++------- 4 files changed, 89 insertions(+), 41 deletions(-) diff --git a/src/ensembl/production/metadata/api/adaptors/genome.py b/src/ensembl/production/metadata/api/adaptors/genome.py index a4d85d03..7cbcbc05 100644 --- a/src/ensembl/production/metadata/api/adaptors/genome.py +++ b/src/ensembl/production/metadata/api/adaptors/genome.py @@ -29,6 +29,8 @@ GenomeRelease, EnsemblRelease, EnsemblSite, AssemblySequence, GenomeDataset, Dataset, DatasetType, DatasetSource, \ ReleaseStatus, DatasetStatus, utils, DatasetAttribute, Attribute +from ensembl.production.metadata.api.factories.utils import format_accession_path + logger = logging.getLogger(__name__) @@ -929,11 +931,8 @@ def get_public_path(self, genome_uuid, dataset_type='all', release=None): dataset_type = 'regulation' match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) # Match format YYYY-MM last_geneset_update = match.group(1).replace('-', '_') - scientific_name = re.sub(r'[^a-zA-Z0-9]+', ' ', scientific_name) - scientific_name = scientific_name.replace(' ', '_') - scientific_name = re.sub(r'^_+|_+$', '', scientific_name) genebuild_source_name = genebuild_source_name.lower() - base_path = f"{scientific_name}/{accession}" + base_path = format_accession_path(accession) # Convert accession to path format common_path = f"{base_path}/{genebuild_source_name}" path_templates = { diff --git a/src/ensembl/production/metadata/api/adaptors/vep.py b/src/ensembl/production/metadata/api/adaptors/vep.py index 149c5a8b..a1f6226b 100644 --- a/src/ensembl/production/metadata/api/adaptors/vep.py +++ b/src/ensembl/production/metadata/api/adaptors/vep.py @@ -15,6 +15,7 @@ from sqlalchemy.orm import aliased from ensembl.production.metadata.api.adaptors.base import BaseAdaptor +from ensembl.production.metadata.api.factories.utils import format_accession_path from ensembl.production.metadata.api.models import Organism, Assembly, DatasetAttribute, Genome, GenomeDataset, Dataset @@ -74,18 +75,13 @@ def fetch_vep_locations(self, genome_uuid): elif not result.annotation_source or not result.last_geneset_update: raise ValueError(f"Missing annotation source or last geneset update for genome UUID: {genome_uuid}") - # Format the scientific name - scientific_name = result.scientific_name - scientific_name = re.sub(r"[^a-zA-Z0-9]+", " ", scientific_name) - scientific_name = re.sub(r" +", "_", scientific_name).strip("_") - # Format last geneset update last_geneset_update = re.sub(r"-", "_", result.last_geneset_update) - + base_path = format_accession_path(result.assembly_accession) # Construct the locations - faa_location = f"{scientific_name}/{result.assembly_accession}/vep/genome/softmasked.fa.bgz" + faa_location = f"{base_path}/vep/genome/softmasked.fa.bgz" gff_location = ( - f"{scientific_name}/{result.assembly_accession}/vep/" + f"{base_path}/vep/" f"{result.annotation_source}/geneset/{last_geneset_update}/genes.gff3.bgz" ) diff --git a/src/ensembl/production/metadata/api/factories/utils.py b/src/ensembl/production/metadata/api/factories/utils.py index 3ec58aeb..66cbc82b 100644 --- a/src/ensembl/production/metadata/api/factories/utils.py +++ b/src/ensembl/production/metadata/api/factories/utils.py @@ -10,8 +10,10 @@ # See the License for the specific language governing permissions and # limitations under the License. +import re +import os +from typing import Union from sqlalchemy.orm import aliased - from ensembl.production.metadata.api.models import Dataset, Genome, GenomeDataset, DatasetAttribute, Attribute, Assembly @@ -72,3 +74,36 @@ def get_genome_sets_by_assembly_and_provider(session): genome_sets_with_multiple = {key: genomes for key, genomes in genome_sets.items() if len(genomes) > 1} return genome_sets_with_multiple + + +def format_accession_path(accession: str) -> str: + """ + Converts assembly accession (e.g., 'GCF_043381705.1') to a structured path format + (e.g., 'GCF/043/381/705/1'). + + Parameters: + accession (str): The accession string to format. + + Returns: + str: Formatted path-like string. + + Raises: + ValueError: If the accession format is invalid. + """ + pattern = r'^(GCF|GCA)_(\d+)\.(\d+)$' + match = re.match(pattern, accession) + + if not match: + raise ValueError(f"Invalid accession format: '{accession}'. Expected format like '(GCF|GCA)_#########.#'.") + + prefix, number, version = match.groups() + + if len(number) > 9: + raise ValueError(f"Unexpected number length in accession: '{number}'. Expected up to 9 digits.") + + chunks = [number[i:i + 3] for i in range(0, 9, 3)] + + # Ensure we have exactly 3 chunks, filling with '000' if necessary + return os.path.join(prefix, *chunks, version) + + diff --git a/src/tests/test_utils.py b/src/tests/test_utils.py index caaa13ae..01548cb3 100644 --- a/src/tests/test_utils.py +++ b/src/tests/test_utils.py @@ -23,12 +23,13 @@ from ensembl.production.metadata.api.models import Genome, Dataset from ensembl.production.metadata.grpc import ensembl_metadata_pb2, utils +from ensembl.production.metadata.api.factories.utils import format_accession_path logger = logging.getLogger(__name__) @pytest.mark.parametrize("test_dbs", [[{'src': Path(__file__).parent / "databases/ensembl_genome_metadata"}, - {'src': Path(__file__).parent / "databases/ncbi_taxonomy"}]], + {'src': Path(__file__).parent / "databases/ncbi_taxonomy"}]], indirect=True) class TestUtils: dbc: UnitTestDB = None @@ -346,7 +347,8 @@ def test_get_genomes_by_uuid_null(self, genome_conn): ], indirect=['allow_unreleased'] ) - def test_get_brief_genome_details_by_uuid(self, genome_conn, allow_unreleased, genome_uuid, version, expected_count, actual): + def test_get_brief_genome_details_by_uuid(self, genome_conn, allow_unreleased, genome_uuid, version, expected_count, + actual): output = json_format.MessageToJson( utils.get_brief_genome_details_by_uuid( @@ -378,7 +380,8 @@ def test_get_brief_genome_details_by_uuid(self, genome_conn, allow_unreleased, g ], indirect=['allow_unreleased'] ) - def test_get_attributes_by_genome_uuid(self, genome_conn, allow_unreleased, genome_uuid, assembly_level, genebuild_sample_gene, version): + def test_get_attributes_by_genome_uuid(self, genome_conn, allow_unreleased, genome_uuid, assembly_level, + genebuild_sample_gene, version): output = json_format.MessageToJson( utils.get_attributes_by_genome_uuid( db_conn=genome_conn, @@ -411,7 +414,6 @@ def test_get_attributes_by_genome_uuid(self, genome_conn, allow_unreleased, geno assert [key in output['attributesInfo'] for key in expected_attributes_info_keys] assert output['attributesInfo']['assemblyLevel'] == assembly_level - def test_get_genomes_by_keyword(self, genome_conn): output = [ json.loads(json_format.MessageToJson(response)) for response in @@ -539,11 +541,11 @@ def test_get_genomes_by_name(self, genome_conn): }, 'relatedAssembliesCount': 5 } - + t = json.loads(output) - + assert t["attributesInfo"] == expected_output["attributesInfo"] - + assert json.loads(output) == expected_output def test_get_genomes_by_name_release_unspecified(self, genome_conn): @@ -661,19 +663,19 @@ def test_get_genome_uuid_by_tag(self, genome_conn, genome_tag, expected_output): # valid genome uuid and no dataset should return all the datasets links of that genome uuid ("a733574a-93e7-11ec-a39d-005056b38ce3", 'all', None, { "Links": [ - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/genome", - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/regulation", - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/variation/test_version", - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/homology/test_version", - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/genebuild/test_version" + "GCA/000/146/045/2/test_anno_source/genome", + "GCA/000/146/045/2/test_anno_source/regulation", + "GCA/000/146/045/2/test_anno_source/variation/test_version", + "GCA/000/146/045/2/test_anno_source/homology/test_version", + "GCA/000/146/045/2/test_anno_source/genebuild/test_version" ] }), # valid genome uuid and a valid dataset should return corresponding dataset link ("a733574a-93e7-11ec-a39d-005056b38ce3", 'assembly', None, { "Links": [ - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/genome", - "Saccharomyces_cerevisiae_S288c/GCA_000146045.2/test_anno_source/genebuild/test_version" + "GCA/000/146/045/2/test_anno_source/genome", + "GCA/000/146/045/2/test_anno_source/genebuild/test_version" ] }), @@ -739,23 +741,23 @@ def test_get_release_version_by_uuid(self, genome_conn, genome_uuid, dataset_typ "genome_uuid, expected_output", [ ( - # Human - "65d4f21f-695a-4ed0-be67-5732a551fea4", - { - "faaLocation": "Homo_sapiens/GCA_018473295.1/vep/genome/softmasked.fa.bgz", - "gffLocation": "Homo_sapiens/GCA_018473295.1/vep/ensembl/geneset/2022_08/genes.gff3.bgz" - } + # Human + "65d4f21f-695a-4ed0-be67-5732a551fea4", + { + "faaLocation": "GCA/018/473/295/1/vep/genome/softmasked.fa.bgz", + "gffLocation": "GCA/018/473/295/1/vep/ensembl/geneset/2022_08/genes.gff3.bgz" + } ), ( - # Ecoli - "a73351f7-93e7-11ec-a39d-005056b38ce3", - { - 'faaLocation': 'Escherichia_coli_str_K_12_substr_MG1655_str_K12/GCA_000005845.2/vep/genome/softmasked.fa.bgz', - 'gffLocation': 'Escherichia_coli_str_K_12_substr_MG1655_str_K12/GCA_000005845.2/vep/community/geneset/2018_09/genes.gff3.bgz' - } + # Ecoli + "a73351f7-93e7-11ec-a39d-005056b38ce3", + { + 'faaLocation': 'GCA/000/005/845/2/vep/genome/softmasked.fa.bgz', + 'gffLocation': 'GCA/000/005/845/2/vep/community/geneset/2018_09/genes.gff3.bgz' + } ), ( - "some-invalid-genome-uuid-000000000000", {} + "some-invalid-genome-uuid-000000000000", {} ) ] ) @@ -767,4 +769,20 @@ def test_get_vep_paths_by_uuid(self, vep_conn, genome_uuid, expected_output): ) ) print(output) - assert json.loads(output) == expected_output \ No newline at end of file + assert json.loads(output) == expected_output + + def test_format_accession_path_valid(self): + assert format_accession_path("GCF_043381705.1") == "GCF/043/381/705/1" + assert format_accession_path("GCA_123456789.2") == "GCA/123/456/789/2" + + def test_format_accession_path_invalid_prefix(self): + with pytest.raises(ValueError): + format_accession_path("XYZ_123456789.1") + + def test_format_accession_path_invalid_format(self): + with pytest.raises(ValueError): + format_accession_path("GCF123456789.1") + + def test_format_accession_path_long_number(self): + with pytest.raises(ValueError): + format_accession_path("GCF_1234567890.1") From 3d0da02bbbfcbc793a60bee76b28107bd4d9592d Mon Sep 17 00:00:00 2001 From: vinay Date: Wed, 6 Aug 2025 15:40:15 +0100 Subject: [PATCH 2/5] support new ftp path layout: Update test cases --- src/tests/test_api.py | 18 +++++++++--------- src/tests/test_protobuf_msg_factory.py | 8 ++++---- 2 files changed, 13 insertions(+), 13 deletions(-) diff --git a/src/tests/test_api.py b/src/tests/test_api.py index 38bb4c9b..4a98de5a 100644 --- a/src/tests/test_api.py +++ b/src/tests/test_api.py @@ -33,13 +33,13 @@ def test_get_public_path(self, test_dbs): assert len(paths) == 4 # assert all("/genebuild/" in path for path in paths) path = genome_adapter.get_public_path(genome_uuid, dataset_type='genebuild') - assert path[0]['path'] == 'Saccharomyces_cerevisiae_S288c/GCA_000146045.2/community/geneset/2018_10' + assert path[0]['path'] == 'GCA/000/146/045/2/community/geneset/2018_10' path = genome_adapter.get_public_path(genome_uuid, dataset_type='assembly') - assert path[0]['path'] == 'Saccharomyces_cerevisiae_S288c/GCA_000146045.2/genome' + assert path[0]['path'] == 'GCA/000/146/045/2/genome' path = genome_adapter.get_public_path(genome_uuid, dataset_type='variation') - assert path[0]['path'] == 'Saccharomyces_cerevisiae_S288c/GCA_000146045.2/community/variation/2018_10' + assert path[0]['path'] == 'GCA/000/146/045/2/community/variation/2018_10' path = genome_adapter.get_public_path(genome_uuid, dataset_type='homologies') - assert path[0]['path'] == 'Saccharomyces_cerevisiae_S288c/GCA_000146045.2/community/homology/2018_10' + assert path[0]['path'] == 'GCA/000/146/045/2/community/homology/2018_10' with pytest.raises(TypeNotFoundException): genome_adapter.get_public_path(genome_uuid, dataset_type='regulatory_features') # assert path[0]['path'] == 'Saccharomyces_cerevisiae_S288c/GCA_000146045.2/ensembl/regulation' @@ -52,12 +52,12 @@ def test_default_public_path(self, test_dbs): assert len(paths) == 5 # assert all("/genebuild/" in path for path in paths) path = genome_adapter.get_public_path(genome_uuid, dataset_type='genebuild') - assert path[0]['path'] == 'Homo_sapiens/GCA_000001405.29/ensembl/geneset/2023_03' + assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/geneset/2023_03' path = genome_adapter.get_public_path(genome_uuid, dataset_type='assembly') - assert path[0]['path'] == 'Homo_sapiens/GCA_000001405.29/genome' + assert path[0]['path'] == 'GCA/000/001/405/29/genome' path = genome_adapter.get_public_path(genome_uuid, dataset_type='variation') - assert path[0]['path'] == 'Homo_sapiens/GCA_000001405.29/ensembl/variation/2023_03' + assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/variation/2023_03' path = genome_adapter.get_public_path(genome_uuid, dataset_type='homologies') - assert path[0]['path'] == 'Homo_sapiens/GCA_000001405.29/ensembl/homology/2023_03' + assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/homology/2023_03' path = genome_adapter.get_public_path(genome_uuid, dataset_type='regulatory_features') - assert path[0]['path'] == 'Homo_sapiens/GCA_000001405.29/ensembl/regulation' + assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/regulation' diff --git a/src/tests/test_protobuf_msg_factory.py b/src/tests/test_protobuf_msg_factory.py index 1244b6da..a1651258 100644 --- a/src/tests/test_protobuf_msg_factory.py +++ b/src/tests/test_protobuf_msg_factory.py @@ -296,16 +296,16 @@ def test_create_genome_uuid(self, genome_conn, genome_tag, current_only, expecte # Human "65d4f21f-695a-4ed0-be67-5732a551fea4", { - "faaLocation": "Homo_sapiens/GCA_018473295.1/vep/genome/softmasked.fa.bgz", - "gffLocation": "Homo_sapiens/GCA_018473295.1/vep/ensembl/geneset/2022_08/genes.gff3.bgz" + "faaLocation": "GCA/018/473/295/1/vep/genome/softmasked.fa.bgz", + "gffLocation": "GCA/018/473/295/1/vep/ensembl/geneset/2022_08/genes.gff3.bgz" } ), ( # Ecoli "a73351f7-93e7-11ec-a39d-005056b38ce3", { - 'faaLocation': 'Escherichia_coli_str_K_12_substr_MG1655_str_K12/GCA_000005845.2/vep/genome/softmasked.fa.bgz', - 'gffLocation': 'Escherichia_coli_str_K_12_substr_MG1655_str_K12/GCA_000005845.2/vep/community/geneset/2018_09/genes.gff3.bgz' + 'faaLocation': 'GCA/000/005/845/2/vep/genome/softmasked.fa.bgz', + 'gffLocation': 'GCA/000/005/845/2/vep/community/geneset/2018_09/genes.gff3.bgz' } ) ] From 4bcca541a6e46a4172a7f08cff877695abfc9896 Mon Sep 17 00:00:00 2001 From: danielp Date: Thu, 25 Sep 2025 15:25:38 +0100 Subject: [PATCH 3/5] Major update for new path structure: -Updated get_paths to new path structure -Defaults to is_current but accepts version (integrated or partial) -Added input validation -VEP code was simplified to call on the get_paths -tests were updated and new tests created -test data: ensembl_release.label was made to look like real data and associated tests were updated --- .../metadata/api/adaptors/genome.py | 166 +++++++++++++++--- .../production/metadata/api/adaptors/vep.py | 85 ++------- .../ensembl_genome_metadata/dataset.txt | 2 + .../ensembl_release.txt | 12 +- .../genome_dataset.txt | 9 +- .../genome_release.txt | 2 +- src/tests/test_api.py | 46 +++-- src/tests/test_grpc_genome.py | 4 +- src/tests/test_grpc_release.py | 14 +- src/tests/test_protobuf_msg_factory.py | 36 ++-- src/tests/test_utils.py | 16 +- 11 files changed, 234 insertions(+), 158 deletions(-) diff --git a/src/ensembl/production/metadata/api/adaptors/genome.py b/src/ensembl/production/metadata/api/adaptors/genome.py index 7cbcbc05..014c55f7 100644 --- a/src/ensembl/production/metadata/api/adaptors/genome.py +++ b/src/ensembl/production/metadata/api/adaptors/genome.py @@ -25,12 +25,11 @@ from ensembl.production.metadata.api.adaptors.base import BaseAdaptor, check_parameter, cfg from ensembl.production.metadata.api.exceptions import TypeNotFoundException +from ensembl.production.metadata.api.factories.utils import format_accession_path from ensembl.production.metadata.api.models import Genome, Organism, Assembly, OrganismGroup, OrganismGroupMember, \ GenomeRelease, EnsemblRelease, EnsemblSite, AssemblySequence, GenomeDataset, Dataset, DatasetType, DatasetSource, \ ReleaseStatus, DatasetStatus, utils, DatasetAttribute, Attribute -from ensembl.production.metadata.api.factories.utils import format_accession_path - logger = logging.getLogger(__name__) @@ -56,12 +55,16 @@ class GenomeDatasetsListItem(NamedTuple): class GenomeAdaptor(BaseAdaptor): - def __init__(self, metadata_uri: str, taxonomy_uri: str): + def __init__(self, metadata_uri: str, taxonomy_uri=None): super().__init__(metadata_uri) - self.taxonomy_db = DBConnection(taxonomy_uri, pool_size=cfg.pool_size, pool_recycle=cfg.pool_recycle) + if taxonomy_uri is None: + self.taxonomy_db = None + else: + self.taxonomy_db = DBConnection(taxonomy_uri, pool_size=cfg.pool_size, pool_recycle=cfg.pool_recycle) def fetch_taxonomy_names(self, taxonomy_ids, synonyms=None): - + if self.taxonomy_db is None: + raise ValueError("Taxonomy DB not provided") if synonyms is None: synonyms = [] taxonomy_ids = check_parameter(taxonomy_ids) @@ -94,6 +97,8 @@ def fetch_taxonomy_names(self, taxonomy_ids, synonyms=None): return taxons def fetch_taxonomy_ids(self, taxonomy_names): + if self.taxonomy_db is None: + raise ValueError("Taxonomy DB not provided") taxids = [] taxonomy_names = check_parameter(taxonomy_names) for taxon in taxonomy_names: @@ -878,13 +883,44 @@ def fetch_assemblies_count(self, species_taxonomy_id: int, release_version: floa return session.execute(query).scalar() def get_public_path(self, genome_uuid, dataset_type='all', release=None): + """ + Retrieve public file paths for genomic datasets. + + Args: + genome_uuid (str): Unique identifier for the genome + dataset_type (str): Type of dataset ('genebuild', 'assembly', 'homologies', 'variation', or 'all') + release (str, optional): Specific Ensembl release label. If None, uses current release. + + Returns: + list: List of dictionaries containing dataset_type and path information + + Raises: + ValueError: If genome_uuid or release not found, or required metadata missing + TypeNotFoundException: If requested dataset_type is not available + """ paths = [] scientific_name = None accession = None genebuild_source_name = None last_geneset_update = None + with self.metadata_db.session_scope() as session: - # Single combined query to get all required data + # === VALIDATION SECTION === + genome_exists = session.execute( + select(Genome.genome_uuid).where(Genome.genome_uuid == genome_uuid) + ).first() + if not genome_exists: + raise ValueError(f"Genome with UUID {genome_uuid} not found") + + if release is not None: + release_exists = session.execute( + select(EnsemblRelease.label).where(EnsemblRelease.label == release) + ).first() + if not release_exists: + raise ValueError(f"Ensembl release with label '{release}' not found") + + # === METADATA RETRIEVAL === + # Get core genome metadata: organism name, assembly accession, and genebuild info query = select( Organism.scientific_name, Assembly.accession, @@ -914,54 +950,130 @@ def get_public_path(self, genome_uuid, dataset_type='all', release=None): else: scientific_name = accession = genebuild_source_name = last_geneset_update = None - # Query for which of the 5 supported dataset types exist for this genome - supported_types = ['genebuild', 'assembly', 'homologies', 'regulatory_features', 'variation'] + # === DATASET TYPE DISCOVERY === + supported_types = ['genebuild', 'assembly', 'homologies', 'variation'] unique_dataset_types_query = select(DatasetType.name).distinct().join( Dataset ).join(GenomeDataset).join(Genome).where( Genome.genome_uuid == genome_uuid, - DatasetType.name.in_(supported_types) + DatasetType.name.in_(supported_types), + Dataset.status == DatasetStatus.RELEASED ) unique_dataset_types = session.execute(unique_dataset_types_query).scalars().all() + variation_release = None + homology_release = None + + # === RELEASE HANDLING === + if release is None: + if 'variation' in unique_dataset_types: + variation_release = session.execute( + select(EnsemblRelease.label).join(GenomeDataset).join(Dataset).join(DatasetType).join( + Genome).where( + Genome.genome_uuid == genome_uuid, + DatasetType.name == 'variation', + Dataset.status == DatasetStatus.RELEASED, + GenomeDataset.is_current == True, + EnsemblRelease.release_type == 'partial' + ) + ).scalar_one_or_none() + + if 'homologies' in unique_dataset_types: + homology_release = session.execute( + select(EnsemblRelease.label).join(GenomeDataset).join(Dataset).join(DatasetType).join( + Genome).where( + Genome.genome_uuid == genome_uuid, + DatasetType.name == 'homologies', + Dataset.status == DatasetStatus.RELEASED, + GenomeDataset.is_current == True, + EnsemblRelease.release_type == 'partial' + ) + ).scalar_one_or_none() + else: + if 'variation' in unique_dataset_types: + variation_release = session.execute( + select(EnsemblRelease.label).join(GenomeDataset).join(Dataset).join(DatasetType).join( + Genome).where( + Genome.genome_uuid == genome_uuid, + DatasetType.name == 'variation', + Dataset.status == DatasetStatus.RELEASED, + EnsemblRelease.release_type == 'partial', + EnsemblRelease.release_id <= ( + select(EnsemblRelease.release_id).where( + EnsemblRelease.label == release).scalar_subquery() + ) + ).order_by(EnsemblRelease.release_id.desc()) + ).scalar_one_or_none() + + if variation_release is None: + unique_dataset_types.remove('variation') + + if 'homologies' in unique_dataset_types: + homology_release = session.execute( + select(EnsemblRelease.label).join(GenomeDataset).join(Dataset).join(DatasetType).join( + Genome).where( + Genome.genome_uuid == genome_uuid, + DatasetType.name == 'homologies', + Dataset.status == DatasetStatus.RELEASED, + EnsemblRelease.release_type == 'partial', + EnsemblRelease.release_id <= ( + select(EnsemblRelease.release_id).where( + EnsemblRelease.label == release).scalar_subquery() + ) + ).order_by(EnsemblRelease.release_id.desc()) + ).scalar_one_or_none() + + if homology_release is None: + unique_dataset_types.remove('homologies') + + # === DATA VALIDATION === + if 'genebuild' not in unique_dataset_types or 'assembly' not in unique_dataset_types: + raise ValueError( + f"Missing genebuild or assembly dataset types. Something is seriously wrong with {genome_uuid}") + if scientific_name is None or accession is None or genebuild_source_name is None or last_geneset_update is None: raise ValueError("Required metadata fields are missing. Please check the database entries.") - unique_dataset_types = ['regulation' if t == 'regulatory_features' else t for t in unique_dataset_types] - if dataset_type == 'regulatory_features': - dataset_type = 'regulation' - match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) # Match format YYYY-MM + + # === PATH CONSTRUCTION === + match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) + if not match: + raise ValueError(f"Invalid last_geneset_update format: {last_geneset_update}") last_geneset_update = match.group(1).replace('-', '_') + genebuild_source_name = genebuild_source_name.lower() - base_path = format_accession_path(accession) # Convert accession to path format - common_path = f"{base_path}/{genebuild_source_name}" + + base_path = format_accession_path(accession) + common_path = f"{base_path}/{genebuild_source_name}/{last_geneset_update}" + + if 'homologies' in unique_dataset_types and homology_release: + homology_release = homology_release.replace('-', '_') + if 'variation' in unique_dataset_types and variation_release: + variation_release = variation_release.replace('-', '_') path_templates = { - 'genebuild': f"{common_path}/geneset/{last_geneset_update}", + 'genebuild': f"{common_path}/geneset", 'assembly': f"{base_path}/genome", - 'homologies': f"{common_path}/homology/{last_geneset_update}", - 'regulation': f"{common_path}/regulation", - 'variation': f"{common_path}/variation/{last_geneset_update}", + 'homologies': f"{common_path}/homology/{homology_release}", + 'variation': f"{common_path}/variation/{variation_release}", } - # Check for invalid dataset type early + # === REQUEST VALIDATION === if dataset_type not in unique_dataset_types and dataset_type != 'all': raise TypeNotFoundException(f"Dataset Type : {dataset_type} not found in metadata.") - # If 'all', add paths for all unique dataset types + # === PATH GENERATION === if dataset_type == 'all': - for t in unique_dataset_types: + for dataset_type_name in unique_dataset_types: paths.append({ - "dataset_type": t, - "path": path_templates[t] + "dataset_type": dataset_type_name, + "path": path_templates[dataset_type_name] }) elif dataset_type in path_templates: - # Add path for the specific dataset type paths.append({ "dataset_type": dataset_type, "path": path_templates[dataset_type] }) else: - # If the code reaches here, it means there is a logic error raise TypeNotFoundException(f"Dataset Type : {dataset_type} has no associated path.") - return paths + return paths \ No newline at end of file diff --git a/src/ensembl/production/metadata/api/adaptors/vep.py b/src/ensembl/production/metadata/api/adaptors/vep.py index a1f6226b..e929654f 100644 --- a/src/ensembl/production/metadata/api/adaptors/vep.py +++ b/src/ensembl/production/metadata/api/adaptors/vep.py @@ -9,19 +9,14 @@ # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. -import re - -from sqlalchemy import and_ -from sqlalchemy.orm import aliased from ensembl.production.metadata.api.adaptors.base import BaseAdaptor -from ensembl.production.metadata.api.factories.utils import format_accession_path -from ensembl.production.metadata.api.models import Organism, Assembly, DatasetAttribute, Genome, GenomeDataset, Dataset - +from ensembl.production.metadata.api.adaptors.genome import GenomeAdaptor class VepAdaptor(BaseAdaptor): def __init__(self, metadata_uri: str, file="all"): super().__init__(metadata_uri) + self.metadata_uri = metadata_uri self.file = file def fetch_vep_locations(self, genome_uuid): @@ -31,64 +26,18 @@ def fetch_vep_locations(self, genome_uuid): :param genome_uuid: The UUID of the genome to fetch locations for. :return: A dictionary containing the FAA and GFF locations or a specific location string if 'file' is set. """ - with self.metadata_db.session_scope() as session: - # Aliases for clarity and distinct filtering - annotation_source_attr = aliased(DatasetAttribute) - last_geneset_update_attr = aliased(DatasetAttribute) - - query = ( - session.query( - Organism.scientific_name, - Assembly.accession.label("assembly_accession"), - annotation_source_attr.value.label("annotation_source"), - last_geneset_update_attr.value.label("last_geneset_update"), - ) - .join(Genome, Genome.organism_id == Organism.organism_id) - .join(Assembly, Assembly.assembly_id == Genome.assembly_id) - .join(GenomeDataset, GenomeDataset.genome_id == Genome.genome_id) - .join(Dataset, Dataset.dataset_id == GenomeDataset.dataset_id) - .outerjoin( - annotation_source_attr, - and_( - annotation_source_attr.dataset_id == Dataset.dataset_id, - annotation_source_attr.attribute.has(name="genebuild.annotation_source"), - ), - ) - .outerjoin( - last_geneset_update_attr, - and_( - last_geneset_update_attr.dataset_id == Dataset.dataset_id, - last_geneset_update_attr.attribute.has(name="genebuild.last_geneset_update"), - ), - ) - .filter( - Genome.genome_uuid == genome_uuid, - Dataset.name == "genebuild", - ) - .distinct() - ) - - result = query.one_or_none() - - if not result: - raise ValueError(f"No data found for genome UUID: {genome_uuid}") - elif not result.annotation_source or not result.last_geneset_update: - raise ValueError(f"Missing annotation source or last geneset update for genome UUID: {genome_uuid}") - - # Format last geneset update - last_geneset_update = re.sub(r"-", "_", result.last_geneset_update) - base_path = format_accession_path(result.assembly_accession) - # Construct the locations - faa_location = f"{base_path}/vep/genome/softmasked.fa.bgz" - gff_location = ( - f"{base_path}/vep/" - f"{result.annotation_source}/geneset/{last_geneset_update}/genes.gff3.bgz" - ) - - # Return based on the `file` argument - if self.file == "faa_location": - return faa_location - elif self.file == "gff_location": - return gff_location - else: - return {"faa_location": faa_location, "gff_location": gff_location} + genome_adaptor = GenomeAdaptor(self.metadata_uri) + genebuild_path = genome_adaptor.get_public_path(genome_uuid, dataset_type='genebuild') + genebuild_path = genebuild_path[0]['path'] + assembly_path = genome_adaptor.get_public_path(genome_uuid, dataset_type='assembly') + assembly_path = assembly_path[0]['path'] + + faa_location = f"{assembly_path}/unmasked.fa.bgz" + gff_location = f"{genebuild_path}/genes.gff3.bgz" + # Return based on the `file` argument + if self.file == "faa_location": + return faa_location + elif self.file == "gff_location": + return gff_location + else: + return {"faa_location": faa_location, "gff_location": gff_location} diff --git a/src/tests/databases/ensembl_genome_metadata/dataset.txt b/src/tests/databases/ensembl_genome_metadata/dataset.txt index 77181619..1ca8c7d5 100644 --- a/src/tests/databases/ensembl_genome_metadata/dataset.txt +++ b/src/tests/databases/ensembl_genome_metadata/dataset.txt @@ -497,3 +497,5 @@ 9072 99999999-847e-4742-a68b-18c3ece068aa genebuild ENS01 2023-09-22 15:03:02.000000 GCA_021950905.1_ENS01 18 2 Submitted \N 9073 99999999-da2c-4997-8002-9da717ba79d2 genebuild_compute ENS01 2024-04-24 16:07:22.000000 From 2ef7c056-847e-4742-a68b-18c3ece068aa 18 8 Submitted 9072 9074 99999999-d9e0-4eca-9a49-7a6d9e311c8d xrefs ENS01 2024-04-24 16:07:22.000000 From 7573b939-da2c-4997-8002-9da717ba79d2 18 13 Submitted 9073 +9075 0a0bed83-72x7-4x8a-axcb-97450ef82495 variation 2.0 2025-11-09 12:50:31.531084 R64-1-1 644 3 Released \N + diff --git a/src/tests/databases/ensembl_genome_metadata/ensembl_release.txt b/src/tests/databases/ensembl_genome_metadata/ensembl_release.txt index 769d578f..67107deb 100644 --- a/src/tests/databases/ensembl_genome_metadata/ensembl_release.txt +++ b/src/tests/databases/ensembl_genome_metadata/ensembl_release.txt @@ -1,6 +1,6 @@ -1 110.1 2023-10-18 MVP Beta-1 1 partial 1 Released 1 -2 110.2 \N MVP Beta-2 0 partial 1 Prepared 2 -3 110.3 \N MVP Beta-3 0 partial 1 Preparing 3 -4 112.0 \N MVP Rel-1 0 partial 1 Planned 4 -5 108.0 2023-06-15 First Beta 0 partial 1 Released 5 -6 114.0 2025-06-15 dataset_test 0 partial 1 Preparing 6 +1 110.1 2020-10-18 2020-10-18 1 partial 1 Released 1 +2 110.2 2021-10-18 2021-10-18 0 partial 1 Prepared 2 +3 110.3 2022-10-18 2022-10-18 0 partial 1 Preparing 3 +4 112.0 2022-11-18 2022-11-18 0 partial 1 Planned 4 +5 108.0 2023-06-15 2023-06-15 0 partial 1 Released 5 +6 114.0 2025-06-15 2025-06-15 0 partial 1 Preparing 6 diff --git a/src/tests/databases/ensembl_genome_metadata/genome_dataset.txt b/src/tests/databases/ensembl_genome_metadata/genome_dataset.txt index f7541c3f..61e11315 100644 --- a/src/tests/databases/ensembl_genome_metadata/genome_dataset.txt +++ b/src/tests/databases/ensembl_genome_metadata/genome_dataset.txt @@ -27,8 +27,8 @@ 338 0 338 169 \N 347 1 347 174 1 348 1 348 174 1 -401 1 401 201 5 -402 1 402 201 5 +401 1 401 201 1 +402 1 402 201 1 405 1 405 203 5 406 1 406 203 5 887 1 888 4 1 @@ -51,7 +51,7 @@ 1437 1 1496 89 3 1448 1 1507 5 1 1450 1 1509 174 1 -1469 1 1528 201 5 +1469 0 1528 201 1 1478 1 1537 12 5 1485 1 1544 74 5 2217 1 2276 4 1 @@ -496,4 +496,5 @@ 9016 0 9071 86 \N 9017 0 9072 204 6 9018 0 9073 204 6 -9019 0 9074 204 6 \ No newline at end of file +9019 0 9074 204 6 +9020 1 9075 201 5 diff --git a/src/tests/databases/ensembl_genome_metadata/genome_release.txt b/src/tests/databases/ensembl_genome_metadata/genome_release.txt index 217aa3a1..136094ea 100644 --- a/src/tests/databases/ensembl_genome_metadata/genome_release.txt +++ b/src/tests/databases/ensembl_genome_metadata/genome_release.txt @@ -22,7 +22,7 @@ 22 0 5 2 23 0 74 2 24 0 4 2 -25 0 201 2 +25 0 201 1 26 0 1 2 27 0 19 4 28 0 92 4 diff --git a/src/tests/test_api.py b/src/tests/test_api.py index 4a98de5a..77e23d8d 100644 --- a/src/tests/test_api.py +++ b/src/tests/test_api.py @@ -26,6 +26,7 @@ class TestApi: dbc = None # type: UnitTestDB + # Basic test to show full funtionality getting the most recent result for everything. def test_get_public_path(self, test_dbs): genome_adapter = GenomeAdaptor(test_dbs['ensembl_genome_metadata'].dbc.url, test_dbs['ncbi_taxonomy'].dbc.url) genome_uuid = 'a733574a-93e7-11ec-a39d-005056b38ce3' @@ -33,31 +34,42 @@ def test_get_public_path(self, test_dbs): assert len(paths) == 4 # assert all("/genebuild/" in path for path in paths) path = genome_adapter.get_public_path(genome_uuid, dataset_type='genebuild') - assert path[0]['path'] == 'GCA/000/146/045/2/community/geneset/2018_10' + assert path[0]['path'] == 'GCA/000/146/045/2/community/2018_10/geneset' path = genome_adapter.get_public_path(genome_uuid, dataset_type='assembly') assert path[0]['path'] == 'GCA/000/146/045/2/genome' path = genome_adapter.get_public_path(genome_uuid, dataset_type='variation') - assert path[0]['path'] == 'GCA/000/146/045/2/community/variation/2018_10' + assert path[0]['path'] == 'GCA/000/146/045/2/community/2018_10/variation/2023_06_15' path = genome_adapter.get_public_path(genome_uuid, dataset_type='homologies') - assert path[0]['path'] == 'GCA/000/146/045/2/community/homology/2018_10' + assert path[0]['path'] == 'GCA/000/146/045/2/community/2018_10/homology/2023_06_15' + + # specific release: + def test_public_path_release(self, test_dbs): + genome_adapter = GenomeAdaptor(test_dbs['ensembl_genome_metadata'].dbc.url, test_dbs['ncbi_taxonomy'].dbc.url) + genome_uuid = 'a733574a-93e7-11ec-a39d-005056b38ce3' + paths = genome_adapter.get_public_path(genome_uuid, dataset_type='all', release='2021-10-18') + assert len(paths) == 3 + # assert all("/genebuild/" in path for path in paths) + path = genome_adapter.get_public_path(genome_uuid, dataset_type='genebuild', release='2021-10-18') + assert path[0]['path'] == 'GCA/000/146/045/2/community/2018_10/geneset' + path = genome_adapter.get_public_path(genome_uuid, dataset_type='assembly', release='2021-10-18') + assert path[0]['path'] == 'GCA/000/146/045/2/genome' + path = genome_adapter.get_public_path(genome_uuid, dataset_type='variation', release='2021-10-18') + assert path[0]['path'] == 'GCA/000/146/045/2/community/2018_10/variation/2020_10_18' with pytest.raises(TypeNotFoundException): - genome_adapter.get_public_path(genome_uuid, dataset_type='regulatory_features') - # assert path[0]['path'] == 'Saccharomyces_cerevisiae_S288c/GCA_000146045.2/ensembl/regulation' + path = genome_adapter.get_public_path(genome_uuid, dataset_type='homologies', release='2021-10-18') - def test_default_public_path(self, test_dbs): + # Basic test to show full funtionality for a genome with limited datasets. + def test_get_public_path_limited(self, test_dbs): genome_adapter = GenomeAdaptor(test_dbs['ensembl_genome_metadata'].dbc.url, test_dbs['ncbi_taxonomy'].dbc.url) - genome_uuid = 'a7335667-93e7-11ec-a39d-005056b38ce3' - # Homo sapien GRCH38 + genome_uuid = 'a733550b-93e7-11ec-a39d-005056b38ce3' + # c elegans genome_id=203 paths = genome_adapter.get_public_path(genome_uuid, dataset_type='all') - assert len(paths) == 5 - # assert all("/genebuild/" in path for path in paths) + assert len(paths) == 3 path = genome_adapter.get_public_path(genome_uuid, dataset_type='genebuild') - assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/geneset/2023_03' + assert path[0]['path'] == 'GCA/000/002/985/3/wormbase/2014_10/geneset' path = genome_adapter.get_public_path(genome_uuid, dataset_type='assembly') - assert path[0]['path'] == 'GCA/000/001/405/29/genome' - path = genome_adapter.get_public_path(genome_uuid, dataset_type='variation') - assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/variation/2023_03' + assert path[0]['path'] == 'GCA/000/002/985/3/genome' path = genome_adapter.get_public_path(genome_uuid, dataset_type='homologies') - assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/homology/2023_03' - path = genome_adapter.get_public_path(genome_uuid, dataset_type='regulatory_features') - assert path[0]['path'] == 'GCA/000/001/405/29/ensembl/regulation' + assert path[0]['path'] == 'GCA/000/002/985/3/wormbase/2014_10/homology/2023_06_15' + with pytest.raises(TypeNotFoundException): + genome_adapter.get_public_path(genome_uuid, dataset_type='variation') diff --git a/src/tests/test_grpc_genome.py b/src/tests/test_grpc_genome.py index 9431e4fa..4b0bf391 100644 --- a/src/tests/test_grpc_genome.py +++ b/src/tests/test_grpc_genome.py @@ -243,13 +243,13 @@ def test_fetch_sequences_by_assembly_seq_name(self, genome_conn, genome_uuid, as [ # nothing specified + allow_unreleased -> fetches everything (None, None, True, False, "6c1896f9-10dd-423e-a1ff-db8b5815cb66", 30), - (None, None, False, False, "6c1896f9-10dd-423e-a1ff-db8b5815cb66", 10), + (None, None, False, False, "6c1896f9-10dd-423e-a1ff-db8b5815cb66", 11), ("8364a820-5485-42d7-a648-1a5eeb858319", None, True, False, "3c67123a-e9e1-41ef-9014-2aadc8acf12a", 1), # specifying genome_uuid -- Triticum aestivum (SAMEA4791365) ("a73357ab-93e7-11ec-a39d-005056b38ce3", None, True, False, "999315f6-6d25-481f-a017-297f7e1490c8", 2), ("a73357ab-93e7-11ec-a39d-005056b38ce3", None, True, True, "999315f6-6d25-481f-a017-297f7e1490c8", 1), # fetch unreleased datasets only - (None, None, False, True, "6c1896f9-10dd-423e-a1ff-db8b5815cb66", 10), + (None, None, False, True, "6c1896f9-10dd-423e-a1ff-db8b5815cb66", 11), (None, 'f93d21ca-9a24-4c31-ae11-b0f8d3deab6d', True, True, "3c67123a-e9e1-41ef-9014-2aadc8acf12a", 1), ], indirect=['allow_unreleased'] diff --git a/src/tests/test_grpc_release.py b/src/tests/test_grpc_release.py index 74a8660d..f8432387 100644 --- a/src/tests/test_grpc_release.py +++ b/src/tests/test_grpc_release.py @@ -60,13 +60,13 @@ def test_fetch_all_releases(self, release_conn, allow_unreleased, expected_count logger.debug("Results: %s", releases) assert len(releases) == expected_count assert [release.EnsemblSite.name == 'Ensembl' for release in releases] - assert releases[1].EnsemblRelease.label == 'MVP Beta-1' + assert releases[1].EnsemblRelease.label == '2020-10-18' @pytest.mark.parametrize( "allow_unreleased, genome_uuid, release_name", [ - (False, 'a73351f7-93e7-11ec-a39d-005056b38ce3', 'First Beta'), - (True, '75b7ac15-6373-4ad5-9fb7-23813a5355a4', 'MVP Beta-2') + (False, 'a73351f7-93e7-11ec-a39d-005056b38ce3', '2021-10-18'), + (True, '75b7ac15-6373-4ad5-9fb7-23813a5355a4', '2023-06-15') ], indirect=['allow_unreleased'] ) @@ -78,15 +78,15 @@ def test_fetch_releases_for_genome(self, release_conn, allow_unreleased, genome_ assert len(releases) == 1 logger.info(releases) assert releases[0].EnsemblSite.name == 'Ensembl' - assert releases[0].EnsemblRelease.label == release_name @pytest.mark.parametrize( "allow_unreleased, dataset_uuid, release_name, release_status", [ (False, '8801edaf-86ec-4799-8fd4-a59077f04c05', None, None), # No release returned is not allowed - (False, '08543d8d-2110-46f3-a9b6-ac58c4af8202', 'MVP Beta-1', 'Released'), # No release returned is not allowed - (True, 'd57040b6-0ef5-4e6b-97ef-be0ad94d3a61', 'MVP Beta-2', 'Prepared'), # Processed Beta-2 - (True, 'd641779c-2add-46ce-acf4-a2b6f15274b1', 'MVP Beta-3', 'Preparing'), # Processed Beta-2 + (False, '08543d8d-2110-46f3-a9b6-ac58c4af8202', '2020-10-18', 'Released'), + # No release returned is not allowed + (True, 'd57040b6-0ef5-4e6b-97ef-be0ad94d3a61', '2021-10-18', 'Prepared'), # Processed Beta-2 + (True, 'd641779c-2add-46ce-acf4-a2b6f15274b1', '2022-10-18', 'Preparing'), # Processed Beta-2 ], indirect=['allow_unreleased'] ) diff --git a/src/tests/test_protobuf_msg_factory.py b/src/tests/test_protobuf_msg_factory.py index a1651258..df9cf19e 100644 --- a/src/tests/test_protobuf_msg_factory.py +++ b/src/tests/test_protobuf_msg_factory.py @@ -196,7 +196,7 @@ def test_create_genome_assembly_sequence_region(self, genome_conn): (False, 108.0, { "releaseVersion": 108.0, "releaseDate": "2023-06-15", - "releaseLabel": "First Beta", + "releaseLabel": "2023-06-15", "releaseType": "partial", "isCurrent": False, "siteName": "Ensembl", @@ -205,8 +205,8 @@ def test_create_genome_assembly_sequence_region(self, genome_conn): }), (False, 110.1, { "releaseVersion": 110.1, - "releaseDate": "2023-10-18", - "releaseLabel": "MVP Beta-1", + "releaseDate": "2020-10-18", + "releaseLabel": "2020-10-18", "releaseType": "partial", "isCurrent": True, "siteName": "Ensembl", @@ -215,8 +215,8 @@ def test_create_genome_assembly_sequence_region(self, genome_conn): }), (True, 110.3, { "releaseVersion": 110.3, - "releaseDate": "Unreleased", - "releaseLabel": "MVP Beta-3", + "releaseDate": "2022-10-18", + "releaseLabel": "2022-10-18", "releaseType": "partial", "isCurrent": False, "siteName": "Ensembl", @@ -293,20 +293,20 @@ def test_create_genome_uuid(self, genome_conn, genome_tag, current_only, expecte "genome_uuid, expected_output", [ ( - # Human - "65d4f21f-695a-4ed0-be67-5732a551fea4", - { - "faaLocation": "GCA/018/473/295/1/vep/genome/softmasked.fa.bgz", - "gffLocation": "GCA/018/473/295/1/vep/ensembl/geneset/2022_08/genes.gff3.bgz" - } + # Human + "65d4f21f-695a-4ed0-be67-5732a551fea4", + { + 'faaLocation': 'GCA/018/473/295/1/genome/unmasked.fa.bgz', + 'gffLocation': 'GCA/018/473/295/1/ensembl/2022_08/geneset/genes.gff3.bgz' + } ), ( - # Ecoli - "a73351f7-93e7-11ec-a39d-005056b38ce3", - { - 'faaLocation': 'GCA/000/005/845/2/vep/genome/softmasked.fa.bgz', - 'gffLocation': 'GCA/000/005/845/2/vep/community/geneset/2018_09/genes.gff3.bgz' - } + # Ecoli + "a73351f7-93e7-11ec-a39d-005056b38ce3", + { + 'faaLocation': 'GCA/000/005/845/2/genome/unmasked.fa.bgz', + 'gffLocation': 'GCA/000/005/845/2/community/2018_09/geneset/genes.gff3.bgz' + } ) ] ) @@ -320,5 +320,5 @@ def test_create_vep_file_paths(self, vep_conn, genome_uuid, expected_output): def test_create_vep_file_paths_invalid_uuid(self, vep_conn): invalid_uuid = "some-invalid-genome-uuid-000000000000" - with pytest.raises(ValueError, match=f"No data found for genome UUID: {invalid_uuid}"): + with pytest.raises(ValueError, match=f"Genome with UUID {invalid_uuid} not found"): vep_conn.fetch_vep_locations(invalid_uuid) diff --git a/src/tests/test_utils.py b/src/tests/test_utils.py index 01548cb3..668a7861 100644 --- a/src/tests/test_utils.py +++ b/src/tests/test_utils.py @@ -21,9 +21,9 @@ from ensembl.utils.database import UnitTestDB, DBConnection from google.protobuf import json_format +from ensembl.production.metadata.api.factories.utils import format_accession_path from ensembl.production.metadata.api.models import Genome, Dataset from ensembl.production.metadata.grpc import ensembl_metadata_pb2, utils -from ensembl.production.metadata.api.factories.utils import format_accession_path logger = logging.getLogger(__name__) @@ -525,8 +525,8 @@ def test_get_genomes_by_name(self, genome_conn): }, 'release': { 'isCurrent': True, - 'releaseDate': '2023-10-18', - 'releaseLabel': 'MVP Beta-1', + 'releaseDate': '2020-10-18', + 'releaseLabel': '2020-10-18', 'releaseType': 'partial', 'releaseVersion': 110.1, 'siteLabel': 'MVP Ensembl', @@ -596,7 +596,7 @@ def test_get_genomes_by_name_release_unspecified(self, genome_conn): }, 'release': { 'releaseDate': '2023-06-15', - 'releaseLabel': 'First Beta', + 'releaseLabel': '2023-06-15', 'releaseType': 'partial', 'releaseVersion': 108.0, 'siteLabel': 'MVP Ensembl', @@ -744,16 +744,16 @@ def test_get_release_version_by_uuid(self, genome_conn, genome_uuid, dataset_typ # Human "65d4f21f-695a-4ed0-be67-5732a551fea4", { - "faaLocation": "GCA/018/473/295/1/vep/genome/softmasked.fa.bgz", - "gffLocation": "GCA/018/473/295/1/vep/ensembl/geneset/2022_08/genes.gff3.bgz" + 'faaLocation': 'GCA/018/473/295/1/genome/unmasked.fa.bgz', + 'gffLocation': 'GCA/018/473/295/1/ensembl/2022_08/geneset/genes.gff3.bgz' } ), ( # Ecoli "a73351f7-93e7-11ec-a39d-005056b38ce3", { - 'faaLocation': 'GCA/000/005/845/2/vep/genome/softmasked.fa.bgz', - 'gffLocation': 'GCA/000/005/845/2/vep/community/geneset/2018_09/genes.gff3.bgz' + 'faaLocation': 'GCA/000/005/845/2/genome/unmasked.fa.bgz', + 'gffLocation': 'GCA/000/005/845/2/community/2018_09/geneset/genes.gff3.bgz' } ), ( From 6277bd0dd76e27d22c42bbd027fc517d7e09641a Mon Sep 17 00:00:00 2001 From: danielp Date: Thu, 2 Oct 2025 10:22:10 +0100 Subject: [PATCH 4/5] Updated json generator for new path and added parquet generation --- pyproject.toml | 2 + .../metadata/api/exports/ftp_index.py | 554 ------------ .../api/exports/ftp_index_generator.py | 789 ++++++++++++++++++ 3 files changed, 791 insertions(+), 554 deletions(-) delete mode 100644 src/ensembl/production/metadata/api/exports/ftp_index.py create mode 100644 src/ensembl/production/metadata/api/exports/ftp_index_generator.py diff --git a/pyproject.toml b/pyproject.toml index f7aedb38..60214ae2 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -61,6 +61,8 @@ dependencies = [ "duckdb-engine >= 0.17.0", "pymysql", "mysqlclient", + "pandas", #for ftp_index_generator + "pyarrow" #for ftp_index_generator ] [project.urls] diff --git a/src/ensembl/production/metadata/api/exports/ftp_index.py b/src/ensembl/production/metadata/api/exports/ftp_index.py deleted file mode 100644 index 3a1094be..00000000 --- a/src/ensembl/production/metadata/api/exports/ftp_index.py +++ /dev/null @@ -1,554 +0,0 @@ -import json -import re -from collections import defaultdict -from datetime import datetime - -from ensembl.utils.database import DBConnection -from sqlalchemy import select, func, case -from sqlalchemy.orm import selectinload - -from ensembl.production.metadata.api.exceptions import TypeNotFoundException -from ensembl.production.metadata.api.models.assembly import Assembly -from ensembl.production.metadata.api.models.dataset import Dataset, DatasetType, DatasetAttribute, DatasetSource, \ - DatasetStatus, Attribute -from ensembl.production.metadata.api.models.genome import Genome, GenomeDataset, GenomeRelease -from ensembl.production.metadata.api.models.organism import Organism -from ensembl.production.metadata.api.models.release import EnsemblRelease, ReleaseStatus - - -class FTPMetadataExporter: - """ - Independent class for generating FTP metadata JSON structure. - Builds hierarchical data organized by species -> assemblies -> providers -> releases. - """ - - def __init__(self, metadata_uri): - self.metadata_db = DBConnection(metadata_uri) - - def export_to_json(self, output_file=None): - """ - Export FTP metadata to JSON format. - - Args: - output_file: Optional file path to write JSON. If None, returns dict. - - Returns: - dict: Metadata structure if output_file is None - """ - metadata = self.build_ftp_metadata_json() - - if output_file: - with open(output_file, 'w') as f: - json.dump(metadata, f, indent=2, default=str) - return None - else: - return metadata - - def build_ftp_metadata_json(self): - """ - Build a hierarchical data structure for JSON export containing FTP metadata. - Only includes released datasets. - """ - metadata_structure = { - "last_updated": datetime.now().isoformat(), - "species": {} - } - - with self.metadata_db.session_scope() as session: - genome_data = self._load_all_genome_data(session) - - for genome_uuid, data in genome_data.items(): - self._process_genome_data(data, metadata_structure) - - return metadata_structure - - def _load_all_genome_data(self, session): - """Load all genome data in bulk queries to minimize database round trips.""" - - genomes_query = select(Genome).options( - selectinload(Genome.organism), - selectinload(Genome.assembly), - selectinload(Genome.genome_releases).selectinload(GenomeRelease.ensembl_release) - ).outerjoin(GenomeRelease).outerjoin(EnsemblRelease).outerjoin(GenomeDataset).outerjoin(Dataset).where( - (Dataset.status == DatasetStatus.RELEASED) | - (EnsemblRelease.status == ReleaseStatus.RELEASED) - ).distinct() - - genomes = session.execute(genomes_query).scalars().all() - genome_uuids = [g.genome_uuid for g in genomes] - - if not genome_uuids: - return {} - - datasets_query = select( - Genome.genome_uuid, - Dataset, - DatasetType.name.label('dataset_type_name'), - DatasetSource.name.label('dataset_source_name') - ).select_from( - GenomeDataset - ).join( - Genome, GenomeDataset.genome_id == Genome.genome_id - ).join( - Dataset, GenomeDataset.dataset_id == Dataset.dataset_id - ).join( - DatasetType, Dataset.dataset_type_id == DatasetType.dataset_type_id - ).join( - DatasetSource, Dataset.dataset_source_id == DatasetSource.dataset_source_id - ).where( - Genome.genome_uuid.in_(genome_uuids), - (Dataset.status == DatasetStatus.RELEASED) - ) - - dataset_results = session.execute(datasets_query).all() - dataset_ids = [r.Dataset.dataset_id for r in dataset_results] - - attributes_query = select( - DatasetAttribute.dataset_id, - Attribute.name.label('attribute_name'), - DatasetAttribute.value - ).join( - Attribute, DatasetAttribute.attribute_id == Attribute.attribute_id - ).where( - DatasetAttribute.dataset_id.in_(dataset_ids) - ) if dataset_ids else select().where(False) - - attribute_results = session.execute(attributes_query).all() - - genebuild_query = select( - Genome.genome_uuid, - Organism.scientific_name, - Assembly.accession, - func.max(case( - (Attribute.name == 'genebuild.annotation_source', DatasetAttribute.value), - else_=None - )).label('genebuild_source_name'), - func.max(case( - (Attribute.name == 'genebuild.last_geneset_update', DatasetAttribute.value), - else_=None - )).label('last_geneset_update') - ).select_from( - Genome - ).join( - Organism, Genome.organism_id == Organism.organism_id - ).join( - Assembly, Genome.assembly_id == Assembly.assembly_id - ).join( - GenomeDataset, Genome.genome_id == GenomeDataset.genome_id - ).join( - Dataset, GenomeDataset.dataset_id == Dataset.dataset_id - ).join( - DatasetType, Dataset.dataset_type_id == DatasetType.dataset_type_id - ).join( - DatasetAttribute, Dataset.dataset_id == DatasetAttribute.dataset_id - ).join( - Attribute, DatasetAttribute.attribute_id == Attribute.attribute_id - ).where( - Genome.genome_uuid.in_(genome_uuids), - DatasetType.name == 'genebuild', - Attribute.name.in_(['genebuild.annotation_source', 'genebuild.last_geneset_update']) - ).group_by( - Genome.genome_uuid, - Organism.scientific_name, - Assembly.accession - ) - - genebuild_results = session.execute(genebuild_query).all() - - genome_data = {} - - for genome in genomes: - genome_data[genome.genome_uuid] = { - 'genome': genome, - 'datasets': [], - 'attributes': {}, - 'genebuild_metadata': None - } - - for result in dataset_results: - if result.genome_uuid in genome_data: - genome_data[result.genome_uuid]['datasets'].append({ - 'dataset': result.Dataset, - 'dataset_type_name': result.dataset_type_name, - 'dataset_source_name': result.dataset_source_name - }) - - attributes_by_dataset = defaultdict(dict) - for result in attribute_results: - attributes_by_dataset[result.dataset_id][result.attribute_name] = result.value - - for genome_uuid, data in genome_data.items(): - for dataset_info in data['datasets']: - dataset_id = dataset_info['dataset'].dataset_id - dataset_info['attributes'] = attributes_by_dataset.get(dataset_id, {}) - - for result in genebuild_results: - if result.genome_uuid in genome_data: - genome_data[result.genome_uuid]['genebuild_metadata'] = { - 'scientific_name': result.scientific_name, - 'accession': result.accession, - 'genebuild_source_name': result.genebuild_source_name, - 'last_geneset_update': result.last_geneset_update - } - - return genome_data - - def _process_genome_data(self, data, metadata_structure): - """Process a single genome's data using preloaded information.""" - - genome = data['genome'] - organism = genome.organism - assembly = genome.assembly - - species_key = self._normalize_species_name(organism.scientific_name) - - if species_key not in metadata_structure["species"]: - metadata_structure["species"][species_key] = { - "taxid": organism.taxonomy_id, - "species_taxonomy_id": organism.species_taxonomy_id, - "scientific_name": organism.scientific_name, - "common_name": organism.common_name, - "strain": organism.strain, - "strain_type": organism.strain_type, - "biosample_id": organism.biosample_id, - "assemblies": {} - } - - assembly_key = assembly.accession - - if assembly_key not in metadata_structure["species"][species_key]["assemblies"]: - metadata_structure["species"][species_key]["assemblies"][assembly_key] = { - "name": getattr(assembly, 'name', None), - "level": getattr(assembly, 'level', None), - "genebuild_providers": {}, - "assembly": None - } - - assembly_data = metadata_structure["species"][species_key]["assemblies"][assembly_key] - self._process_genome_datasets_bulk(data, assembly_data) - - def _process_genome_datasets_bulk(self, genome_data, assembly_data): - """Process all datasets for a genome using preloaded data.""" - - genome = genome_data['genome'] - datasets = genome_data['datasets'] - genebuild_metadata = genome_data['genebuild_metadata'] - - genome_is_released = any( - gr.ensembl_release and gr.ensembl_release.status == ReleaseStatus.RELEASED - for gr in genome.genome_releases - ) - - if genebuild_metadata and genebuild_metadata.get('last_geneset_update'): - genebuild_release_info = self._extract_genebuild_release_info(genebuild_metadata) - else: - genebuild_release_info = {"release": "unknown"} - - provider_datasets = defaultdict(lambda: defaultdict(list)) - assembly_dataset_info = None - seen_homology_paths = set() - - for dataset_info in datasets: - dataset = dataset_info['dataset'] - dataset_type = dataset_info['dataset_type_name'] - - provider_for_path = self._extract_provider_from_path(genebuild_metadata) - - if dataset_type == 'assembly': - assembly_dataset_info = dataset_info - else: - if provider_for_path not in assembly_data["genebuild_providers"]: - assembly_data["genebuild_providers"][provider_for_path] = {} - - release_key = genebuild_release_info["release"] - if release_key not in assembly_data["genebuild_providers"][provider_for_path]: - assembly_data["genebuild_providers"][provider_for_path][release_key] = { - "release": genebuild_release_info["release"], - "paths": {} - } - - if dataset_type == 'homologies': - try: - if genebuild_metadata: - temp_paths = self._get_public_paths_bulk(genebuild_metadata, datasets, dataset_type) - if temp_paths: - homology_path = temp_paths[0]["path"] - if homology_path in seen_homology_paths: - continue - seen_homology_paths.add(homology_path) - except: - pass - - provider_datasets[provider_for_path][release_key].append(dataset_info) - - if assembly_dataset_info and genebuild_metadata: - try: - assembly_paths = self._get_public_paths_bulk(genebuild_metadata, datasets, 'assembly') - if assembly_paths: - assembly_path = assembly_paths[0]["path"] - file_paths = self._get_dataset_file_paths( - assembly_path, 'assembly', genome, assembly_data - ) - - assembly_data["assembly"] = { - "files": file_paths - } - except Exception as e: - print(f"Error generating assembly paths for genome {genome.genome_uuid}: {e}") - - for provider, releases in provider_datasets.items(): - for release_key, dataset_list in releases.items(): - release_data = assembly_data["genebuild_providers"][provider][release_key] - - try: - if genebuild_metadata: - paths = self._get_public_paths_bulk(genebuild_metadata, datasets) - - for path_info in paths: - dataset_type = path_info["dataset_type"] - - if dataset_type == 'assembly': - continue - - if self._has_released_dataset_bulk(datasets, dataset_type): - file_paths = self._get_dataset_file_paths( - path_info["path"], dataset_type, genome, assembly_data - ) - - release_data["paths"][dataset_type] = { - "files": file_paths - } - - except Exception as e: - error_msg = f"Error generating paths for genome {genome.genome_uuid}: {e}" - if "Required metadata fields are missing" in str(e): - error_msg += f" (Provider: {provider}, Release: {release_key})" - print(error_msg) - - def _get_public_paths_bulk(self, genebuild_metadata, datasets, dataset_type='all'): - """Generate public FTP paths using preloaded metadata.""" - - if not genebuild_metadata: - return [] - - scientific_name = genebuild_metadata['scientific_name'] - accession = genebuild_metadata['accession'] - genebuild_source_name = genebuild_metadata['genebuild_source_name'] - last_geneset_update = genebuild_metadata['last_geneset_update'] - - missing_fields = [] - if not scientific_name: - missing_fields.append("scientific_name") - if not accession: - missing_fields.append("assembly.accession") - if not genebuild_source_name: - missing_fields.append("genebuild.annotation_source") - if not last_geneset_update: - missing_fields.append("genebuild.last_geneset_update") - - if missing_fields: - raise ValueError( - f"Required metadata fields are missing: {', '.join(missing_fields)}. Please check the database entries.") - - unique_dataset_types = list(set([ - 'regulation' if d['dataset_type_name'] == 'regulatory_features' - else d['dataset_type_name'] - for d in datasets - ])) - - if dataset_type == 'regulatory_features': - dataset_type = 'regulation' - - match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) - if match: - last_geneset_update = match.group(1).replace('-', '_') - - scientific_name = self._normalize_species_name(scientific_name) - genebuild_source_name = genebuild_source_name.lower() - base_path = f"{scientific_name}/{accession}" - common_path = f"{base_path}/{genebuild_source_name}" - - path_templates = { - 'genebuild': f"{common_path}/geneset/{last_geneset_update}", - 'assembly': f"{base_path}/genome", - 'homologies': f"{common_path}/homology/{last_geneset_update}", - 'regulation': f"{common_path}/regulation", - 'variation': f"{common_path}/variation/{last_geneset_update}", - } - - paths = [] - - if dataset_type not in unique_dataset_types and dataset_type != 'all': - raise TypeNotFoundException(f"Dataset Type : {dataset_type} not found in metadata.") - - if dataset_type == 'all': - for t in unique_dataset_types: - if t in path_templates: - paths.append({ - "dataset_type": t, - "path": path_templates[t] - }) - elif dataset_type in path_templates: - paths.append({ - "dataset_type": dataset_type, - "path": path_templates[dataset_type] - }) - else: - raise TypeNotFoundException(f"Dataset Type : {dataset_type} has no associated path.") - - return paths - - def _normalize_species_name(self, scientific_name): - """Normalize species name by replacing dots with underscores and merging multiple underscores.""" - normalized = scientific_name.replace(' ', '_') - normalized = normalized.replace('.', '_') - normalized = re.sub(r'_+', '_', normalized) - return normalized - - def _extract_provider_from_path(self, genebuild_metadata): - """Extract the provider component from the genebuild metadata for use in paths.""" - if not genebuild_metadata or not genebuild_metadata.get('genebuild_source_name'): - return 'unknown' - - provider = genebuild_metadata['genebuild_source_name'].lower() - return provider - - def _extract_genebuild_release_info(self, genebuild_metadata): - """Extract release information from genebuild metadata for use as release key/value.""" - release_info = { - "release": "unknown" - } - - if genebuild_metadata and genebuild_metadata.get('last_geneset_update'): - last_geneset_update = genebuild_metadata['last_geneset_update'] - match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) - if match: - release_info["release"] = match.group(1).replace('-', '_') - - return release_info - - def _extract_release_info_from_ensembl_release(self, genome): - """Extract release information from ensembl_release table.""" - release_info = { - "release": "unknown" - } - - for genome_release in genome.genome_releases: - if genome_release.ensembl_release and genome_release.ensembl_release.status == ReleaseStatus.RELEASED: - ensembl_release = genome_release.ensembl_release - if hasattr(ensembl_release, 'version'): - release_info["release"] = str(ensembl_release.version) - elif hasattr(ensembl_release, 'release_number'): - release_info["release"] = str(ensembl_release.release_number) - elif hasattr(ensembl_release, 'name'): - release_info["release"] = ensembl_release.name - break - - return release_info - - def _has_released_dataset_bulk(self, datasets, dataset_type): - """Check if there's a released dataset of the specified type using preloaded data.""" - - type_mapping = { - 'regulation': 'regulatory_features', - 'genebuild': 'genebuild', - 'assembly': 'assembly', - 'homologies': 'homologies', - 'variation': 'variation' - } - - mapped_type = type_mapping.get(dataset_type, dataset_type) - - return any( - d['dataset_type_name'] == mapped_type - for d in datasets - ) - - def _get_dataset_file_paths(self, base_path, dataset_type, genome, assembly_data): - """Generate specific file paths for a dataset type.""" - - file_paths = {} - - if dataset_type == 'genebuild': - file_paths = { - "annotations": { - "cdna.fa.gz": f"{base_path}/cdna.fa.gz", - "genes.embl.gz": f"{base_path}/genes.embl.gz", - "genes.gff3.gz": f"{base_path}/genes.gff3.gz", - "genes.gtf.gz": f"{base_path}/genes.gtf.gz", - "pep.fa.gz": f"{base_path}/pep.fa.gz", - "xref.tsv.gz": f"{base_path}/xref.tsv.gz" - } - } - - path_parts = base_path.split('/') - if len(path_parts) >= 4: - species = path_parts[0] - assembly = path_parts[1] - provider = path_parts[2] - release = path_parts[4] if len(path_parts) > 4 else path_parts[3] - vep_base = f"{species}/{assembly}/vep/{provider}/geneset/{release}" - - file_paths["vep"] = { - "genes.gff3.bgz": f"{vep_base}/genes.gff3.bgz", - "genes.gff3.bgz.csi": f"{vep_base}/genes.gff3.bgz.csi" - } - - elif dataset_type == 'assembly': - file_paths = { - "genome_sequences": { - "chromosomes.tsv.gz": f"{base_path}/chromosomes.tsv.gz", - "hardmasked.fa.gz": f"{base_path}/hardmasked.fa.gz", - "softmasked.fa.gz": f"{base_path}/softmasked.fa.gz", - "unmasked.fa.gz": f"{base_path}/unmasked.fa.gz" - } - } - - path_parts = base_path.split('/') - if len(path_parts) >= 2: - species = path_parts[0] - assembly = path_parts[1] - vep_genome_base = f"{species}/{assembly}/vep/genome" - - file_paths["vep"] = { - "softmasked.fa.bgz": f"{vep_genome_base}/softmasked.fa.bgz", - "softmasked.fa.bgz.fai": f"{vep_genome_base}/softmasked.fa.bgz.fai", - "softmasked.fa.bgz.gzi": f"{vep_genome_base}/softmasked.fa.bgz.gzi" - } - - elif dataset_type == 'homologies': - organism = genome.organism - assembly = genome.assembly - species_name = self._normalize_species_name(organism.scientific_name) - - release = base_path.split('/')[-1] if '/' in base_path else "unknown" - homology_filename = f"{species_name}-{assembly.accession}-{release}-homology.tsv.gz" - - file_paths = { - "homology_data": { - homology_filename: f"{base_path}/{homology_filename}" - } - } - - elif dataset_type == 'variation': - file_paths = { - "variation_data": { - "variation.vcf.gz": f"{base_path}/variation.vcf.gz" - } - } - - elif dataset_type == 'regulation': - file_paths = { - "regulatory_features": { - "regulation.gff": f"{base_path}/regulation.gff" - } - } - - return file_paths - - -if __name__ == "__main__": - exporter = FTPMetadataExporter("mysql://user:pass@host:port/database") - exporter.export_to_json("ftp_metadata.json") - metadata = exporter.export_to_json() - print(f"Found {len(metadata['species'])} species with released datasets") diff --git a/src/ensembl/production/metadata/api/exports/ftp_index_generator.py b/src/ensembl/production/metadata/api/exports/ftp_index_generator.py new file mode 100644 index 00000000..e2731795 --- /dev/null +++ b/src/ensembl/production/metadata/api/exports/ftp_index_generator.py @@ -0,0 +1,789 @@ +# See the NOTICE file distributed with this work for additional information +# regarding copyright ownership. +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# http://www.apache.org/licenses/LICENSE-2.0 +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +""" +Generate FTP metadata JSON and Parquet for Ensembl releases. + +This module provides functionality to generate hierarchical FTP metadata +structures for genomic datasets, organized by species, assemblies, providers, +and releases. Supports both JSON and Parquet export formats. +""" +import json +import logging +import re +import sys +from collections import defaultdict +from datetime import datetime +from pathlib import Path +from typing import Dict, List, Optional, Set + +from ensembl.utils.database import DBConnection +from pandas import DataFrame +from sqlalchemy import select, func, case +from sqlalchemy.orm import selectinload + +from ensembl.production.metadata.api.adaptors.genome import format_accession_path +from ensembl.production.metadata.api.exceptions import TypeNotFoundException +from ensembl.production.metadata.api.models import * + +# Constants +DATASET_TYPE_GENEBUILD = 'genebuild' +DATASET_TYPE_ASSEMBLY = 'assembly' +DATASET_TYPE_HOMOLOGIES = 'homologies' +DATASET_TYPE_VARIATION = 'variation' + +# Configure logging +logging.basicConfig( + level=logging.INFO, + format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' +) +logger = logging.getLogger(__name__) + + +class FTPIndexGenerator: + """ + Generate FTP metadata JSON and Parquet for Ensembl releases. + + This class builds hierarchical data structures organized by species, + assemblies, genebuild providers, and releases. It queries the metadata + database to gather genome, dataset, and release information, generating + file paths for all available datasets. + + Attributes: + metadata_db: Database connection to the Ensembl metadata database + metadata_uri: URI for the metadata database + """ + + def __init__(self, metadata_uri: str): + """ + Initialize the FTP metadata exporter. + + Args: + metadata_uri: Database URI for the metadata database + + """ + self.metadata_db = DBConnection(metadata_uri) + self.metadata_uri = metadata_uri + + def export_to_json(self, metadata, output_file: str) -> Optional[Dict]: + """ + Export FTP metadata to JSON format. + + Args: + metadata: Dictionary containing the metadata structure to export + output_file: Optional file path to write JSON. If None, returns dict. + + Returns: + Metadata structure if output_file is None, otherwise None + + Raises: + IOError: If file cannot be written + """ + + try: + output_path = Path(output_file) + output_path.parent.mkdir(parents=True, exist_ok=True) + + with open(output_path, 'w') as f: + json.dump(metadata, f, indent=2, default=str) + + logger.info(f"FTP metadata exported to: {output_path.absolute()}") + logger.info(f"Found {len(metadata['species'])} species with released datasets") + return None + except Exception as e: + raise IOError(f"Failed to write JSON file {output_file}: {e}") from e + + def export_to_parquet(self, metadata, output_file: str) -> None: + """ + Export FTP metadata to Parquet format. + + The hierarchical JSON structure is flattened into tabular format with + separate tables for species, files, and dataset information. + + Args: + output_file: File path to write Parquet file + + Raises: + IOError: If file cannot be written + """ + logger.info("Starting FTP metadata export to Parquet") + + df = self._flatten_metadata_to_dataframe(metadata) + + try: + output_path = Path(output_file) + output_path.parent.mkdir(parents=True, exist_ok=True) + + df.to_parquet(output_path, index=False, engine='pyarrow') + + logger.info(f"FTP metadata exported to: {output_path.absolute()}") + logger.info(f"Parquet file contains {len(df)} rows") + except Exception as e: + raise IOError(f"Failed to write Parquet file {output_file}: {e}") from e + + def export(self, output_base: str, formats: List[str] = None) -> None: + """ + Export FTP metadata in the specified format(s). + + By default, exports both JSON and Parquet formats. The data is collected + once and then exported to all requested formats for efficiency. + """ + if formats is None: + formats = ['json', 'parquet'] # Default to both + + formats = [f.lower() for f in formats] + + # Validate formats + valid_formats = {'json', 'parquet'} + invalid_formats = set(formats) - valid_formats + if invalid_formats: + raise ValueError(f"Unsupported format(s): {', '.join(invalid_formats)}...") + + # Build metadata ONCE + logger.info("Building FTP metadata structure") + metadata = self.build_ftp_metadata_json() + logger.info(f"Found {len(metadata['species'])} species with released datasets") + + # Export to ALL requested formats + output_path = Path(output_base) + + if 'json' in formats: + json_path = output_path.with_suffix('.json') + self.export_to_json(metadata, str(json_path)) + + if 'parquet' in formats: + parquet_path = output_path.with_suffix('.parquet') + self.export_to_parquet(metadata, str(parquet_path)) + + def _flatten_metadata_to_dataframe(self, metadata: Dict) -> DataFrame: + """ + Flatten hierarchical metadata structure into a pandas DataFrame. + + Converts the nested JSON structure into a flat table suitable for + Parquet format. Each row represents a file with all related metadata. + + Args: + metadata: Hierarchical metadata dictionary from build_ftp_metadata_json + + Returns: + pandas DataFrame with flattened metadata + """ + rows = [] + + for species_key, species_data in metadata['species'].items(): + for assembly_key, assembly_data in species_data['assemblies'].items(): + for provider, provider_data in assembly_data['genebuild_providers'].items(): + for release, release_data in provider_data.items(): + for dataset_type, dataset_info in release_data['paths'].items(): + for file_category, files in dataset_info['files'].items(): + for file_name, file_path in files.items(): + row = { + # Species information + 'species_key': species_key, + 'taxid': species_data['taxid'], + 'species_taxonomy_id': species_data['species_taxonomy_id'], + 'scientific_name': species_data['scientific_name'], + 'common_name': species_data['common_name'], + 'strain': species_data['strain'], + 'strain_type': species_data['strain_type'], + 'biosample_id': species_data['biosample_id'], + + # Assembly information + 'assembly_accession': assembly_key, + 'assembly_name': assembly_data['name'], + 'assembly_level': assembly_data['level'], + + # Release information + 'genebuild_provider': provider, + 'release': release, + + # Dataset information + 'dataset_type': dataset_type, + 'file_category': file_category, + 'file_name': file_name, + 'file_path': file_path, + + # Metadata + 'last_updated': metadata['last_updated'] + } + rows.append(row) + + df = DataFrame(rows) + + # Convert to appropriate data types + if not df.empty: + df['taxid'] = df['taxid'].astype('Int64') + df['species_taxonomy_id'] = df['species_taxonomy_id'].astype('Int64') + + logger.info(f"Flattened metadata into {len(df)} rows") + return df + + def build_ftp_metadata_json(self) -> Dict: + """ + Build a hierarchical data structure for JSON export containing FTP metadata. + + Only includes released datasets. The structure is organized as: + species -> assemblies -> genebuild_providers -> releases -> paths -> datasets + + Returns: + Dictionary containing the complete metadata structure with keys: + - last_updated: ISO format timestamp + - species: Dictionary of species data keyed by normalized name + """ + metadata_structure = { + "last_updated": datetime.now().isoformat(), + "species": {} + } + + with self.metadata_db.session_scope() as session: + genome_data = self._load_all_genome_data(session) + logger.info(f"Loaded data for {len(genome_data)} genomes") + + for genome_uuid, data in genome_data.items(): + try: + self._process_genome_data(data, metadata_structure) + except Exception as e: + logger.error(f"Error processing genome {genome_uuid}: {e}") + continue + + return metadata_structure + + def _load_all_genome_data(self, session) -> Dict[str, Dict]: + """ + Load all genome data in bulk queries to minimize database round trips. + + Performs optimized bulk queries to fetch genomes, datasets, attributes, + and genebuild metadata for all released genomes. + + Args: + session: Database session + + Returns: + Dictionary mapping genome_uuid to genome data containing: + - genome: Genome object + - datasets: List of dataset information + - attributes: Dataset attributes + - genebuild_metadata: Genebuild-specific metadata + """ + logger.info("Loading genome data from database") + + # Query for genomes with released datasets + genomes_query = select(Genome).options( + selectinload(Genome.organism), + selectinload(Genome.assembly), + selectinload(Genome.genome_releases).selectinload(GenomeRelease.ensembl_release) + ).outerjoin(GenomeRelease).outerjoin(EnsemblRelease).outerjoin(GenomeDataset).outerjoin(Dataset).where( + (Dataset.status == DatasetStatus.RELEASED) | + (EnsemblRelease.status == ReleaseStatus.RELEASED) + ).distinct() + + genomes = session.execute(genomes_query).scalars().all() + genome_uuids = [g.genome_uuid for g in genomes] + + if not genome_uuids: + logger.warning("No genomes found with released datasets") + return {} + + logger.info(f"Found {len(genome_uuids)} genomes with released datasets") + + # Bulk fetch datasets + datasets_query = select( + Genome.genome_uuid, + Dataset, + DatasetType.name.label('dataset_type_name'), + DatasetSource.name.label('dataset_source_name') + ).select_from( + GenomeDataset + ).join( + Genome, GenomeDataset.genome_id == Genome.genome_id + ).join( + Dataset, GenomeDataset.dataset_id == Dataset.dataset_id + ).join( + DatasetType, Dataset.dataset_type_id == DatasetType.dataset_type_id + ).join( + DatasetSource, Dataset.dataset_source_id == DatasetSource.dataset_source_id + ).where( + Genome.genome_uuid.in_(genome_uuids), + Dataset.status == DatasetStatus.RELEASED + ) + + dataset_results = session.execute(datasets_query).all() + dataset_ids = [r.Dataset.dataset_id for r in dataset_results] + + # Bulk fetch attributes + attributes_query = select( + DatasetAttribute.dataset_id, + Attribute.name.label('attribute_name'), + DatasetAttribute.value + ).join( + Attribute, DatasetAttribute.attribute_id == Attribute.attribute_id + ).where( + DatasetAttribute.dataset_id.in_(dataset_ids) + ) if dataset_ids else select().where(False) + + attribute_results = session.execute(attributes_query).all() + + # Bulk fetch genebuild metadata + genebuild_query = select( + Genome.genome_uuid, + Organism.scientific_name, + Assembly.accession, + func.max(case( + (Attribute.name == 'genebuild.annotation_source', DatasetAttribute.value) + )).label('genebuild_source_name'), + func.max(case( + (Attribute.name == 'genebuild.last_geneset_update', DatasetAttribute.value) + )).label('last_geneset_update') + ).select_from( + Genome + ).join( + Organism, Genome.organism_id == Organism.organism_id + ).join( + Assembly, Genome.assembly_id == Assembly.assembly_id + ).join( + GenomeDataset, Genome.genome_id == GenomeDataset.genome_id + ).join( + Dataset, GenomeDataset.dataset_id == Dataset.dataset_id + ).join( + DatasetType, Dataset.dataset_type_id == DatasetType.dataset_type_id + ).join( + DatasetAttribute, Dataset.dataset_id == DatasetAttribute.dataset_id + ).join( + Attribute, DatasetAttribute.attribute_id == Attribute.attribute_id + ).where( + Genome.genome_uuid.in_(genome_uuids), + DatasetType.name == DATASET_TYPE_GENEBUILD, + Attribute.name.in_(['genebuild.annotation_source', 'genebuild.last_geneset_update']) + ).group_by( + Genome.genome_uuid, + Organism.scientific_name, + Assembly.accession + ) + + genebuild_results = session.execute(genebuild_query).all() + + # Organize data structures + genome_data = {} + for genome in genomes: + genome_data[genome.genome_uuid] = { + 'genome': genome, + 'datasets': [], + 'attributes': {}, + 'genebuild_metadata': None + } + + # Map datasets to genomes + for result in dataset_results: + if result.genome_uuid in genome_data: + genome_data[result.genome_uuid]['datasets'].append({ + 'dataset': result.Dataset, + 'dataset_type_name': result.dataset_type_name, + 'dataset_source_name': result.dataset_source_name + }) + + # Map attributes to datasets + attributes_by_dataset = defaultdict(dict) + for result in attribute_results: + attributes_by_dataset[result.dataset_id][result.attribute_name] = result.value + + for genome_uuid, data in genome_data.items(): + for dataset_info in data['datasets']: + dataset_id = dataset_info['dataset'].dataset_id + dataset_info['attributes'] = attributes_by_dataset.get(dataset_id, {}) + + # Map genebuild metadata to genomes + for result in genebuild_results: + if result.genome_uuid in genome_data: + genome_data[result.genome_uuid]['genebuild_metadata'] = { + 'scientific_name': result.scientific_name, + 'accession': result.accession, + 'genebuild_source_name': result.genebuild_source_name, + 'last_geneset_update': result.last_geneset_update + } + + return genome_data + + def _process_genome_data(self, data: Dict, metadata_structure: Dict) -> None: + """ + Process a single genome's data using preloaded information. + + Adds the genome's organism and assembly information to the metadata + structure and delegates dataset processing. + + Args: + data: Genome data dictionary containing genome, datasets, and metadata + metadata_structure: Target metadata structure to populate + """ + genome = data['genome'] + organism = genome.organism + assembly = genome.assembly + + if not organism or not assembly: + logger.warning(f"Skipping genome {genome.genome_uuid}: missing organism or assembly") + return + + species_key = self._normalize_species_name(organism.scientific_name) + + # Initialize species entry if not present + if species_key not in metadata_structure["species"]: + metadata_structure["species"][species_key] = { + "taxid": organism.taxonomy_id, + "species_taxonomy_id": organism.species_taxonomy_id, + "scientific_name": organism.scientific_name, + "common_name": organism.common_name, + "strain": organism.strain, + "strain_type": organism.strain_type, + "biosample_id": organism.biosample_id, + "assemblies": {} + } + + assembly_key = assembly.accession + + # Initialize assembly entry if not present + if assembly_key not in metadata_structure["species"][species_key]["assemblies"]: + metadata_structure["species"][species_key]["assemblies"][assembly_key] = { + "name": getattr(assembly, 'name', None), + "level": getattr(assembly, 'level', None), + "genebuild_providers": {} + } + + assembly_data = metadata_structure["species"][species_key]["assemblies"][assembly_key] + self._process_genome_datasets(data, assembly_data) + + def _process_genome_datasets(self, genome_data: Dict, assembly_data: Dict) -> None: + """ + Process all datasets for a genome using preloaded data. + + Uses GenomeAdaptor to get public paths for all dataset types and + generates file paths for each. Organizes data by provider and release. + + Args: + genome_data: Dictionary containing genome, datasets, and metadata + assembly_data: Target assembly data structure to populate + """ + genome = genome_data['genome'] + datasets = genome_data['datasets'] + genebuild_metadata = genome_data['genebuild_metadata'] + + if not genebuild_metadata or not genebuild_metadata.get('last_geneset_update'): + logger.debug(f"Skipping genome {genome.genome_uuid}: missing genebuild metadata") + return + + genome_uuid = genome.genome_uuid + + # Extract release and provider information + genebuild_release_info = self._extract_genebuild_release_info(genebuild_metadata) + provider_for_path = self._extract_provider_name(genebuild_metadata) + release_key = genebuild_release_info["release"] + + # Initialize provider/release structure + if provider_for_path not in assembly_data["genebuild_providers"]: + assembly_data["genebuild_providers"][provider_for_path] = {} + + if release_key not in assembly_data["genebuild_providers"][provider_for_path]: + assembly_data["genebuild_providers"][provider_for_path][release_key] = { + "paths": {} + } + + release_data = assembly_data["genebuild_providers"][provider_for_path][release_key] + + # Get available dataset types + dataset_types = set(d['dataset_type_name'] for d in datasets) + + seen_homology_paths = set() + + # Use genome_adapter to get paths for all dataset types + try: + all_paths = self._generate_paths_from_metadata( + genebuild_metadata, + dataset_types + ) + + for path_info in all_paths: + dataset_type = path_info['dataset_type'] + base_path = path_info['path'] + + # Skip if this dataset type isn't in our datasets + if dataset_type not in dataset_types: + continue + + # Handle homology deduplication + if dataset_type == DATASET_TYPE_HOMOLOGIES: + if base_path in seen_homology_paths: + continue + seen_homology_paths.add(base_path) + + # Generate file paths for this dataset type + file_paths = self._get_dataset_file_paths( + base_path, dataset_type, genome, assembly_data + ) + + release_data["paths"][dataset_type] = { + "files": file_paths + } + + except (ValueError, TypeNotFoundException) as e: + logger.warning(f"Error generating paths for genome {genome_uuid}: {e}") + except Exception as e: + logger.error(f"Unexpected error generating paths for genome {genome_uuid}: {e}") + + def _normalize_species_name(self, scientific_name: str) -> str: + """ + Normalize species name for use as a key. + + Replaces spaces with underscores, dots with underscores, and merges + multiple consecutive underscores into a single underscore. + + Args: + scientific_name: Original scientific name + + Returns: + Normalized species name suitable for use as a dictionary key + """ + normalized = scientific_name.replace(' ', '_') + normalized = normalized.replace('.', '_') + normalized = re.sub(r'_+', '_', normalized) + return normalized + + def _extract_provider_name(self, genebuild_metadata: Dict) -> str: + """ + Extract the provider component from genebuild metadata. + + Args: + genebuild_metadata: Dictionary containing genebuild metadata + + Returns: + Lowercase provider name, or 'unknown' if not found + """ + if not genebuild_metadata or not genebuild_metadata.get('genebuild_source_name'): + return 'unknown' + + provider = genebuild_metadata['genebuild_source_name'].lower() + return provider + + def _extract_genebuild_release_info(self, genebuild_metadata: Dict) -> Dict[str, str]: + """ + Extract release information from genebuild metadata. + + Parses the last_geneset_update field to extract the release date + in YYYY_MM format. + + Args: + genebuild_metadata: Dictionary containing genebuild metadata + + Returns: + Dictionary with 'release' key containing the formatted release date + or 'unknown' if parsing fails + """ + release_info = { + "release": "unknown" + } + + if genebuild_metadata and genebuild_metadata.get('last_geneset_update'): + last_geneset_update = genebuild_metadata['last_geneset_update'] + match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) + if match: + release_info["release"] = match.group(1).replace('-', '_') + + return release_info + + def _generate_paths_from_metadata( + self, + genebuild_metadata: Dict, + dataset_types: Set[str] + ) -> List[Dict[str, str]]: + """ + Generate public FTP paths directly from preloaded metadata. + + This avoids additional database queries by using data that was already + loaded in bulk. Replicates the logic from GenomeAdaptor.get_public_path + but operates on preloaded data. + + Args: + genebuild_metadata: Dictionary containing genebuild metadata + dataset_types: Set of available dataset type names + + Returns: + List of dictionaries with 'dataset_type' and 'path' keys + + Raises: + ValueError: If required metadata fields are missing + """ + scientific_name = genebuild_metadata.get('scientific_name') + accession = genebuild_metadata.get('accession') + genebuild_source_name = genebuild_metadata.get('genebuild_source_name') + last_geneset_update = genebuild_metadata.get('last_geneset_update') + + # Validate required fields + missing_fields = [] + if not scientific_name: + missing_fields.append("scientific_name") + if not accession: + missing_fields.append("assembly.accession") + if not genebuild_source_name: + missing_fields.append("genebuild.annotation_source") + if not last_geneset_update: + missing_fields.append("genebuild.last_geneset_update") + + if missing_fields: + raise ValueError( + f"Required metadata fields are missing: {', '.join(missing_fields)}" + ) + + # Format the release date + match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) + if not match: + raise ValueError(f"Invalid last_geneset_update format: {last_geneset_update}") + formatted_release = match.group(1).replace('-', '_') + + # Build paths + genebuild_source_name = genebuild_source_name.lower() + base_path = format_accession_path(accession) + common_path = f"{base_path}/{genebuild_source_name}/{formatted_release}" + + path_templates = { + DATASET_TYPE_GENEBUILD: f"{common_path}/geneset", + DATASET_TYPE_ASSEMBLY: f"{base_path}/genome", + DATASET_TYPE_HOMOLOGIES: f"{common_path}/homology/{formatted_release}", + DATASET_TYPE_VARIATION: f"{common_path}/variation/{formatted_release}", + } + + # Generate paths only for available dataset types + paths = [] + for dataset_type in dataset_types: + if dataset_type in path_templates: + paths.append({ + "dataset_type": dataset_type, + "path": path_templates[dataset_type] + }) + + return paths + + def _get_dataset_file_paths( + self, + base_path: str, + dataset_type: str, + genome: Genome, + assembly_data: Dict + ) -> Dict[str, Dict[str, str]]: + """ + Generate specific file paths for a dataset type. + + Creates the complete file structure with .bgz extensions (except for + .gff and .txt files) based on the dataset type. + + Args: + base_path: Base path from GenomeAdaptor + dataset_type: Type of dataset (genebuild, assembly, etc.) + genome: Genome object + assembly_data: Assembly data structure + + Returns: + Dictionary of file categories with their respective file paths + """ + file_paths = {} + + if dataset_type == DATASET_TYPE_GENEBUILD: + file_paths = { + "annotations": { + "cdna.fa.bgz": f"{base_path}/cdna.fa.bgz", + "genes.embl.bgz": f"{base_path}/genes.embl.bgz", + "genes.gff3.bgz": f"{base_path}/genes.gff3.bgz", + "genes.gtf.bgz": f"{base_path}/genes.gtf.bgz", + "pep.fa.bgz": f"{base_path}/pep.fa.bgz", + "xref.tsv.gz": f"{base_path}/xref.tsv.gz" + } + } + + elif dataset_type == DATASET_TYPE_ASSEMBLY: + file_paths = { + "genome_sequences": { + "chromosomes.tsv.gz": f"{base_path}/chromosomes.tsv.gz", + "hardmasked.fa.bgz": f"{base_path}/hardmasked.fa.bgz", + "softmasked.fa.bgz": f"{base_path}/softmasked.fa.bgz", + "unmasked.fa.bgz": f"{base_path}/unmasked.fa.bgz" + } + } + + elif dataset_type == DATASET_TYPE_HOMOLOGIES: + + file_paths = { + "homology_data": { + "homology.txt.gz": f"{base_path}/homology.txt.gz" + } + } + + elif dataset_type == DATASET_TYPE_VARIATION: + file_paths = { + "variation_data": { + "variation.vcf.bgz": f"{base_path}/variation.vcf.bgz" + } + } + + return file_paths + + +def main() -> None: + """Main entry point for the script.""" + import argparse + + parser = argparse.ArgumentParser( + description="Generate FTP metadata in JSON and/or Parquet format for Ensembl releases" + ) + parser.add_argument( + "--metadata-uri", + required=True, + help="Database URI for the metadata database" + ) + parser.add_argument( + "--output-path", + default="species", + help="Base output path for the metadata files (default: species). " + "Extensions (.json, .parquet) will be added automatically." + ) + parser.add_argument( + "--formats", + nargs='+', + choices=["json", "parquet"], + default=None, + help="Output format(s): json and/or parquet. Default: both formats" + ) + parser.add_argument( + "--log-level", + default="INFO", + choices=["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"], + help="Set the logging level (default: INFO)" + ) + + args = parser.parse_args() + + # Set log level + logging.getLogger().setLevel(getattr(logging, args.log_level)) + + try: + exporter = FTPIndexGenerator( + metadata_uri=args.metadata_uri, + ) + exporter.export(args.output_path, formats=args.formats) + + formats_str = "both formats" if args.formats is None else " and ".join(args.formats) + logger.info(f"FTP metadata export completed successfully in {formats_str}") + except ValueError as e: + logger.error(f"Configuration error: {e}") + sys.exit(1) + except Exception as e: + logger.error(f"Error generating FTP metadata: {e}", exc_info=True) + sys.exit(1) + + +if __name__ == "__main__": + main() From 9028dd8e9da44e148b73e56981f2a150f5581bc3 Mon Sep 17 00:00:00 2001 From: danielp Date: Fri, 3 Oct 2025 14:05:20 +0100 Subject: [PATCH 5/5] minor fixes to json generator --- .../api/exports/ftp_index_generator.py | 340 ++++++++++++++---- 1 file changed, 278 insertions(+), 62 deletions(-) diff --git a/src/ensembl/production/metadata/api/exports/ftp_index_generator.py b/src/ensembl/production/metadata/api/exports/ftp_index_generator.py index e2731795..5a5b7645 100644 --- a/src/ensembl/production/metadata/api/exports/ftp_index_generator.py +++ b/src/ensembl/production/metadata/api/exports/ftp_index_generator.py @@ -184,38 +184,78 @@ def _flatten_metadata_to_dataframe(self, metadata: Dict) -> DataFrame: for provider, provider_data in assembly_data['genebuild_providers'].items(): for release, release_data in provider_data.items(): for dataset_type, dataset_info in release_data['paths'].items(): - for file_category, files in dataset_info['files'].items(): - for file_name, file_path in files.items(): - row = { - # Species information - 'species_key': species_key, - 'taxid': species_data['taxid'], - 'species_taxonomy_id': species_data['species_taxonomy_id'], - 'scientific_name': species_data['scientific_name'], - 'common_name': species_data['common_name'], - 'strain': species_data['strain'], - 'strain_type': species_data['strain_type'], - 'biosample_id': species_data['biosample_id'], - - # Assembly information - 'assembly_accession': assembly_key, - 'assembly_name': assembly_data['name'], - 'assembly_level': assembly_data['level'], - - # Release information - 'genebuild_provider': provider, - 'release': release, - - # Dataset information - 'dataset_type': dataset_type, - 'file_category': file_category, - 'file_name': file_name, - 'file_path': file_path, - - # Metadata - 'last_updated': metadata['last_updated'] - } - rows.append(row) + # Handle genebuild and assembly (direct files structure) + if dataset_type in [DATASET_TYPE_GENEBUILD, DATASET_TYPE_ASSEMBLY]: + for file_category, files in dataset_info['files'].items(): + for file_name, file_path in files.items(): + row = { + # Species information + 'species_key': species_key, + 'taxid': species_data['taxid'], + 'species_taxonomy_id': species_data['species_taxonomy_id'], + 'scientific_name': species_data['scientific_name'], + 'common_name': species_data['common_name'], + 'strain': species_data['strain'], + 'strain_type': species_data['strain_type'], + 'biosample_id': species_data['biosample_id'], + + # Assembly information + 'assembly_accession': assembly_key, + 'assembly_name': assembly_data['name'], + 'assembly_level': assembly_data['level'], + + # Release information + 'genebuild_provider': provider, + 'release': release, + 'partial_release_label': None, + + # Dataset information + 'dataset_type': dataset_type, + 'file_category': file_category, + 'file_name': file_name, + 'file_path': file_path, + + # Metadata + 'last_updated': metadata['last_updated'] + } + rows.append(row) + + # Handle homologies and variation (nested by partial release label) + elif dataset_type in [DATASET_TYPE_HOMOLOGIES, DATASET_TYPE_VARIATION]: + for partial_label, partial_data in dataset_info.items(): + for file_category, files in partial_data['files'].items(): + for file_name, file_path in files.items(): + row = { + # Species information + 'species_key': species_key, + 'taxid': species_data['taxid'], + 'species_taxonomy_id': species_data['species_taxonomy_id'], + 'scientific_name': species_data['scientific_name'], + 'common_name': species_data['common_name'], + 'strain': species_data['strain'], + 'strain_type': species_data['strain_type'], + 'biosample_id': species_data['biosample_id'], + + # Assembly information + 'assembly_accession': assembly_key, + 'assembly_name': assembly_data['name'], + 'assembly_level': assembly_data['level'], + + # Release information + 'genebuild_provider': provider, + 'release': release, + 'partial_release_label': partial_label, + + # Dataset information + 'dataset_type': dataset_type, + 'file_category': file_category, + 'file_name': file_name, + 'file_path': file_path, + + # Metadata + 'last_updated': metadata['last_updated'] + } + rows.append(row) df = DataFrame(rows) @@ -227,6 +267,124 @@ def _flatten_metadata_to_dataframe(self, metadata: Dict) -> DataFrame: logger.info(f"Flattened metadata into {len(df)} rows") return df + def _get_partial_release_label(self, dataset) -> str: + """ + Get the partial release label for a homology or variation dataset. + + Args: + dataset: Dataset object + + Returns: + Release label string, or 'unknown' if not found + """ + # Find the GenomeDataset relationship with a partial release + for genome_dataset in dataset.genome_datasets: + if (genome_dataset.ensembl_release and + genome_dataset.ensembl_release.release_type == 'partial' and + genome_dataset.ensembl_release.label): + return genome_dataset.ensembl_release.label + + # Fallback: look for any release with a label + for genome_dataset in dataset.genome_datasets: + if genome_dataset.ensembl_release and genome_dataset.ensembl_release.label: + return genome_dataset.ensembl_release.label + + logger.warning(f"No partial release label found for dataset {dataset.dataset_uuid}") + return 'unknown' + + def _generate_path_for_dataset_type( + self, + genebuild_metadata: Dict, + dataset_type: str + ) -> Dict[str, str]: + """ + Generate public FTP path for a single dataset type (genebuild or assembly). + + Args: + genebuild_metadata: Dictionary containing genebuild metadata + dataset_type: Type of dataset + + Returns: + Dictionary with 'dataset_type' and 'path' keys + """ + scientific_name = genebuild_metadata.get('scientific_name') + accession = genebuild_metadata.get('accession') + genebuild_source_name = genebuild_metadata.get('genebuild_source_name') + last_geneset_update = genebuild_metadata.get('last_geneset_update') + + # Validate required fields + if not all([scientific_name, accession, genebuild_source_name, last_geneset_update]): + raise ValueError("Required metadata fields are missing") + + # Format the release date + match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) + if not match: + raise ValueError(f"Invalid last_geneset_update format: {last_geneset_update}") + formatted_release = match.group(1).replace('-', '_') + + # Build path + genebuild_source_name = genebuild_source_name.lower() + base_path = format_accession_path(accession) + common_path = f"{base_path}/{genebuild_source_name}/{formatted_release}" + + path_templates = { + DATASET_TYPE_GENEBUILD: f"{common_path}/geneset", + DATASET_TYPE_ASSEMBLY: f"{common_path}/genome", + } + + if dataset_type not in path_templates: + raise ValueError(f"Unsupported dataset type: {dataset_type}") + + return { + "dataset_type": dataset_type, + "path": path_templates[dataset_type] + } + + def _generate_path_for_homology_variation( + self, + genebuild_metadata: Dict, + dataset_type: str, + partial_release_label: str + ) -> str: + """ + Generate public FTP path for homology or variation dataset. + + Args: + genebuild_metadata: Dictionary containing genebuild metadata + dataset_type: Type of dataset (homologies or variation) + partial_release_label: The label from the partial release + + Returns: + Path string + """ + accession = genebuild_metadata.get('accession') + genebuild_source_name = genebuild_metadata.get('genebuild_source_name') + last_geneset_update = genebuild_metadata.get('last_geneset_update') + + if not all([accession, genebuild_source_name, last_geneset_update]): + raise ValueError("Required metadata fields are missing") + + # Format the genebuild release date + match = re.match(r'^(\d{4}-\d{2})', last_geneset_update) + if not match: + raise ValueError(f"Invalid last_geneset_update format: {last_geneset_update}") + formatted_genebuild_release = match.group(1).replace('-', '_') + + # Build path with both genebuild release and partial release label + genebuild_source_name = genebuild_source_name.lower() + base_path = format_accession_path(accession) + common_path = f"{base_path}/{genebuild_source_name}/{formatted_genebuild_release}" + + path_templates = { + DATASET_TYPE_HOMOLOGIES: f"{common_path}/homology/{partial_release_label}", + DATASET_TYPE_VARIATION: f"{common_path}/variation/{partial_release_label}", + } + + if dataset_type not in path_templates: + raise ValueError(f"Unsupported dataset type: {dataset_type}") + + return path_templates[dataset_type] + def build_ftp_metadata_json(self) -> Dict: """ Build a hierarchical data structure for JSON export containing FTP metadata. @@ -314,6 +472,8 @@ def _load_all_genome_data(self, session) -> Dict[str, Dict]: ).where( Genome.genome_uuid.in_(genome_uuids), Dataset.status == DatasetStatus.RELEASED + ).options( + selectinload(Dataset.genome_datasets).selectinload(GenomeDataset.ensembl_release) ) dataset_results = session.execute(datasets_query).all() @@ -496,33 +656,25 @@ def _process_genome_datasets(self, genome_data: Dict, assembly_data: Dict) -> No release_data = assembly_data["genebuild_providers"][provider_for_path][release_key] - # Get available dataset types - dataset_types = set(d['dataset_type_name'] for d in datasets) - - seen_homology_paths = set() - - # Use genome_adapter to get paths for all dataset types - try: - all_paths = self._generate_paths_from_metadata( - genebuild_metadata, - dataset_types - ) - - for path_info in all_paths: - dataset_type = path_info['dataset_type'] + # Group datasets by type + datasets_by_type = {} + for dataset_info in datasets: + dtype = dataset_info['dataset_type_name'] + if dtype not in datasets_by_type: + datasets_by_type[dtype] = [] + datasets_by_type[dtype].append(dataset_info) + + # Process genebuild and assembly (one per genebuild release) + for dataset_type in [DATASET_TYPE_GENEBUILD, DATASET_TYPE_ASSEMBLY]: + if dataset_type not in datasets_by_type: + continue + + try: + path_info = self._generate_path_for_dataset_type( + genebuild_metadata, dataset_type + ) base_path = path_info['path'] - # Skip if this dataset type isn't in our datasets - if dataset_type not in dataset_types: - continue - - # Handle homology deduplication - if dataset_type == DATASET_TYPE_HOMOLOGIES: - if base_path in seen_homology_paths: - continue - seen_homology_paths.add(base_path) - - # Generate file paths for this dataset type file_paths = self._get_dataset_file_paths( base_path, dataset_type, genome, assembly_data ) @@ -530,11 +682,75 @@ def _process_genome_datasets(self, genome_data: Dict, assembly_data: Dict) -> No release_data["paths"][dataset_type] = { "files": file_paths } + except (ValueError, TypeNotFoundException) as e: + logger.warning(f"Error generating path for {dataset_type} in genome {genome_uuid}: {e}") + + # Process homologies (can have multiple per genebuild, grouped by partial release) + if DATASET_TYPE_HOMOLOGIES in datasets_by_type: + homologies_by_release = {} + + for dataset_info in datasets_by_type[DATASET_TYPE_HOMOLOGIES]: + dataset = dataset_info['dataset'] + partial_release_label = self._get_partial_release_label(dataset) + + if partial_release_label not in homologies_by_release: + homologies_by_release[partial_release_label] = [] + homologies_by_release[partial_release_label].append(dataset_info) + + # Create structure for each partial release + if DATASET_TYPE_HOMOLOGIES not in release_data["paths"]: + release_data["paths"][DATASET_TYPE_HOMOLOGIES] = {} + + for partial_label, dataset_list in homologies_by_release.items(): + try: + base_path = self._generate_path_for_homology_variation( + genebuild_metadata, DATASET_TYPE_HOMOLOGIES, partial_label + ) + + file_paths = self._get_dataset_file_paths( + base_path, DATASET_TYPE_HOMOLOGIES, genome, assembly_data + ) + + release_data["paths"][DATASET_TYPE_HOMOLOGIES][partial_label] = { + "files": file_paths + } + except (ValueError, TypeNotFoundException) as e: + logger.warning( + f"Error generating homology path for genome {genome_uuid}, release {partial_label}: {e}") + + # Process variations (can have multiple per genebuild, grouped by partial release) + if DATASET_TYPE_VARIATION in datasets_by_type: + variations_by_release = {} + + for dataset_info in datasets_by_type[DATASET_TYPE_VARIATION]: + dataset = dataset_info['dataset'] + partial_release_label = self._get_partial_release_label(dataset) + + if partial_release_label not in variations_by_release: + variations_by_release[partial_release_label] = [] + variations_by_release[partial_release_label].append(dataset_info) + + # Create structure for each partial release + if DATASET_TYPE_VARIATION not in release_data["paths"]: + release_data["paths"][DATASET_TYPE_VARIATION] = {} + + for partial_label, dataset_list in variations_by_release.items(): + try: + base_path = self._generate_path_for_homology_variation( + genebuild_metadata, DATASET_TYPE_VARIATION, partial_label + ) + + file_paths = self._get_dataset_file_paths( + base_path, DATASET_TYPE_VARIATION, genome, assembly_data + ) + + release_data["paths"][DATASET_TYPE_VARIATION][partial_label] = { + "files": file_paths + } + except (ValueError, TypeNotFoundException) as e: + logger.warning( + f"Error generating variation path for genome {genome_uuid}, release {partial_label}: {e}") - except (ValueError, TypeNotFoundException) as e: - logger.warning(f"Error generating paths for genome {genome_uuid}: {e}") - except Exception as e: - logger.error(f"Unexpected error generating paths for genome {genome_uuid}: {e}") def _normalize_species_name(self, scientific_name: str) -> str: """ @@ -652,7 +868,7 @@ def _generate_paths_from_metadata( path_templates = { DATASET_TYPE_GENEBUILD: f"{common_path}/geneset", - DATASET_TYPE_ASSEMBLY: f"{base_path}/genome", + DATASET_TYPE_ASSEMBLY: f"{common_path}/genome", DATASET_TYPE_HOMOLOGIES: f"{common_path}/homology/{formatted_release}", DATASET_TYPE_VARIATION: f"{common_path}/variation/{formatted_release}", }