diff --git a/docs/figshare.md b/docs/figshare.md index 404ab6b..b051d1d 100644 --- a/docs/figshare.md +++ b/docs/figshare.md @@ -68,9 +68,9 @@ The POST parameters look like this: } ``` -We search only for datasets by probiding the parameter `"item_type": 3`. +We search only for datasets by providing the parameter `"item_type": 3`. -Example datasets: +Example of datasets related to molecular dynamics: - [Molecular dynamics of DSB in nucleosome](https://figshare.com/articles/dataset/M1_gro/5840706) - [a-Synuclein short MD simulations:homo-A53T](https://figshare.com/articles/dataset/a-Synuclein_short_MD_simulations_homo-A53T/7007552) @@ -88,10 +88,10 @@ We search for all file types and keywords. Results are paginated by batch of 100 Example dataset "[Molecular dynamics of DSB in nucleosome](https://figshare.com/articles/dataset/M1_gro/5840706)": -- web view: -- API view: +- [web view](https://figshare.com/articles/dataset/M1_gro/5840706) +- [API view](https://api.figshare.com/v2/articles/5840706) -All metadata related to a given dataset is provided, as well as all files metadata. +All metadata related to a given dataset is provided, as well as all files metadata for this dataset. ### Zip files diff --git a/src/mdverse_scrapers/core/toolbox.py b/src/mdverse_scrapers/core/toolbox.py index 058745f..231eda4 100644 --- a/src/mdverse_scrapers/core/toolbox.py +++ b/src/mdverse_scrapers/core/toolbox.py @@ -1,11 +1,9 @@ """Common functions and utilities used in the project.""" import argparse -import pathlib import re import time import warnings -from dataclasses import dataclass from datetime import datetime, timedelta from pathlib import Path from urllib.parse import urlparse @@ -17,7 +15,6 @@ from bs4 import BeautifulSoup from ..models.dataset import DatasetMetadata -from ..models.enums import DataType from ..models.file import FileMetadata from ..models.scraper import ScraperContext @@ -29,16 +26,6 @@ ) -@dataclass(kw_only=True) -class ContextManager: - """ContextManager dataclass.""" - - logger: "loguru.Logger" = loguru.logger - output_path: pathlib.Path = pathlib.Path("") - query_file_name: pathlib.Path = pathlib.Path("") - token: str = "" - - def make_http_get_request_with_retries( url: str, params: dict | None = None, @@ -165,14 +152,14 @@ def get_scraper_cli_arguments(): return parser.parse_args() -def read_query_file(query_file_path, logger: "loguru.Logger" = loguru.logger): +def read_query_file(query_file_path: Path, logger: "loguru.Logger" = loguru.logger): """Read the query definition file. The query definition file is formatted in yaml. Parameters ---------- - query_file_path : str + query_file_path : Path Path to the query definition file. logger : "loguru.Logger" Logger for logging messages. @@ -202,54 +189,24 @@ def read_query_file(query_file_path, logger: "loguru.Logger" = loguru.logger): return file_types, keywords, exclusion_file_patterns, exclusion_path_patterns -def verify_file_exists(filename): - """Verify file exists. - - Parameters - ---------- - filename : str - Name of file to verify existence - - Raises - ------ - FileNotFoundError - If file does not exist or is not a file. - """ - file_in = pathlib.Path(filename) - if not file_in.exists(): - msg = f"File {filename} not found" - raise FileNotFoundError(msg) - if not file_in.is_file(): - msg = f"{filename} is not a file" - raise FileNotFoundError(msg) - - -def verify_output_directory(directory, logger: "loguru.Logger" = loguru.logger): - """Verify output directory exists. - - Create it if necessary. +def remove_duplicates_in_list_of_dicts(input_list: list[dict]) -> list[dict]: + """Remove duplicates in a list while preserving the original order. Parameters ---------- - directory : str - Path to directory to store results - logger : "loguru.Logger" - Logger for logging messages. + input_list : list + List with possible duplicate entries. - Raises - ------ - FileNotFoundError - If directory path is an existing file. + Returns + ------- + list + List without duplicates. """ - directory_path = pathlib.Path(directory) - if directory_path.is_file(): - msg = f"{directory} is an existing file." - raise FileNotFoundError(msg) - if directory_path.is_dir(): - logger.info(f"Output directory {directory} already exists.") - else: - directory_path.mkdir(parents=True, exist_ok=True) - logger.info(f"Created output directory {directory}") + output_list = [] + for dict_item in input_list: + if dict_item not in output_list: + output_list.append(dict_item) + return output_list def clean_text(string): @@ -276,53 +233,12 @@ def clean_text(string): return text_decode -def extract_file_extension(file_path: str) -> str: - """Extract file extension from file path. - - Parameters - ---------- - file_path : str - File path - Example: "/something/here/file.txt" - - Returns - ------- - str - File extension without a dot. - Example: "txt" - """ - # Extract the file name for its path. - file_name = file_path.split("/")[-1] - file_type = "none" - if "." in file_name: - file_type = file_name.split(".")[-1].lower() - return file_type - - -def extract_date(date_str): - """Extract and format date from a string. - - Parameters - ---------- - date_str : str - Date as a string in ISO 8601. - For example: 2020-07-29T19:22:57.752335+00:00 - - Returns - ------- - str - Date as in string in YYYY-MM-DD format. - For example: 2020-07-29 - """ - date = datetime.fromisoformat(date_str) - return f"{date:%Y-%m-%d}" - - def remove_excluded_files( - files_df: pd.DataFrame, + files_metadata: list[FileMetadata], exclusion_file_patterns: list[str], exclusion_path_patterns: list[str], -) -> pd.DataFrame: + logger: "loguru.Logger" = loguru.logger, +) -> list[FileMetadata]: """Remove excluded files. Excluded files are, for example: @@ -331,45 +247,57 @@ def remove_excluded_files( Parameters ---------- - files_df : Pandas dataframe - Pandas dataframe with files metadata. - exclusion_file_patterns : list + files_metadata : list[FileMetadata] + List of files metadata. + exclusion_file_patterns : list[str] Patterns for file exclusion. - exclusion_path_patterns : list + exclusion_path_patterns : list[str] Patterns for path exclusion. + logger : "loguru.Logger" + Logger for logging messages. Returns ------- - Pandas dataframe - Dataframe without excluded files and paths. + list[FileMetadata] + List of files metadata without excluded files and paths. """ - df_tmp = files_df.copy(deep=True) - # For file names with path, extract file name only: - df_tmp["name"] = df_tmp["file_name"].apply(lambda x: x.split("/")[-1]) - - boolean_mask = pd.Series(data=False, index=files_df.index) - print("-" * 30) - - for pattern in exclusion_path_patterns: - print(f"Selecting file paths containing: {pattern}") - selection = df_tmp["file_name"].str.contains(pat=pattern, regex=False) - print(f"Found {sum(selection)} files") - boolean_mask = boolean_mask | selection - - for pattern in exclusion_file_patterns: - print(f"Selecting file names starting with: {pattern}") - selection = df_tmp["name"].str.startswith(pattern) - print(f"Found {sum(selection)} files") - boolean_mask = boolean_mask | selection - - print(f"Removed {sum(boolean_mask)} excluded files") - print(f"Remaining files: {sum(~boolean_mask)}") - print("-" * 30) - return files_df[~boolean_mask] + excluded_files_count = {} + files_remaining = [] + for file_meta in files_metadata: + is_excluded = False + # Search exclusion patterns in file path. + for pattern in exclusion_path_patterns: + if pattern in file_meta.file_name: + pattern_label = f"in path: {pattern}" + excluded_files_count[pattern_label] = ( + excluded_files_count.get(pattern_label, 0) + 1 + ) + is_excluded = True + break + # Don't check file name patterns if already excluded by path. + if is_excluded: + continue + # Search exclusion patterns in file name. + name = file_meta.file_name.split("/")[-1] + for pattern in exclusion_file_patterns: + if name.startswith(pattern): + pattern_label = f"starting with: {pattern}" + excluded_files_count[pattern_label] = ( + excluded_files_count.get(pattern_label, 0) + 1 + ) + is_excluded = True + break + if not is_excluded: + files_remaining.append(file_meta) + logger.info(f"Removed {len(files_metadata) - len(files_remaining)} excluded files") + for pattern_label, count in excluded_files_count.items(): + logger.info(f"- {count} files excluded for pattern -> {pattern_label}") + logger.info(f"Remaining files: {len(files_remaining)}") + return files_remaining def find_false_positive_datasets( - files_df: pd.DataFrame, + files_metadata: list[FileMetadata], md_file_types: list[str], logger: "loguru.Logger" = loguru.logger, ) -> list[str]: @@ -380,21 +308,22 @@ def find_false_positive_datasets( Parameters ---------- - files_df : pd.DataFrame - Dataframe which contains all files metadata from a given repo. - md_file_types: list - List containing molecular dynamics file types. + files_metadata : list[FileMetadata] + List of files metadata. + md_file_types: list[str] + List of molecular dynamics file types. logger : "loguru.Logger" Logger for logging messages. Returns ------- - list + list[str] List of false positive dataset ids. """ # Get total number of files and unique file types per dataset. + files_df = pd.DataFrame([file_meta.model_dump() for file_meta in files_metadata]) unique_file_types_per_dataset = ( - files_df.groupby("dataset_id")["file_type"] + files_df.groupby("dataset_id_in_repository")["file_type"] .agg(["count", "unique"]) .rename(columns={"count": "total_files", "unique": "unique_file_types"}) .sort_values(by="total_files", ascending=False) @@ -405,9 +334,9 @@ def find_false_positive_datasets( unique_file_types_per_dataset.loc[dataset_id, "unique_file_types"] ) number_of_files = unique_file_types_per_dataset.loc[dataset_id, "total_files"] - dataset_url = files_df.query(f"dataset_id == '{dataset_id}'").iloc[0][ - "dataset_url" - ] + dataset_url = files_df.query( + f"dataset_id_in_repository == '{dataset_id}'" + ).iloc[0]["dataset_url_in_repository"] # For a given dataset, if there is no MD file types in the entire set # of the dataset file types, then we probably have a false-positive dataset, # i.e. a dataset that does not contain any MD data. @@ -419,149 +348,93 @@ def find_false_positive_datasets( logger.info(f"Dataset will be removed with its {number_of_files} files.") logger.info("List of the first file types:") logger.info(" ".join(dataset_file_types[:20])) - logger.info("---") + logger.info("-" * 30) false_positive_datasets.append(dataset_id) logger.info( - f"In total, {len(false_positive_datasets):,} false positive datasets " - "will be removed." + f"In total, {len(false_positive_datasets):,} " + "false positive datasets will be removed." ) - logger.info("---") + logger.info("-" * 30) return false_positive_datasets def remove_false_positive_datasets( - df_to_clean: pd.DataFrame, + metadata: list[DatasetMetadata] | list[FileMetadata], dataset_ids_to_remove: list[str], logger: "loguru.Logger" = loguru.logger, -) -> pd.DataFrame: - """Remove false positive datasets from file. +) -> list[DatasetMetadata] | list[FileMetadata]: + """Remove false positive datasets from datasets or files metadata. Parameters ---------- - df_to_clean : pd.DataFrame - Dataframe to clean. - dataset_ids_to_remove : list + metadata : list[DatasetMetadata] | list[FileMetadata] + List of metadata to clean (datasets or files). + dataset_ids_to_remove : list[str] List of dataset ids to remove. logger : "loguru.Logger" Logger for logging messages. Returns ------- - pd.DataFrame - Cleaned dataframe. + list[DatasetMetadata] | list[FileMetadata] + Cleaned metadata. """ - records_count_old = len(df_to_clean) - # We keep rows NOT associated to false-positive dataset ids - df_clean = df_to_clean[~df_to_clean["dataset_id"].isin(dataset_ids_to_remove)] - records_count_clean = len(df_clean) - logger.info( - f"Removing {records_count_old - records_count_clean:,} lines " - f"({records_count_old:,} -> {records_count_clean:,}) in dataframe." - ) - return df_clean + metadata_clean = [ + meta + for meta in metadata + if meta.dataset_id_in_repository not in dataset_ids_to_remove + ] + logger.info(f"Removed: {len(metadata) - len(metadata_clean):,}") + logger.info(f"Remaining: {len(metadata_clean):,}") + return metadata_clean def find_remove_false_positive_datasets( - datasets_df: pd.DataFrame, - files_df: pd.DataFrame, - ctx: ContextManager, -) -> tuple[pd.DataFrame, pd.DataFrame]: + datasets_metadata: list[DatasetMetadata], + files_metadata: list[FileMetadata], + scraper: ScraperContext, + logger: "loguru.Logger" = loguru.logger, +) -> tuple[list[DatasetMetadata], list[FileMetadata]]: """Find and remove false-positive datasets. False-positive datasets do not contain MD-related files. Parameters ---------- - datasets_df : pd.DataFrame - Dataframe with information about datasets. - files_df : pd.DataFrame - Dataframe with information about files. - ctx : toolbox.ContextManager - ContextManager object. + datasets_metadata : list[DatasetMetadata] + List of datasets metadata. + files_metadata : list[FileMetadata] + List of files metadata. + scraper : ScraperContext + ScraperContext object. + logger : "loguru.Logger" + Logger for logging messages. Returns ------- - tuple[pd.DataFrame, pd.DataFrame] - Cleaned dataframes for: - - datasets - - files + tuple[list[DatasetMetadata], list[FileMetadata]] + Cleaned lists of metadata for datasets and files. """ # Read parameter file. - file_types, _, _, _ = read_query_file(ctx.query_file_name) + file_types, _, _, _ = read_query_file(scraper.query_file_path, logger=logger) # List file types from the query parameter file. file_types_lst = [file_type["type"] for file_type in file_types] # Zip is not a MD-specific file type. file_types_lst.remove("zip") # Find false-positive datasets. false_positive_datasets = find_false_positive_datasets( - files_df, file_types_lst, logger=ctx.logger + files_metadata, file_types_lst, logger=logger ) # Remove false-positive datasets from all dataframes. - ctx.logger.info("Removing false-positive datasets in the datasets dataframe...") - datasets_df = remove_false_positive_datasets( - datasets_df, false_positive_datasets, logger=ctx.logger + logger.info("Removing false-positive datasets in datasets...") + datasets_metadata = remove_false_positive_datasets( + datasets_metadata, false_positive_datasets, logger=logger ) - ctx.logger.info("Removing false-positive datasets in the files dataframe...") - files_df = remove_false_positive_datasets( - files_df, false_positive_datasets, logger=ctx.logger + logger.info("Removing false-positive datasets in files...") + files_metadata = remove_false_positive_datasets( + files_metadata, false_positive_datasets, logger=logger ) - return datasets_df, files_df - - -def export_dataframe_to_parquet( - repository_name: str, suffix: DataType, df: pd.DataFrame, ctx: ContextManager -) -> None: - """Export dataframes to parquet file. - - Parameters - ---------- - repository_name : str - Name of the data repository. - suffix : DataType - Suffix for the parquet file name. - df : pd.DataFrame - Dataframe to export. - ctx : ContextManager - ContextManager object. - """ - parquet_name = ctx.output_path / f"{repository_name}_{suffix}.parquet" - df.to_parquet(parquet_name, index=False) - ctx.logger.success(f"Dataframe with {len(df):,} rows exported to:") - ctx.logger.success(parquet_name) - - -def export_list_of_models_to_parquet( - parquet_path: Path, - list_of_models: list[DatasetMetadata] | list[FileMetadata], - logger: "loguru.Logger" = loguru.logger, -) -> int: - """Export list of Pydantic models to parquet file. - - Parameters - ---------- - parquet_path : Path - Path to the output parquet file. - list_of_models : list[DatasetMetadata] | list[FileMetadata] - List of Pydantic models to export. - logger : "loguru.Logger" - Logger for logging messages. - - Returns - ------- - int - Number of exported models. - """ - logger.info("Exporting models to parquet.") - try: - df = pd.DataFrame([model.model_dump() for model in list_of_models]) - df.to_parquet(parquet_path, index=False) - logger.success(f"Exported {len(df):,} rows to:") - logger.success(parquet_path) - return len(df) - except (ValueError, TypeError, OSError) as e: - logger.error("Failed to export models to parquet.") - logger.error(e) - return 0 + return datasets_metadata, files_metadata def validate_http_url(v: str) -> str: @@ -650,7 +523,7 @@ def print_statistics( logger: "loguru.Logger" Logger for logging messages. """ - logger.info("-" * 40) + logger.info("-" * 30) logger.success( f"Number of datasets scraped: {scraper.number_of_datasets_scraped:,}" ) @@ -662,3 +535,5 @@ def print_statistics( f"Scraped {scraper.data_source_name} in: {timedelta(seconds=elapsed_time)} 🎉" ) logger.info(f"Saved log file in: {scraper.log_file_path}") + if scraper.is_in_debug_mode: + logger.warning("---Debug mode was ON---") diff --git a/src/mdverse_scrapers/models/dataset.py b/src/mdverse_scrapers/models/dataset.py index cd6e6cf..785c7ee 100644 --- a/src/mdverse_scrapers/models/dataset.py +++ b/src/mdverse_scrapers/models/dataset.py @@ -41,6 +41,7 @@ class DatasetCoreMetadata(BaseModel): ) dataset_id_in_repository: str = Field( ..., + min_length=1, description="Identifier of the dataset in the source repository.", ) dataset_url_in_repository: str = Field( @@ -79,20 +80,6 @@ class DatasetMetadata(SimulationMetadata, DatasetCoreMetadata): description="URL to access the dataset in the project.", ) - # ------------------------------------------------------------------ - # Statistics metadata - # ------------------------------------------------------------------ - download_number: int | None = Field( - None, - ge=0, - description="Total number of downloads for the dataset.", - ) - view_number: int | None = Field( - None, - ge=0, - description="Total number of views for the dataset.", - ) - # ------------------------------------------------------------------ # Temporal metadata # ------------------------------------------------------------------ @@ -143,6 +130,20 @@ class DatasetMetadata(SimulationMetadata, DatasetCoreMetadata): description="External links to publications or other databases.", ) + # ------------------------------------------------------------------ + # Statistics metadata + # ------------------------------------------------------------------ + download_number: int | None = Field( + None, + ge=0, + description="Total number of downloads for the dataset.", + ) + view_number: int | None = Field( + None, + ge=0, + description="Total number of views for the dataset.", + ) + # ------------------------------------------------------------------ # File-level metadata # ------------------------------------------------------------------ diff --git a/src/mdverse_scrapers/models/file.py b/src/mdverse_scrapers/models/file.py index 6756d47..3dbb507 100644 --- a/src/mdverse_scrapers/models/file.py +++ b/src/mdverse_scrapers/models/file.py @@ -21,14 +21,14 @@ class FileMetadata(DatasetCoreMetadata): # ------------------------------------------------------------------ # Descriptive metadata # ------------------------------------------------------------------ - file_url_in_repository: str = Field( - ..., - description="URL to access the file in the repository.", - ) file_name: str = Field( ..., description="File name.", ) + file_url_in_repository: str = Field( + ..., + description="URL to access the file in the repository.", + ) file_size_in_bytes: ByteSize | None = Field(None, description="File size in bytes.") file_md5: str | None = Field(None, description="MD5 checksum.") containing_archive_file_name: str | None = Field( diff --git a/src/mdverse_scrapers/models/scraper.py b/src/mdverse_scrapers/models/scraper.py index 4396aa5..163764c 100644 --- a/src/mdverse_scrapers/models/scraper.py +++ b/src/mdverse_scrapers/models/scraper.py @@ -57,6 +57,10 @@ class ScraperContext(BaseModel): default_factory=lambda: datetime.now(), description="Datetime when the scraper started.", ) + is_in_debug_mode: bool = Field( + False, # noqa: FBT003 + description="Flag indicating if the scraper is running in debug mode.", + ) @model_validator(mode="after") def define_output_dir_file_paths(self) -> Self: diff --git a/src/mdverse_scrapers/models/utils.py b/src/mdverse_scrapers/models/utils.py index 80a1a68..052d5e8 100644 --- a/src/mdverse_scrapers/models/utils.py +++ b/src/mdverse_scrapers/models/utils.py @@ -1,8 +1,10 @@ """Utils for Pydantic models.""" +from pathlib import Path from typing import Any import loguru +import pandas as pd from pydantic import ValidationError from .dataset import DatasetMetadata @@ -45,3 +47,129 @@ def validate_metadata_against_model( else: logger.debug("Input is a complex structure. Skipping value display.") return None + + +def normalize_datasets_metadata( + datasets_list: list[dict], + logger: "loguru.Logger" = loguru.logger, +) -> list[DatasetMetadata]: + """ + Normalize dataset metadata with a Pydantic model. + + Parameters + ---------- + datasets_list : list[dict] + List of dataset metadata dictionaries. + logger: "loguru.Logger" + Logger for logging messages. + + Returns + ------- + list[DatasetMetadata] + List of successfully validated `DatasetMetadata` objects. + """ + datasets_metadata = [] + for dataset in datasets_list: + dataset_id = dataset.get("dataset_id_in_repository") + logger.info(f"Normalizing metadata for dataset: {dataset_id}") + normalized_metadata = validate_metadata_against_model( + dataset, DatasetMetadata, logger=logger + ) + if not normalized_metadata: + logger.error( + f"Metadata normalization failed for dataset " + f"{dataset_id} " + f"from {dataset.get('dataset_repository_name', 'Unknown')}" + ) + continue + datasets_metadata.append(normalized_metadata) + logger.info( + "Normalized metadata for " + f"{len(datasets_metadata):,}/{len(datasets_list):,} " + f"({len(datasets_metadata) / len(datasets_list):.0%}) datasets." + ) + return datasets_metadata + + +def normalize_files_metadata( + files_list: list[dict], + logger: "loguru.Logger" = loguru.logger, +) -> list[FileMetadata]: + """ + Normalize file metadata with a Pydantic model. + + Parameters + ---------- + files_list : list[dict] + List of file metadata dictionaries. + logger: "loguru.Logger" + Logger for logging messages. + + Returns + ------- + list[FileMetadata] + List of successfully validated `FileMetadata` objects. + """ + files_metadata = [] + previous_dataset_id = "" + for file_meta in files_list: + dataset_id = file_meta.get("dataset_id_in_repository") + # Print info only when changing dataset. + if dataset_id != previous_dataset_id: + logger.info(f"Normalizing metadata for files in dataset: {dataset_id}") + normalized_metadata = validate_metadata_against_model( + file_meta, FileMetadata, logger=logger + ) + if not normalized_metadata: + logger.error( + "Metadata normalization failed for file: " + f"{file_meta.get('file_name', 'Unknown')}" + ) + logger.info( + f"In dataset: {dataset_id} from " + f"{file_meta.get('dataset_repository_name', 'Unknown')}" + ) + continue + files_metadata.append(normalized_metadata) + # Print info only when changing dataset. + if dataset_id != previous_dataset_id: + previous_dataset_id = dataset_id + logger.info( + "Normalized metadata for " + f"{len(files_metadata):,}/{len(files_list):,} " + f"({len(files_metadata) / len(files_list):.0%}) files." + ) + return files_metadata + + +def export_list_of_models_to_parquet( + parquet_path: Path, + list_of_models: list[DatasetMetadata] | list[FileMetadata], + logger: "loguru.Logger" = loguru.logger, +) -> int: + """Export list of Pydantic models to parquet file. + + Parameters + ---------- + parquet_path : Path + Path to the output parquet file. + list_of_models : list[DatasetMetadata] | list[FileMetadata] + List of Pydantic models to export. + logger : "loguru.Logger" + Logger for logging messages. + + Returns + ------- + int + Number of exported models. + """ + try: + df = pd.DataFrame([model.model_dump() for model in list_of_models]) + df.to_parquet(parquet_path, index=False) + logger.success(f"Exported {len(df):,} rows to:") + logger.success(parquet_path) + return len(df) + except (ValueError, TypeError, OSError) as e: + logger.error("Failed to export models to parquet.") + logger.error(e) + return 0 diff --git a/src/mdverse_scrapers/scrapers/figshare.py b/src/mdverse_scrapers/scrapers/figshare.py index 197d4c9..aa52b15 100644 --- a/src/mdverse_scrapers/scrapers/figshare.py +++ b/src/mdverse_scrapers/scrapers/figshare.py @@ -1,33 +1,34 @@ """Scrape molecular dynamics datasets and files from Figshare.""" -from arrow import get import json import os import sys -import time -from datetime import datetime, timedelta from pathlib import Path import click import loguru -import pandas as pd from dotenv import load_dotenv +from mdverse_scrapers.models.file import FileMetadata + from ..core.figshare_api import FigshareAPI from ..core.logger import create_logger from ..core.network import get_html_page_with_selenium from ..core.toolbox import ( - ContextManager, - DataType, clean_text, - export_dataframe_to_parquet, - extract_date, - extract_file_extension, find_remove_false_positive_datasets, make_http_get_request_with_retries, + print_statistics, read_query_file, remove_excluded_files, ) +from ..models.enums import DatasetSourceName +from ..models.scraper import ScraperContext +from ..models.utils import ( + export_list_of_models_to_parquet, + normalize_datasets_metadata, + normalize_files_metadata, +) def extract_files_from_json_response( @@ -63,7 +64,8 @@ def extract_files_from_json_response( def extract_files_from_zip_file( - file_id: str, logger: "loguru.Logger" = loguru.logger) -> list[str]: + file_id: str, logger: "loguru.Logger" = loguru.logger +) -> list[str]: """Extract files from a zip file content. No endpoint is available in the Figshare API. @@ -79,7 +81,7 @@ def extract_files_from_zip_file( file_id : str ID of the zip file to get content from. logger : "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -123,7 +125,7 @@ def get_stats_for_dataset( dataset_id: str Dataset id. logger: loguru.Logger - Logger object. + Logger for logging messages. Returns ------- @@ -150,66 +152,60 @@ def get_stats_for_dataset( return stats -def scrap_zip_files( - files_df: pd.DataFrame, logger: "loguru.Logger" = loguru.logger -) -> pd.DataFrame: +def scrap_zip_files_content( + all_files_metadata: list[FileMetadata], logger: "loguru.Logger" = loguru.logger +) -> list[dict]: """Scrap information from files contained in zip archives. - Uncertain how many files can be fetched from the preview. Only get file name and file type. File size and MD5 checksum are not available. Arguments --------- - files_df: Pandas dataframe - Dataframe with information about files. - logger: loguru.Logger - Logger object. + all_files_metadata: list[FileMetadata] + List of files metadata. + logger: "loguru.Logger" + Logger for logging messages. Returns ------- - zip_df: Pandas dataframe - Dataframe with information about files found in zip archive. + list[dict] + List of dictionaries with files metadata found in zip archive. """ - files_in_zip_lst = [] - zip_files_counter = 0 - zip_files_df = files_df[files_df["file_type"] == "zip"] - logger.info(f"Number of zip files to scrap content from: {zip_files_df.shape[0]}") - for zip_files_counter, zip_idx in enumerate(zip_files_df.index, start=1): - zip_file = zip_files_df.loc[zip_idx] - file_id = zip_file["file_url"].split("/")[-1] + files_in_zip_metadata = [] + # Select zip files only. + zip_files = [f_meta for f_meta in all_files_metadata if f_meta.file_type == "zip"] + logger.info(f"Number of zip files to scrap content from: {len(zip_files):,}") + for zip_files_counter, zip_file in enumerate(zip_files, start=1): + zip_file_id = zip_file.file_url_in_repository.split("/")[-1] logger.info("Extracting files from zip file:") - logger.info(zip_file["file_url"]) - file_names = extract_files_from_zip_file(file_id, logger) - if file_names == []: + logger.info(zip_file.file_url_in_repository) + file_names = extract_files_from_zip_file(zip_file_id, logger) + if not file_names: logger.warning("No file found!") continue # Add other metadata. for name in file_names: - file_metadata = {} - file_metadata["dataset_origin"] = zip_file["dataset_origin"] - file_metadata["dataset_id"] = zip_file["dataset_id"] - file_metadata["dataset_url"] = zip_file["dataset_url"] - file_metadata["file_name"] = name - file_metadata["file_type"] = extract_file_extension(name) - file_metadata["file_size"] = None - file_metadata["file_md5"] = None - file_metadata["is_from_zip_file"] = True - file_metadata["containing_zip_file_name"] = zip_file["file_name"] - file_metadata["file_url"] = zip_file["file_url"] - files_in_zip_lst.append(file_metadata) + file_meta = {} + file_meta["dataset_repository_name"] = zip_file.dataset_repository_name + file_meta["dataset_id_in_repository"] = zip_file.dataset_id_in_repository + file_meta["dataset_url_in_repository"] = zip_file.dataset_url_in_repository + file_meta["file_name"] = name + file_meta["file_url_in_repository"] = zip_file.file_url_in_repository + file_meta["containing_archive_file_name"] = zip_file.file_name + files_in_zip_metadata.append(file_meta) logger.info( - f"{zip_files_counter} Figshare zip files processed " - f"({zip_files_counter}/{len(zip_files_df)}" - f":{zip_files_counter / len(zip_files_df):.0%})" + f"{zip_files_counter} zip files from Figshare processed " + f"({zip_files_counter:,}/{len(zip_files):,}" + f":{zip_files_counter / len(zip_files):.0%})" ) - files_in_zip_df = pd.DataFrame(files_in_zip_lst) logger.success("Done extracting files from zip archives.") - return files_in_zip_df + return files_in_zip_metadata def extract_metadata_from_single_dataset_record( record_json: dict, + scraper: ScraperContext, ) -> tuple[dict, list[dict]]: """Extract information from a Figshare dataset/article record. @@ -222,6 +218,8 @@ def extract_metadata_from_single_dataset_record( --------- record_json: dict JSON object obtained after a request on FigShare API. + scraper: ScraperContext + ScraperContext object. Returns ------- @@ -230,59 +228,56 @@ def extract_metadata_from_single_dataset_record( list List of files metadata. """ - dataset_info = {} - files_info = [] - if record_json["is_embargoed"]: - return dataset_info, files_info - dataset_id = str(record_json["id"]) + dataset_metadata = {} + files_metadata = [] + if record_json.get("is_embargoed"): + return dataset_metadata, files_metadata + dataset_id = str(record_json.get("id", "")) # Disable stats for now. # dataset_stats = get_stats_for_dataset(dataset_id) - dataset_stats = {"downloads": None, "views": None} - dataset_info = { - "dataset_origin": "figshare", - "dataset_id": dataset_id, - "doi": record_json["doi"], - "date_creation": extract_date(record_json["created_date"]), - "date_last_modified": extract_date(record_json["modified_date"]), - "date_fetched": datetime.now().isoformat(timespec="seconds"), - "file_number": len(record_json["files"]), - "download_number": dataset_stats["downloads"], - "view_number": dataset_stats["views"], - "license": record_json["license"]["name"], - "dataset_url": record_json["url_public_html"], - "title": clean_text(record_json["title"]), - "author": clean_text(record_json["authors"][0]["full_name"]), - "keywords": "", - "description": clean_text(record_json["description"]), + dataset_stats = {"download_number": None, "view_number": None} + dataset_metadata = { + "dataset_repository_name": scraper.data_source_name, + "dataset_id_in_repository": dataset_id, + "dataset_url_in_repository": record_json.get("url_public_html"), + "date_created": record_json.get("created_date"), + "date_last_updated": record_json.get("modified_date"), + "title": clean_text(record_json.get("title")), + "author_names": [ + clean_text(author.get("full_name")) + for author in record_json.get("authors", []) + ], + "description": clean_text(record_json.get("description")), + "license": record_json.get("license", {}).get("name"), + "doi": record_json.get("doi"), + "download_number": dataset_stats["download_number"], + "view_number": dataset_stats["view_number"], + "number_of_files": len(record_json.get("files", [])), } - # Add keywords only if any. - if "tags" in record_json: - dataset_info["keywords"] = ";".join( - [clean_text(keyword) for keyword in record_json["tags"]] - ) - for file_in in record_json["files"]: - if len(file_in["name"].split(".")) == 1: - filetype = "none" - else: - filetype = file_in["name"].split(".")[-1].lower() - file_dict = { - "dataset_origin": dataset_info["dataset_origin"], - "dataset_id": dataset_info["dataset_id"], - "dataset_url": dataset_info["dataset_url"], - "file_name": file_in["name"], - "file_type": filetype, - "file_size": file_in["size"], - "file_md5": file_in["computed_md5"], - "is_from_zip_file": False, - "containing_zip_file_name": None, - "file_url": file_in["download_url"], + # Add keywords. + dataset_metadata["keywords"] = [ + clean_text(keyword) for keyword in record_json.get("keywords", []) + ] + for file_in in record_json.get("files", []): + file_meta = { + "dataset_repository_name": dataset_metadata["dataset_repository_name"], + "dataset_id_in_repository": dataset_metadata["dataset_id_in_repository"], + "dataset_url_in_repository": dataset_metadata["dataset_url_in_repository"], + "file_name": file_in.get("name"), + "file_url_in_repository": file_in.get("download_url"), + "file_size_in_bytes": file_in.get("size"), + "file_md5": file_in.get("computed_md5"), + "containing_archive_file_name": None, } - files_info.append(file_dict) - return dataset_info, files_info + files_metadata.append(file_meta) + return dataset_metadata, files_metadata def search_all_datasets( - api: FigshareAPI, ctx: ContextManager, max_hits_per_page: int = 100 + api: FigshareAPI, + scraper: ScraperContext, + max_hits_per_page: int = 100, + logger: "loguru.Logger" = loguru.logger, ) -> list[str]: """Search all Figshare datasets. @@ -295,10 +290,12 @@ def search_all_datasets( ---------- api : FigshareAPI Figshare API object. - ctx : ContextManager - ContextManager object. + scraper : ScraperContext + ScraperContext object. max_hits_per_page : int Maximum number of hits per page. + logger : loguru.Logger + Logger for logging messages. Returns ------- @@ -306,14 +303,13 @@ def search_all_datasets( List of Figshare datasets ids. """ # Read parameter file - file_types, keywords, _, _ = read_query_file(ctx.query_file_name, logger=ctx.logger) + file_types, keywords, _, _ = read_query_file(scraper.query_file_path, logger=logger) # We use paging to fetch all results. - # we query max_hits_per_page hits per page. - + # We query max_hits_per_page hits per page. unique_datasets = [] - ctx.logger.info("-" * 30) + logger.info("-" * 30) for file_type in file_types: - ctx.logger.info(f"Looking for filetype: {file_type['type']}") + logger.info(f"Looking for filetype: {file_type['type']}") base_query = f":extension: {file_type['type']}" target_keywords = [""] if file_type["keywords"] == "keywords": @@ -328,11 +324,11 @@ def search_all_datasets( ) else: query = base_query - ctx.logger.info("Search query:") - ctx.logger.info(query) + logger.info("Search query:") + logger.info(query) page = 1 found_datasets_per_keyword = [] - # Search endpoint: + # Search endpoint: /articles/search # https://docs.figshare.com/#articles_search # Iterate seach on pages. while True: @@ -346,13 +342,13 @@ def search_all_datasets( } results = api.query(endpoint="/articles/search", data=data_query) if results["status_code"] >= 400: - ctx.logger.warning( + logger.warning( f"Failed to fetch page {page} " f"for file extension {file_type['type']}" ) - ctx.logger.warning(f"Status code: {results['status_code']}") - ctx.logger.warning(f"Response headers: {results['headers']}") - ctx.logger.warning(f"Response body: {results['response']}") + logger.warning(f"Status code: {results['status_code']}") + logger.warning(f"Response headers: {results['headers']}") + logger.warning(f"Response body: {results['response']}") break response = results["response"] if not response or len(response) == 0: @@ -360,30 +356,34 @@ def search_all_datasets( # Extract datasets ids. found_datasets_per_keyword_per_page = [hit["id"] for hit in response] found_datasets_per_keyword += found_datasets_per_keyword_per_page - ctx.logger.info( + logger.info( f"Page {page} fetched " f"({len(found_datasets_per_keyword_per_page)} datasets)." ) page += 1 found_datasets_per_filetype.update(found_datasets_per_keyword) - ctx.logger.success( + logger.success( f"Found {len(found_datasets_per_filetype)} datasets " f"for filetype: {file_type['type']}" ) - # For debugging purpose, we want unique datasets only, ordered by file types. + # For debugging purpose, we want unique datasets only, + # ordered by file types query. # Instead of a set (sets are unordered), # we use a list and remove duplicates later. unique_datasets += list(found_datasets_per_filetype) # Get unique datasets. # This trick preserves the order datasets were found. unique_datasets = list(dict.fromkeys(unique_datasets)) - ctx.logger.success(f"Found {len(unique_datasets)} unique datasets.") + logger.success(f"Found {len(unique_datasets)} unique datasets.") return unique_datasets -def get_metadata_for_datasets( - api: FigshareAPI, found_datasets: list[str], ctx: ContextManager -) -> tuple[pd.DataFrame, pd.DataFrame]: +def get_metadata_for_datasets_and_files( + api: FigshareAPI, + found_datasets: list[str], + scraper: ScraperContext, + logger: "loguru.Logger" = loguru.logger, +) -> tuple[list[dict], list[dict]]: """Get metadata for all selected datasets. Parameters @@ -392,13 +392,15 @@ def get_metadata_for_datasets( Figshare API object. found_datasets : list[str] List of Figshare dataset ids. - ctx : ContextManager - ContextManager object. + scraper : ScraperContext + ScraperContext object. + logger : "loguru.Logger" + Logger for logging messages. Returns ------- - tuple[pd.DataFrame, pd.DataFrame] - Dataframes for: + tuple[list[dict], list[dict]] + Lists of dictionaries for: - datasets - files """ @@ -408,26 +410,26 @@ def get_metadata_for_datasets( # One dataset at a time. datasets_counter = 0 for datasets_counter, dataset_id in enumerate(found_datasets, start=1): - ctx.logger.info( + logger.info( f"Fetching metadata for Figshare dataset id: {dataset_id} " - f"({datasets_counter}/{len(found_datasets)}" + f"({datasets_counter:,}/{len(found_datasets):,}" f":{datasets_counter / len(found_datasets):.0%})" ) results = api.query(endpoint=f"/articles/{dataset_id}") if results["status_code"] >= 400 or results["response"] is None: - ctx.logger.warning("Failed to fetch dataset.") + logger.warning("Failed to fetch dataset.") continue resp_json_article = results["response"] - dataset_info, files_info = extract_metadata_from_single_dataset_record( - resp_json_article + dataset_metadata, files_metadata = extract_metadata_from_single_dataset_record( + resp_json_article, scraper ) - ctx.logger.info("Done.") - datasets_lst.append(dataset_info) - files_lst += files_info - # Prepare dataframes for export. - datasets_df = pd.DataFrame(data=datasets_lst) - files_df = pd.DataFrame(data=files_lst) - return datasets_df, files_df + logger.info("Done.") + # Append non-empty metadata to datasets and files lists. + if dataset_metadata: + datasets_lst.append(dataset_metadata) + if files_metadata: + files_lst += files_metadata + return datasets_lst, files_lst @click.command( @@ -437,7 +439,7 @@ def get_metadata_for_datasets( @click.option( "--output-dir", "output_dir_path", - type=click.Path(exists=False, file_okay=False, dir_okay=True, path_type=Path), + type=click.Path(exists=True, file_okay=False, dir_okay=True, path_type=Path), required=True, help="Output directory path to save results.", ) @@ -448,70 +450,100 @@ def get_metadata_for_datasets( required=True, help="Query parameters file (YAML format).", ) -def main(output_dir_path: Path, query_file_path: Path) -> None: +@click.option( + "--debug", + "is_in_debug_mode", + is_flag=True, + default=False, + help="Enable debug mode.", +) +def main( + output_dir_path: Path, + query_file_path: Path, + *, + is_in_debug_mode: bool = False, +) -> None: """Scrape Figshare datasets and files.""" - # Define data repository name. - repository_name = "figshare" - # Keep track of script duration. - start_time = time.perf_counter() - # Create context manager. - output_path = output_dir_path / repository_name - output_path.mkdir(parents=True, exist_ok=True) - context = ContextManager( - logger=create_logger(logpath=f"{output_path}/{repository_name}_scraping.log"), - output_path=output_path, - query_file_name=query_file_path, + # Create scraper context. + scraper = ScraperContext( + data_source_name=DatasetSourceName.FIGSHARE, + output_dir_path=output_dir_path, + query_file_path=query_file_path, + is_in_debug_mode=is_in_debug_mode, ) + logger = create_logger(logpath=scraper.log_file_path, level="INFO") # Log script name and doctring. - context.logger.info(__file__) - context.logger.info(__doc__) + logger.info(__file__) + logger.info(__doc__) # Load API tokens. load_dotenv() # Create API object. api = FigshareAPI( token=os.getenv("FIGSHARE_TOKEN"), - logger=context.logger, + logger=logger, ) # Test API token validity. if api.is_token_valid(): - context.logger.success("Figshare token is valid!") + logger.success("Figshare token is valid!") else: - context.logger.error("Figshare token is invalid!") - context.logger.error("Exiting.") + logger.error("Figshare token is invalid!") + logger.error("Exiting.") sys.exit(1) # Search datasets. - found_datasets = search_all_datasets(api, context) + found_datasets = search_all_datasets(api, scraper, logger=logger) + if scraper.is_in_debug_mode: + # Limit number of datasets for debugging purpose. + found_datasets = found_datasets[:20] + found_datasets[-20:] + logger.warning("Debug mode is ON.") + logger.warning("Limiting number of datasets to 40.") # Extract information for all found datasets. - datasets_df, files_df = get_metadata_for_datasets(api, found_datasets, context) - context.logger.success(f"Total number of datasets found: {datasets_df.shape[0]}") - context.logger.success(f"Total number of files found: {files_df.shape[0]}") - + datasets_all, files_all = get_metadata_for_datasets_and_files( + api, found_datasets, scraper, logger=logger + ) + logger.success(f"Total number of datasets found: {len(datasets_all)}") + logger.success(f"Total number of files found: {len(files_all)}") + # Normalize datasets and files metadata. + datasets_all_normalized = normalize_datasets_metadata(datasets_all, logger=logger) + files_all_normalized = normalize_files_metadata(files_all, logger=logger) + logger.success(f"Number of normalized datasets: {len(datasets_all_normalized)}") + logger.success(f"Number of normalized files: {len(files_all_normalized)}") + if (len(datasets_all_normalized) == 0) or (len(files_all_normalized) == 0): + logger.error("No dataset or file left after normalization. Exiting.") + sys.exit(1) # Add files inside zip archives. - zip_df = scrap_zip_files(files_df, context.logger) - context.logger.success(f"Number of files found inside zip files: {zip_df.shape[0]}") - files_df = pd.concat([files_df, zip_df], ignore_index=True) - context.logger.success(f"Total number of files found: {files_df.shape[0]}") - + files_zip = scrap_zip_files_content(files_all_normalized, logger) + logger.success(f"Number of files found inside zip files: {len(files_zip)}") + # Normalize files metadata found inside zip archives. + files_zip_normalized = normalize_files_metadata(files_zip, logger=logger) + files_all_normalized += files_zip_normalized + logger.success(f"Total number of files found: {len(files_all_normalized)}") # Remove unwanted files based on exclusion lists. - context.logger.info("Removing unwanted files...") + logger.info("Removing unwanted files...") _, _, exclude_files, exclude_paths = read_query_file( - query_file_path, context.logger + scraper.query_file_path, logger=logger ) - files_df = remove_excluded_files(files_df, exclude_files, exclude_paths) - context.logger.info("-" * 30) + files_all_normalized = remove_excluded_files( + files_all_normalized, exclude_files, exclude_paths, logger=logger + ) + logger.info("-" * 30) # Remove datasets that contain non-MD related files. - datasets_df, files_df = find_remove_false_positive_datasets( - datasets_df, files_df, context + datasets_all_normalized, files_all_normalized = find_remove_false_positive_datasets( + datasets_all_normalized, files_all_normalized, scraper, logger=logger ) - - # Save dataframes to disk. - export_dataframe_to_parquet("figshare", DataType.DATASETS, datasets_df, context) - export_dataframe_to_parquet("figshare", DataType.FILES, files_df, context) - - # Script duration. - elapsed_time = int(time.perf_counter() - start_time) - context.logger.info(f"Scraped Figshare in: {timedelta(seconds=elapsed_time)}") + # Save metadata to parquet files. + scraper.number_of_datasets_scraped = export_list_of_models_to_parquet( + scraper.datasets_parquet_file_path, + datasets_all_normalized, + logger=logger, + ) + scraper.number_of_files_scraped = export_list_of_models_to_parquet( + scraper.files_parquet_file_path, + files_all_normalized, + logger=logger, + ) + # Print scraping statistics. + print_statistics(scraper, logger=logger) if __name__ == "__main__": diff --git a/src/mdverse_scrapers/scrapers/nomad.py b/src/mdverse_scrapers/scrapers/nomad.py index 1ba530d..be638a0 100644 --- a/src/mdverse_scrapers/scrapers/nomad.py +++ b/src/mdverse_scrapers/scrapers/nomad.py @@ -19,13 +19,16 @@ create_httpx_client, make_http_request_with_retries, ) -from ..core.toolbox import export_list_of_models_to_parquet, print_statistics +from ..core.toolbox import print_statistics from ..models.dataset import DatasetMetadata from ..models.enums import DatasetSourceName -from ..models.file import FileMetadata from ..models.scraper import ScraperContext from ..models.simulation import Molecule, Software -from ..models.utils import validate_metadata_against_model +from ..models.utils import ( + export_list_of_models_to_parquet, + normalize_datasets_metadata, + normalize_files_metadata, +) BASE_NOMAD_URL = "http://nomad-lab.eu/prod/v1/api/v1" JSON_PAYLOAD_NOMAD_REQUEST: dict[str, Any] = { @@ -49,7 +52,7 @@ def is_nomad_connection_working( url : str The URL endpoint. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -72,6 +75,7 @@ def scrape_all_datasets( page_size: int = 50, json_payload: dict[str, Any] | None = None, logger: "loguru.Logger" = loguru.logger, + scraper: ScraperContext | None = None, ) -> list[dict]: """ Scrape Molecular Dynamics-related datasets from the NOMAD API. @@ -88,7 +92,7 @@ def scrape_all_datasets( page_size : int Number of entries to fetch per page. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -175,6 +179,9 @@ def scrape_all_datasets( f"({len(all_datasets):,}/{total_datasets:,}" f":{len(all_datasets) / total_datasets:.0%})" ) + if scraper and scraper.is_in_debug_mode and len(all_datasets) >= 120: + logger.warning("Debug mode is ON: stopping after 120 datasets.") + return all_datasets logger.success(f"Scraped {len(all_datasets):,} datasets in NOMAD.") return all_datasets @@ -183,7 +190,7 @@ def scrape_files_for_all_datasets( client: httpx.Client, datasets: list[DatasetMetadata], logger: "loguru.Logger" = loguru.logger, -) -> list[FileMetadata]: +) -> list[dict]: """Scrape files metadata for all datasets in NOMAD. Parameters @@ -193,12 +200,12 @@ def scrape_files_for_all_datasets( datasets : list[DatasetMetadata] List of datasets to scrape files metadata for. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- - list[FileMetadata] - List of successfully validated `FileMetadata` objects. + list[dict] + List of files metadata dictionaries. """ all_files_metadata = [] for dataset_count, dataset in enumerate(datasets, start=1): @@ -212,27 +219,13 @@ def scrape_files_for_all_datasets( if not files_metadata: continue # Extract relevant files metadata. - files_selected_metadata = extract_files_metadata(files_metadata, logger=logger) + logger.info(f"Getting files metadata for dataset: {dataset_id}") + files_metadata = extract_files_metadata(files_metadata, logger=logger) + all_files_metadata += files_metadata # Normalize files metadata with pydantic model (FileMetadata) - logger.info(f"Validating files metadata for dataset: {dataset_id}") - for file_metadata in files_selected_metadata: - normalized_metadata = validate_metadata_against_model( - file_metadata, - FileMetadata, - logger=logger, - ) - if not normalized_metadata: - logger.error( - f"Normalization failed for metadata of file " - f"{file_metadata.get('file_name')} " - f"in dataset {dataset_id}" - ) - continue - all_files_metadata.append(normalized_metadata) - logger.info("Done.") logger.info(f"Total files found: {len(all_files_metadata):,}") logger.info( - "Extracted and validated files metadata for " + "Extracted files metadata for " f"{dataset_count:,}/{len(datasets):,} " f"({dataset_count / len(datasets):.0%}) datasets." ) @@ -259,7 +252,7 @@ def scrape_files_for_one_dataset( dataset_id : str The unique identifier of the dataset in NOMAD. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -293,7 +286,7 @@ def extract_software_and_version( entry_id : str Identifier of the dataset entry, used for logging. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -316,7 +309,7 @@ def extract_software_and_version( def extract_molecules_and_total_atoms( dataset: dict, entry_id: str, logger: "loguru.Logger" = loguru.logger -) -> tuple[int | None, list[Molecule]]: +) -> tuple[list[Molecule], int | None]: """ Extract molecules and total number of atoms from a dataset entry. @@ -327,16 +320,16 @@ def extract_molecules_and_total_atoms( entry_id : str Identifier of the dataset entry, used for logging. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- - tuple[int | None, list[Molecule]] - total_atoms: Number of atoms for the "original" label, or None if not found. + tuple[list[Molecule], int | None] molecules: List of Molecule objects extracted from the topology. + total_atoms: Number of atoms for the "original" label, or None if not found. """ - total_atoms = None molecules = [] + total_atoms = None try: topologies = dataset.get("results", {}).get("material", {}).get("topology", []) @@ -346,7 +339,7 @@ def extract_molecules_and_total_atoms( if topology.get("label") == "original": total_atoms = topology.get("n_atoms") break - # Extract molecules + # Extract molecules. for topology in topologies: if topology.get("structural_type") == "molecule": molecules.append( # noqa: PERF401 @@ -360,9 +353,10 @@ def extract_molecules_and_total_atoms( logger.warning(f"Topologies is not a list for entry {entry_id}.") logger.warning("Skipping molecules extraction.") except (ValueError, KeyError) as e: - logger.warning(f"Error parsing molecules for entry {entry_id}: {e}") + logger.warning(f"Error parsing molecules for entry: {entry_id}") + logger.warning(e) - return total_atoms, molecules + return molecules, total_atoms def extract_time_step( @@ -380,7 +374,7 @@ def extract_time_step( entry_id : str Identifier of the dataset entry, used for logging. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -418,7 +412,7 @@ def extract_datasets_metadata( datasets : list[dict] List of raw NOMAD datasets metadata. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -427,7 +421,7 @@ def extract_datasets_metadata( """ datasets_metadata = [] for dataset in datasets: - entry_id = dataset.get("entry_id") + entry_id = str(dataset.get("entry_id")) logger.info(f"Extracting relevant metadata for dataset: {entry_id}") entry_url = ( f"https://nomad-lab.eu/prod/v1/gui/search/entries?entry_id={entry_id}" @@ -449,7 +443,7 @@ def extract_datasets_metadata( # Software names with their versions. metadata["software"] = extract_software_and_version(dataset, entry_id, logger) # Molecules with their nb of atoms and number total of atoms. - total_atoms, molecules = extract_molecules_and_total_atoms( + molecules, total_atoms = extract_molecules_and_total_atoms( dataset, entry_id, logger ) metadata["total_number_of_atoms"] = total_atoms @@ -466,48 +460,6 @@ def extract_datasets_metadata( return datasets_metadata -def normalize_datasets_metadata( - datasets: list[dict], - logger: "loguru.Logger" = loguru.logger, -) -> list[DatasetMetadata]: - """ - Normalize dataset metadata with a Pydantic model. - - Parameters - ---------- - datasets : list[dict] - List of dataset metadata dictionaries. - logger: "loguru.Logger" - Logger object. - - Returns - ------- - list[DatasetMetadata] - List of successfully validated `DatasetMetadata` objects. - """ - datasets_metadata = [] - for dataset in datasets: - logger.info( - f"Normalizing metadata for dataset: {dataset['dataset_id_in_repository']}" - ) - normalized_metadata = validate_metadata_against_model( - dataset, DatasetMetadata, logger=logger - ) - if not normalized_metadata: - logger.error( - f"Normalization failed for metadata of dataset " - f"{dataset['dataset_id_in_repository']}" - ) - continue - datasets_metadata.append(normalized_metadata) - logger.info( - "Normalized metadata for " - f"{len(datasets_metadata)}/{len(datasets)} " - f"({len(datasets_metadata) / len(datasets):.0%}) datasets." - ) - return datasets_metadata - - def extract_files_metadata( raw_metadata: dict, logger: "loguru.Logger" = loguru.logger, @@ -520,7 +472,7 @@ def extract_files_metadata( raw_metadata: dict Raw files metadata. logger: "loguru.Logger" - Logger object. + Logger for logging messages. Returns ------- @@ -564,12 +516,20 @@ def extract_files_metadata( required=True, help="Output directory path to save results.", ) -def main(output_dir_path: Path) -> None: +@click.option( + "--debug", + "is_in_debug_mode", + is_flag=True, + default=False, + help="Enable debug mode.", +) +def main(output_dir_path: Path, *, is_in_debug_mode: bool = False) -> None: """Scrape molecular dynamics datasets and files from NOMAD.""" # Create scraper context. scraper = ScraperContext( data_source_name=DatasetSourceName.NOMAD, output_dir_path=output_dir_path, + is_in_debug_mode=is_in_debug_mode, ) logger = create_logger(logpath=scraper.log_file_path, level="INFO") logger.debug(scraper.model_dump_json(indent=4, exclude={"token"})) @@ -590,6 +550,7 @@ def main(output_dir_path: Path) -> None: query_entry_point="entries/query", json_payload=JSON_PAYLOAD_NOMAD_REQUEST, logger=logger, + scraper=scraper, ) if not datasets_raw_metadata: logger.critical("No datasets found in NOMAD.") @@ -599,7 +560,7 @@ def main(output_dir_path: Path) -> None: datasets_selected_metadata = extract_datasets_metadata( datasets_raw_metadata, logger=logger ) - # Parse and validate NOMAD dataset metadata with a pydantic model (DatasetMetadata) + # Validate NOMAD datasets metadata with the DatasetMetadata Pydantic model. datasets_normalized_metadata = normalize_datasets_metadata( datasets_selected_metadata, logger=logger ) @@ -610,18 +571,18 @@ def main(output_dir_path: Path) -> None: logger=logger, ) # Scrape NOMAD files metadata. - files_normalized_metadata = scrape_files_for_all_datasets( + files_metadata = scrape_files_for_all_datasets( client, datasets_normalized_metadata, logger=logger ) - + # Validate NOMAD files metadata with the FileMetadata Pydantic model. + files_normalized_metadata = normalize_files_metadata(files_metadata, logger=logger) # Save files metadata to parquet file. scraper.number_of_files_scraped = export_list_of_models_to_parquet( scraper.files_parquet_file_path, files_normalized_metadata, logger=logger, ) - - # Print script duration. + # Print scraping statistics. print_statistics(scraper, logger=logger) diff --git a/src/mdverse_scrapers/scrapers/zenodo.py b/src/mdverse_scrapers/scrapers/zenodo.py index acda120..59077f1 100644 --- a/src/mdverse_scrapers/scrapers/zenodo.py +++ b/src/mdverse_scrapers/scrapers/zenodo.py @@ -3,8 +3,6 @@ import json import os import sys -import time -from datetime import datetime, timedelta from pathlib import Path import click @@ -15,20 +13,22 @@ from ..core.logger import create_logger from ..core.toolbox import ( - ContextManager, clean_text, - export_dataframe_to_parquet, - extract_date, - extract_file_extension, find_remove_false_positive_datasets, make_http_get_request_with_retries, + print_statistics, read_query_file, + remove_duplicates_in_list_of_dicts, remove_excluded_files, - verify_output_directory, ) -from ..models.enums import DataType - -# logging.getLogger("httpx").setLevel(logging.DEBUG) +from ..models.enums import DatasetSourceName +from ..models.file import FileMetadata +from ..models.scraper import ScraperContext +from ..models.utils import ( + export_list_of_models_to_parquet, + normalize_datasets_metadata, + normalize_files_metadata, +) def get_rate_limit_info( @@ -73,55 +73,6 @@ def get_rate_limit_info( logger.info(f"Header retry-after: {response.headers.get('retry-after', None)}") -def normalize_file_size(file_str): - """Normalize file size in bytes. - - Parameters - ---------- - file_str : str - File size with unit. - For instance: 1.8 GB, 108.7 kB - - Returns - ------- - int - File size in bytes. - """ - size, unit = file_str.split() - if unit == "GB": - size_in_bytes = float(size) * 1_000_000_000 - elif unit == "MB": - size_in_bytes = float(size) * 1_000_000 - elif unit == "kB": - size_in_bytes = float(size) * 1_000 - elif unit == "Bytes": - size_in_bytes = float(size) - else: - size_in_bytes = 0 - return int(size_in_bytes) - - -def extract_license(metadata): - """Extract license from metadata. - - Parameters - ---------- - metadata : dict - Metadata from Zenodo API. - - Returns - ------- - str - License. - Empty string if no license found. - """ - try: - license_name = metadata["license"]["id"] - except KeyError: - license_name = "" - return license_name - - def get_files_structure_from_zip(ul): """Get files structure from zip file preview. @@ -153,7 +104,9 @@ def get_files_structure_from_zip(ul):
  • -
    flow_00001.dat
    +
    + flow_00001.dat +
    4.6 kB
    @@ -161,7 +114,9 @@ def get_files_structure_from_zip(ul):
  • -
    flow_00003.dat
    +
    + flow_00003.dat +
    4.6 kB
    @@ -235,12 +190,12 @@ def extract_data_from_zip_file(url, logger: "loguru.Logger" = loguru.logger): # Normalize file size. for path, size in files_dict.items(): if size: - file_dict = { - "file_name": path, - "file_size": normalize_file_size(size), - "file_type": extract_file_extension(path), - } - file_lst.append(file_dict) + file_lst.append( + { + "file_name": path, + "file_size_in_bytes": size, + } + ) logger.success(f"Found {len(file_lst)} files.") return file_lst @@ -280,8 +235,8 @@ def is_zenodo_connection_working( def scrap_zip_content( - files_df, logger: "loguru.Logger" = loguru.logger -) -> pd.DataFrame: + files_metadata: list[FileMetadata], logger: "loguru.Logger" = loguru.logger +) -> list[dict]: """Scrap information from files contained in zip archives. Zenodo provides a preview only for the first 1000 files within a zip file. @@ -292,55 +247,50 @@ def scrap_zip_content( Arguments --------- - files_df: dataframe - Dataframe with information about files. + files_metadata: list[FileMetadata] + List of files metadata. Returns ------- - zip_df: dataframe - Dataframe with information about files in zip archive. + list[dict] + List of dictionaries with metadata of files found in zip archive. """ files_in_zip_lst = [] - zip_counter = 0 - zip_files_df = files_df[files_df["file_type"] == "zip"] - logger.info(f"Number of zip files to scrap content from: {zip_files_df.shape[0]}") + # Select zip files only. + zip_files = [f_meta for f_meta in files_metadata if f_meta.file_type == "zip"] + logger.info(f"Number of zip files to scrap content from: {len(zip_files)}") # The Zenodo API does not provide any endpoint to get the content of zip files. # We use direct GET requests on the HTML preview of the zip files. - # We wait 1.5 seconds between each request, - # to be gentle with the Zenodo servers. - for zip_counter, zip_idx in enumerate(zip_files_df.index, start=1): - zip_file = zip_files_df.loc[zip_idx] + for zip_counter, zip_file in enumerate(zip_files, start=1): url = ( - f"https://zenodo.org/records/{zip_file['dataset_id']}" - f"/preview/{zip_file.loc['file_name']}" + f"https://zenodo.org/records/{zip_file.dataset_id_in_repository}" + f"/preview/{zip_file.file_name}" ) files_tmp = extract_data_from_zip_file( url, logger=logger, ) - if files_tmp == []: + if not files_tmp: continue # Add common extra fields - for idx in range(len(files_tmp)): - files_tmp[idx]["dataset_origin"] = zip_file["dataset_origin"] - files_tmp[idx]["dataset_id"] = zip_file["dataset_id"] - files_tmp[idx]["is_from_zip_file"] = True - files_tmp[idx]["containing_zip_file_name"] = zip_file["file_name"] - files_tmp[idx]["file_url"] = "" - files_tmp[idx]["file_md5"] = "" - files_in_zip_lst += files_tmp + for file_meta in files_tmp: + file_meta["dataset_repository_name"] = zip_file.dataset_repository_name + file_meta["dataset_id_in_repository"] = zip_file.dataset_id_in_repository + file_meta["dataset_url_in_repository"] = zip_file.dataset_url_in_repository + file_meta["containing_archive_file_name"] = zip_file.file_name + file_meta["file_url_in_repository"] = zip_file.file_url_in_repository + files_in_zip_lst.append(file_meta) logger.info( "Zenodo zip files scraped: " - f"{zip_counter}/{len(zip_files_df)} " - f"({zip_counter / len(zip_files_df):.0%})" + f"{zip_counter}/{len(zip_files)} " + f"({zip_counter / len(zip_files):.0%})" ) - files_in_zip_df = pd.DataFrame(files_in_zip_lst) - return files_in_zip_df + return files_in_zip_lst -def extract_records( - response_json, logger: "loguru.Logger" = loguru.logger -) -> tuple[list, list]: +def extract_metadata_from_json( + response_json: dict, logger: "loguru.Logger" = loguru.logger +) -> tuple[list[dict], list[dict]]: """Extract information from the Zenodo records. Arguments @@ -350,76 +300,70 @@ def extract_records( Returns ------- - datasets: list - List of dictionnaries. Information on datasets. - files: list - List of dictionaries. Information on files. + datasets: list[dict] + List of datasets metadata. + files: list[dict] + List of files metadata. """ datasets = [] files = [] - if response_json["hits"]["hits"]: - for hit in response_json["hits"]["hits"]: - # 'hit' is a Python dictionary. - if hit["metadata"]["access_right"] != "open": - continue - dataset_id = str(hit["id"]) - logger.info(f"Extracting metadata for dataset id: {dataset_id}") - dataset_dict = { - "dataset_origin": "zenodo", - "dataset_id": dataset_id, - "doi": hit["doi"], - "date_creation": extract_date(hit["created"]), - "date_last_modified": extract_date(hit["updated"]), - "date_fetched": datetime.now().isoformat(timespec="seconds"), - "file_number": len(hit["files"]), - "download_number": int(hit["stats"]["downloads"]), - "view_number": int(hit["stats"]["views"]), - "license": extract_license(hit["metadata"]), - "dataset_url": hit["links"]["self_html"], - "title": clean_text(hit["metadata"]["title"]), - "author": clean_text(hit["metadata"]["creators"][0]["name"]), - "keywords": "none", - "description": clean_text(hit["metadata"].get("description", "")), + try: + _ = response_json["hits"]["hits"] + except KeyError: + logger.warning("Cannot extract hits from the response JSON.") + return datasets, files + for hit in response_json["hits"]["hits"]: + # 'hit' is a Python dictionary. + if hit.get("metadata", {}).get("access_right", "") != "open": + continue + dataset_id = str(hit["id"]) + logger.info(f"Extracting metadata for dataset id: {dataset_id}") + dataset_dict = { + "dataset_repository_name": DatasetSourceName.ZENODO, + "dataset_id_in_repository": dataset_id, + "dataset_url_in_repository": hit.get("links", {}).get("self_html", ""), + "date_created": hit.get("created", None), + "date_last_updated": hit.get("modified", None), + "title": clean_text(hit.get("metadata", {}).get("title", "")), + "author_names": [ + author.get("name") + for author in hit.get("metadata", {}).get("creators", []) + if author.get("name", None) + ], + "description": clean_text(hit.get("metadata", {}).get("description", "")), + "keywords": [ + str(keyword) for keyword in hit.get("metadata", {}).get("keywords", []) + ], + "license": hit.get("metadata", {}).get("license", {}).get("id", None), + "doi": hit.get("doi", None), + "number_of_files": len(hit.get("files", [])), + "download_number": hit.get("stats", {}).get("downloads", None), + "view_number": hit.get("stats", {}).get("views", None), + } + datasets.append(dataset_dict) + logger.info(f"Dataset URL: {dataset_dict['dataset_url_in_repository']}") + for file_in in hit.get("files", []): + file_dict = { + "dataset_repository_name": dataset_dict["dataset_repository_name"], + "dataset_id_in_repository": dataset_dict["dataset_id_in_repository"], + "dataset_url_in_repository": dataset_dict["dataset_url_in_repository"], + "file_name": file_in.get("key", ""), + "file_url_in_repository": file_in.get("links", {}).get("self", ""), + # File size in bytes. + "file_size_in_bytes": file_in.get("size", None), + "file_md5": file_in.get("checksum", "").removeprefix("md5:"), + "containing_archive_file_name": None, } - if "keywords" in hit["metadata"]: - dataset_dict["keywords"] = ";".join( - [str(keyword) for keyword in hit["metadata"]["keywords"]] - ) - # Handle existing but empty keywords. - # For instance: https://zenodo.org/records/3741678 - if dataset_dict["keywords"] == "": - dataset_dict["keywords"] = "none" - datasets.append(dataset_dict) - logger.info(f"Dataset URL: {dataset_dict['dataset_url']}") - for file_in in hit["files"]: - file_dict = { - "dataset_origin": dataset_dict["dataset_origin"], - "dataset_id": dataset_dict["dataset_id"], - "dataset_url": dataset_dict["dataset_url"], - "file_size": int(file_in["size"]), # File size in bytes. - "file_md5": file_in["checksum"].removeprefix("md5:"), - "is_from_zip_file": False, - "file_name": file_in["key"], - "file_type": extract_file_extension(file_in["key"]), - "file_url": file_in["links"]["self"], - "containing_zip_file_name": "none", - } - # Some file types could be empty. - # See for instance file "lmp_mpi" in: - # https://zenodo.org/record/5797177 - # https://zenodo.org/api/records/5797177 - # Set these file types to "none". - if file_dict["file_type"] == "": - file_dict["file_type"] = "none" - files.append(file_dict) + files.append(file_dict) return datasets, files def search_zenodo( query: str, - ctx: ContextManager, + scraper: ScraperContext, page: int = 1, number_of_results: int = 1, + logger: "loguru.Logger" = loguru.logger, ) -> dict | None: """Get total number of hits for a given query. @@ -427,8 +371,14 @@ def search_zenodo( ---------- query : str The search query string. - ctx : ContextManager - Context manager containing configuration and logger. + scraper: ScraperContext + Scraper context manager containing configuration. + page : int, optional + The page number to retrieve. Default is 1. + number_of_results : int, optional + Number of results per page. Default is 1. + logger : loguru.Logger, optional + Logger for logging messages. Returns ------- @@ -440,69 +390,46 @@ def search_zenodo( "size": number_of_results, "page": page, "status": "published", - "access_token": ctx.token, + "access_token": scraper.token, } response_json = None response = make_http_get_request_with_retries( url="https://zenodo.org/api/records", params=params, timeout=60, - logger=ctx.logger, + logger=logger, delay_before_request=2, max_attempts=5, ) if response is None: - ctx.logger.warning("Failed to get response from the Zenodo API.") - ctx.logger.warning("Getting next file type...") + logger.warning("Failed to get response from the Zenodo API.") + logger.warning("Getting next file type...") return None # Try to decode JSON response. try: response_json = response.json() except (json.decoder.JSONDecodeError, ValueError) as exc: - ctx.logger.warning("Failed to decode JSON response from the Zenodo API.") - ctx.logger.warning(f"Error: {exc}") + logger.warning("Failed to decode JSON response from the Zenodo API.") + logger.warning(f"Error: {exc}") return None # Try to extract hits (= results). try: _ = response_json["hits"] _ = int(response_json["hits"]["total"]) except (KeyError, ValueError): - ctx.logger.warning("Cannot extract hits for HTTP response.") - ctx.logger.debug("Response JSON") - ctx.logger.debug(response_json) + logger.warning("Cannot extract hits for HTTP response.") + logger.debug("Response JSON") + logger.debug(response_json) return None return response_json -def merge_dataframes_remove_duplicates( - df1: pd.DataFrame, df2: pd.DataFrame, on_columns: list[str] | None = None -) -> pd.DataFrame: - """Merge two dataframes and remove duplicates. - - Parameters - ---------- - df1 : pd.DataFrame - First dataframe. - df2 : pd.DataFrame - Second dataframe. - on_columns : list of str, optional - List of columns to consider for duplicates. - If None, all columns are considered. - - Returns - ------- - pd.DataFrame - Merged dataframe with duplicates removed. - """ - df_concat = pd.concat([df1, df2], ignore_index=True) - return df_concat.drop_duplicates(subset=on_columns, keep="first") - - def search_all_datasets( file_types: list[dict], keywords: list[str], - ctx: ContextManager, -) -> tuple[pd.DataFrame, pd.DataFrame]: + scraper: ScraperContext, + logger: "loguru.Logger" = loguru.logger, +) -> tuple[list[dict], list[dict]]: """Search all datasets on Zenodo. Parameters @@ -514,15 +441,17 @@ def search_all_datasets( - keywords: str, "keywords" or "none" keywords : list of str List of keywords to use in the search. - ctx : ContextManager + scraper: ScraperContext Context manager containing configuration and logger. + logger : "loguru.Logger" + Logger for logging messages. Returns ------- - datasets_df : pd.DataFrame - Dataframe with information on datasets. - files_df : pd.DataFrame - Dataframe with information on files. + datasets : list of dict + List with datasets metadata. + files : list of dict + List with files metadata. """ # There is a hard limit of the number of hits # one can get from a single query. @@ -532,69 +461,74 @@ def search_all_datasets( # Build query part with keywords. We want something like: # AND ("KEYWORD 1" OR "KEYWORD 2" OR "KEYWORD 3") query_keywords = ' AND ("' + '" OR "'.join(keywords) + '")' - # Create empty dataframes to store results. - datasets_df = pd.DataFrame() - files_df = pd.DataFrame() - ctx.logger.info("-" * 30) + # Create empty lists to store results. + datasets = [] + files = [] + logger.info("-" * 30) for file_type in file_types: - ctx.logger.info(f"Looking for filetype: {file_type['type']}") - datasets_count_old = datasets_df.shape[0] + logger.info(f"Looking for filetype: {file_type['type']}") + datasets_count_old = len(datasets) # Build query with file type and optional keywords. query = f"""resource_type.type:"dataset" AND filetype:"{file_type["type"]}" """ if file_type["keywords"] == "keywords": query += query_keywords - ctx.logger.info("Query:") - ctx.logger.info(f"{query}") + logger.info("Query:") + logger.info(f"{query}") # First, get the total number of hits for a given query. # This is needed to compute the number of pages of results to get. - json_response = search_zenodo(query, ctx, page=1, number_of_results=1) + json_response = search_zenodo( + query, scraper, page=1, number_of_results=1, logger=logger + ) if json_response is None or int(json_response["hits"]["total"]) == 0: - ctx.logger.error("Getting next file type...") - ctx.logger.info("-" * 30) + logger.error("Getting next file type...") + logger.info("-" * 30) continue total_hits = int(json_response["hits"]["total"]) - ctx.logger.info(f"Total hits for this query: {total_hits}") + logger.info(f"Total hits for this query: {total_hits}") page_max = total_hits // max_hits_per_page + 1 + if scraper.is_in_debug_mode: + logger.warning("Debug mode is ON") + logger.warning("Limiting the number of pages to 1 with 10 hits per page.") + page_max = 1 + max_hits_per_page = 10 # Then, slice the query by page. for page in range(1, page_max + 1): - ctx.logger.info( + logger.info( f"Starting page {page}/{page_max} for filetype: {file_type['type']}" ) json_response = search_zenodo( - query, ctx, page=page, number_of_results=max_hits_per_page + query, + scraper, + page=page, + number_of_results=max_hits_per_page, + logger=logger, ) if json_response is None: - ctx.logger.warning("Failed to get response from the Zenodo API.") - ctx.logger.warning("Getting next page...") + logger.warning("Failed to get response from the Zenodo API.") + logger.warning("Getting next page...") continue - datasets_tmp, files_tmp = extract_records(json_response, logger=ctx.logger) - # Merge dataframes - datasets_df = merge_dataframes_remove_duplicates( - datasets_df, - pd.DataFrame(datasets_tmp), - on_columns=["dataset_origin", "dataset_id"], + datasets_tmp, files_tmp = extract_metadata_from_json( + json_response, logger=logger ) - files_df = merge_dataframes_remove_duplicates( - files_df, - pd.DataFrame(files_tmp), - on_columns=["dataset_id", "file_name"], - ) - ctx.logger.success( - f"Found so far: {datasets_df.shape[0]:,} datasets " - f"and {len(files_df):,} files" + # Merge datasets and remove duplicates. + datasets = remove_duplicates_in_list_of_dicts(datasets + datasets_tmp) + # Merge files and remove duplicates. + files = remove_duplicates_in_list_of_dicts(files + files_tmp) + logger.success( + f"Found so far: {len(datasets):,} datasets and {len(files):,} files" ) if page * max_hits_per_page >= max_hits_per_query: - ctx.logger.info("Max hits per query reached!") + logger.info("Max hits per query reached!") break - ctx.logger.info( + logger.info( f"Number of datasets found: {len(datasets_tmp):,} " - f"({datasets_df.shape[0] - datasets_count_old} new)" + f"({len(datasets) - datasets_count_old} new)" ) - ctx.logger.info(f"Number of files found: {len(files_tmp):,}") - ctx.logger.info("-" * 30) - ctx.logger.info(f"Total number of datasets found: {len(datasets_df):,}") - ctx.logger.info(f"Total number of files found: {len(files_df):,}") - return datasets_df, files_df + logger.info(f"Number of files found: {len(files_tmp):,}") + logger.info("-" * 30) + logger.info(f"Total number of datasets found: {len(datasets):,}") + logger.info(f"Total number of files found: {len(files):,}") + return datasets, files @click.command( @@ -604,7 +538,7 @@ def search_all_datasets( @click.option( "--output-dir", "output_dir_path", - type=click.Path(exists=False, file_okay=False, dir_okay=True, path_type=Path), + type=click.Path(exists=True, file_okay=False, dir_okay=True, path_type=Path), required=True, help="Output directory path to save results.", ) @@ -615,39 +549,44 @@ def search_all_datasets( required=True, help="Query parameters file (YAML format).", ) -def main(output_dir_path: Path, query_file_path: Path): +@click.option( + "--debug", + "is_in_debug_mode", + is_flag=True, + default=False, + help="Enable debug mode.", +) +def main( + output_dir_path: Path, query_file_path: Path, *, is_in_debug_mode: bool = False +) -> None: """Scrape Zenodo datasets and files.""" - # Define data repository name. - repository_name = "zenodo" - # Keep track of script duration. - start_time = time.perf_counter() - # Create context manager. - output_path = output_dir_path / repository_name - output_path.mkdir(parents=True, exist_ok=True) - context = ContextManager( - logger=create_logger(logpath=f"{output_path}/{repository_name}_scraping.log"), - output_path=output_path, - query_file_name=query_file_path, + # Create scraper context. + scraper = ScraperContext( + data_source_name=DatasetSourceName.ZENODO, + output_dir_path=output_dir_path, + query_file_path=query_file_path, + is_in_debug_mode=is_in_debug_mode, ) + logger = create_logger(logpath=scraper.log_file_path, level="INFO") # Log script name and doctring. - context.logger.info(__file__) - context.logger.info(__doc__) + logger.info(__file__) + logger.info(__doc__) # Read and verify Zenodo token. load_dotenv() zenodo_token = os.environ.get("ZENODO_TOKEN", "") if not zenodo_token: - context.logger.critical("No Zenodo token found.") - context.logger.critical("Aborting.") + logger.critical("No Zenodo token found.") + logger.critical("Aborting.") sys.exit(1) else: - context.logger.success("Found Zenodo token.") - context.token = zenodo_token + logger.success("Found Zenodo token.") + scraper.token = zenodo_token # Test connection to Zenodo API. - if is_zenodo_connection_working(context.token, logger=context.logger): - context.logger.success("Connection to Zenodo API successful.") + if is_zenodo_connection_working(scraper.token, logger=logger): + logger.success("Connection to Zenodo API successful.") else: - context.logger.critical("Connection to Zenodo API failed.") - context.logger.critical("Aborting.") + logger.critical("Connection to Zenodo API failed.") + logger.critical("Aborting.") sys.exit(1) # Get rate limit information. get_rate_limit_info( @@ -656,43 +595,60 @@ def main(output_dir_path: Path, query_file_path: Path): "https://zenodo.org/records/4444751/preview/code.zip", ], zenodo_token, - logger=context.logger, + logger=logger, ) # Read parameter file (file_types, keywords, excluded_files, excluded_paths) = read_query_file( - context.query_file_name, - logger=context.logger, + scraper.query_file_path, + logger=logger, ) - # Verify output directory exists - verify_output_directory(context.output_path) - - datasets_df, files_df = search_all_datasets(file_types, keywords, context) - + datasets_metadata, files_metadata = search_all_datasets( + file_types, keywords, scraper, logger=logger + ) + # Normalize datasets and files metadata. + datasets_normalized_metadata = normalize_datasets_metadata( + datasets_metadata, logger=logger + ) + files_normalized_metadata = normalize_files_metadata(files_metadata, logger=logger) # Scrap zip files content. - context.logger.info("-" * 30) - zip_df = scrap_zip_content(files_df, logger=context.logger) - # We don't remove duplicates here because - # one zip file can contain several files with the same name - # but within different folders. - files_df = pd.concat([files_df, zip_df], ignore_index=True) - context.logger.info(f"Number of files found inside zip files: {len(zip_df)}") - context.logger.info(f"Total number of files found: {len(files_df)}") - files_df = remove_excluded_files(files_df, excluded_files, excluded_paths) - context.logger.info("-" * 30) + logger.info("-" * 30) + files_zip_metadata = scrap_zip_content(files_normalized_metadata, logger=logger) + logger.info(f"Number of files found inside zip files: {len(files_zip_metadata)}") + # Normalize files metadata from zip files. + zip_normalized_metadata = normalize_files_metadata( + files_zip_metadata, logger=logger + ) + # Merge all metadata files. + files_normalized_metadata += zip_normalized_metadata + logger.info(f"Total number of files found: {len(files_normalized_metadata)}") + files_normalized_metadata = remove_excluded_files( + files_normalized_metadata, excluded_files, excluded_paths, logger=logger + ) + logger.info("-" * 30) # Remove datasets that contain non-MD related files # that come from zip files. - datasets_df, files_df = find_remove_false_positive_datasets( - datasets_df, files_df, context + datasets_normalized_metadata, files_normalized_metadata = ( + find_remove_false_positive_datasets( + datasets_normalized_metadata, + files_normalized_metadata, + scraper, + logger=logger, + ) ) - - # Save dataframes to disk. - export_dataframe_to_parquet("zenodo", DataType.DATASETS, datasets_df, context) - export_dataframe_to_parquet("zenodo", DataType.FILES, files_df, context) - - # Script duration. - elapsed_time = int(time.perf_counter() - start_time) - context.logger.info(f"Scraped Zenodo in: {timedelta(seconds=elapsed_time)}") + # Save metadata to parquet files. + scraper.number_of_datasets_scraped = export_list_of_models_to_parquet( + scraper.datasets_parquet_file_path, + datasets_normalized_metadata, + logger=logger, + ) + scraper.number_of_files_scraped = export_list_of_models_to_parquet( + scraper.files_parquet_file_path, + files_normalized_metadata, + logger=logger, + ) + # Print scraping statistics. + print_statistics(scraper, logger=logger) if __name__ == "__main__": diff --git a/tests/scrapers/test_figshare.py b/tests/scrapers/test_figshare.py index 7d6adfe..431a9a3 100644 --- a/tests/scrapers/test_figshare.py +++ b/tests/scrapers/test_figshare.py @@ -1,32 +1,9 @@ """Tests for the figshare scraper module.""" -from pathlib import Path - -import loguru -import pytest - -import mdverse_scrapers.core.toolbox as toolbox import mdverse_scrapers.scrapers.figshare as figshare_scraper -@pytest.fixture -def create_context(): - """Create a context for testing. - - Returns - ------- - toolbox.ContextManager - A context manager. - """ - ctx = toolbox.ContextManager( - logger=loguru.logger, - output_path=Path(""), - query_file_name=Path(""), - ) - return ctx - - -def test_extract_files_from_zip_file(create_context): +def test_extract_files_from_zip_file(): """Test the extract_files_from_zip_file function.""" expected_file_names = [ "topologies/martini/betacarotene-CG.itp",