diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f6f60e0..27bb46d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -14,7 +14,7 @@ jobs: working-directory: cli_tools/mcdi strategy: matrix: - python-version: ["3.9", "3.11"] + python-version: ["3.10", "3.12"] steps: - uses: actions/checkout@v4 - uses: actions/setup-python@v5 diff --git a/cli_tools/mcdi/Dockerfile b/cli_tools/mcdi/Dockerfile index 24c16f5..2c8a766 100644 --- a/cli_tools/mcdi/Dockerfile +++ b/cli_tools/mcdi/Dockerfile @@ -1,4 +1,4 @@ -FROM python:3.12-alpine +FROM python:3.12-slim WORKDIR /app # build from the mcdi parent directory COPY . . diff --git a/cli_tools/mcdi/README.md b/cli_tools/mcdi/README.md index dda7dd8..c0c9f57 100644 --- a/cli_tools/mcdi/README.md +++ b/cli_tools/mcdi/README.md @@ -4,10 +4,10 @@ Galaxy Cancer Data Importers (GaCDI) provides Galaxy tools for importing cancer datasets from major public and controlled-access cancer data repositories into Galaxy histories. This package provides one command, `mcdi` (Multi-Commons Data Importer), with two subcommands: `mcdi manifest` builds a manifest from -filters, and `mcdi download` downloads the files a GDC or PDC manifest lists -(whether built here or exported from a portal). +filters, and `mcdi download` downloads the files a GDC, PDC, or IDC manifest +lists (whether built here or exported from a portal). -## Manifest Builder (this branch) +## Manifest Builder `gacdi_manifest_gdc` generates the **manifests** that drive the importers. Instead of downloading a whole dataset, the user filters the NCI @@ -15,7 +15,7 @@ of downloading a whole dataset, the user filters the NCI want, described in two complementary outputs: - **GDC manifest** — strict `id / filename / md5 / size / state`, consumable - directly by `gdc-client` and by the GaCDI GDC importer. + directly by `gdc-client` and by `mcdi download` (the GaCDI Manifest Downloader). - **metadata table** — the same files joined to sample barcodes and, optionally, clinical/molecular annotations (GDC fields, cBioPortal subtypes like PAM50 and ER/PR/HER2, and/or a user-uploaded annotation TSV). It also carries @@ -58,40 +58,34 @@ repeatable custom facets (`--extra-filter "field=…;op=in|exclude;values=a,b"`) a raw GDC filters JSON (`--raw-filters`). The manifest is emitted in a deterministic (sorted) order for reproducible workflows. -### End-to-end with the GaCDI GDC importer +### Feeding the manifest into `mcdi download` -This tool is designed to feed directly into the **GaCDI GDC importer** (the -manifest-download branch), so a single Galaxy workflow goes *filter → manifest → -download → analysis*: +A single Galaxy workflow goes *filter → manifest → download → analysis*: build +the manifest here, run **GaCDI Manifest Downloader** (`mcdi download`, below) +to bring the files into the history, then join `metadata.tsv` to those history +datasets to attach clinical labels, the `galaxy_ext` datatype hint, and subtype +annotations to each sample. -``` -[GaCDI Manifest Builder]--gdc_manifest.txt-->[GaCDI GDC importer]--collection-->[analysis tools] - \--metadata.tsv------------------------------(join)----/ -``` - -**Compatibility contract (locked by `tests/test_importer_contract.py`):** +This works because of a contract between the two tools (locked by +`tests/test_importer_contract.py`): -1. **Manifest → importer.** `gdc_manifest.txt` is a TSV whose header - (`id, filename, md5, size, state`) is a superset of what the importer's - `parse_gdc_manifest` requires (`id/filename/md5/size`); its datatype (`txt`) - is accepted by the importer's manifest input (`tabular,txt`). Rows with no - `id` are dropped so the manifest and metadata stay row-aligned. The same file - also works with `gdc-client download -m gdc_manifest.txt`. +1. **Manifest → downloader.** `gdc_manifest.txt` is a TSV whose header + (`id, filename, md5, size, state`) is a superset of what `mcdi download`'s + GDC parsing requires (`id/filename/md5/size`); its datatype (`txt`) is + accepted by the downloader's manifest input (`tabular,txt`). Rows with no + `id` are dropped so the manifest and metadata stay row-aligned. The same + file also works with `gdc-client download -m gdc_manifest.txt`. 2. **Metadata ↔ history.** `metadata.tsv` leads with `file_id` and `filename` — - the exact keys of the importer's history **summary** — so after download you - join the metadata to the imported collection/summary on `file_id` (stable - UUID) or `filename` to attach clinical labels, the `galaxy_ext` datatype hint, - and subtype annotations to each sample in the history. - -So: build the manifest here, run the importer to bring the samples into the -history, then join `metadata.tsv` to give those history datasets their -annotations (e.g. labels for an image ML model). + the same keys the downloaded collection's datasets are named/identified by + — so after download you join the metadata to the collection on `file_id` + (stable UUID) or `filename` to attach clinical labels, the `galaxy_ext` + datatype hint, and subtype annotations to each sample in the history. ## Downloading files from a manifest -`mcdi download` fetches the files listed in a GDC or PDC manifest — either -one built by `mcdi manifest gdc` above, or one exported directly from a -portal: +`mcdi download` fetches the files listed in a GDC, PDC, or IDC manifest — +either one built by `mcdi manifest gdc` above, or one exported directly from +a portal: - **GDC**: build a file cart in the [GDC portal](https://portal.gdc.cancer.gov) and download the manifest (TSV), or generate one via the API: @@ -100,22 +94,32 @@ portal: to the files you want, and use "Export File Manifest" (CSV or TSV). Note that the signed download links embedded in a PDC manifest expire after 7 days — re-export if downloads start failing. +- **IDC**: build a cohort in the [IDC portal](https://portal.imaging.datacommons.cancer.gov) + and export either manifest format it offers — the **s5cmd manifest** (a + `.s5cmd` script of `cp s3://...` lines, one per DICOM series) or the + **CSV/TSV cohort manifest** (a table keyed by `SeriesInstanceUID`). Both are + auto-detected and accepted. A manifest entry references a whole series (many + files, not one), so the actual file list is resolved at download time + against [idc-index](https://github.com/ImagingDataCommons/idc-index)'s + metadata and the public `idc-open-data` S3 bucket; only AWS-hosted series are + currently supported. No token is needed — all IDC data is open access. ```bash mcdi download --manifest gdc_manifest.txt --output-dir downloads/ mcdi download --manifest pdc_manifest.csv --output-dir downloads/ --verify-checksum +mcdi download --manifest idc_manifest.s5cmd --output-dir downloads/ ``` -The data commons is auto-detected from the manifest's header row; pass -`--source {gdc,pdc}` to override. +The data commons is auto-detected from the manifest's content; pass +`--source {gdc,idc,pdc}` to override. | Flag | Description | |---|---| | `--manifest PATH` | Path to the manifest (required) | | `--output-dir DIR` | Directory to download files into (required) | -| `--source {gdc,pdc}` | Skip auto-detection | +| `--source {gdc,idc,pdc}` | Skip auto-detection | | `--workers N` | Concurrent downloads (default: 4). PDC always runs at 1 to respect its rate limit, regardless of this flag. | -| `--verify-checksum` | Verify each file's md5 against the manifest after download | +| `--verify-checksum` | Verify each file's md5 against the manifest after download. No-op for IDC entries: IDC has no reliable per-object checksum (S3's ETag isn't sound for multipart uploads), so `md5` is always unset and verification is silently skipped, same as any manifest entry without one. | | `-x`, `--extract` | Extract recognized archives (`.tar.gz`, `.tgz`, `.tar.bz2`, `.tar.xz`, `.tar`, `.zip`, `.gz`, `.bz2`, `.xz`) in place after download, keeping only the extracted contents at the manifest's output path. Off by default. | | `--token-file PATH` | File containing a GDC auth token, for controlled-access files | | `--retries N` | Extra attempts for files that fail transiently within this run (default: 2) | @@ -163,11 +167,14 @@ mcdi download --manifest gdc_manifest.txt --output-dir downloads/ ``` PDC downloads use pre-signed URLs embedded in the manifest and need no token. +IDC downloads are open access and also need no token. Output layout: - **GDC**: `downloads/gdc//` - **PDC**: `downloads/pdc////[/]/` +- **IDC**: `downloads/idc////_/` + (mirrors idc-index's own default download hierarchy) Re-running against the same manifest and output directory skips files already downloaded (verifying checksums too, if `--verify-checksum` is set), so @@ -176,33 +183,27 @@ so a rerun over mostly-complete output is cheap, not another full pass. ## Runtime environment -The tool ships a pinned container (`quay.io//mcdi`) referenced from -the wrapper, with Python + `requests` Conda requirements as a fallback. The Quay -namespace (`paulocilasjr`) is a placeholder — update `@QUAY_ORG@` in -`tools/manifest_gdc/macros.xml`, `containers/Dockerfile.manifest`, and the workflow -before publishing. +Both Galaxy tools (`manifest_gdc`, `manifest_downloader`) reference the same +pinned container, `quay.io/goeckslab/mcdi:` — there's no Conda +package, so the container is the sole runtime. It's built from +`cli_tools/mcdi/Dockerfile`, tagged from `mcdi/__init__.py`'s `__version__`, +and built/pushed automatically by `.github/workflows/containers.yml` on +merges to `main` that touch `cli_tools/mcdi/**`. ```bash -docker build -f containers/Dockerfile.manifest -t mcdi:dev . +docker build -t mcdi:dev cli_tools/mcdi docker run --rm mcdi:dev mcdi manifest gdc --help +docker run --rm mcdi:dev mcdi download --help ``` ## Development ```bash python -m pip install -e '.[dev]' -pytest -q # mocked; no network -planemo lint tools/manifest_gdc +pytest -q -m "not network" # mocked; add `-m network` for live-API tests +planemo lint tools/manifest_gdc tools/manifest_downloader ``` -## Roadmap - -- **Phase 1 (this branch):** GDC manifest builder + enrichment + join/QC. -- **Phase 2:** CRDC GDC-style commons (PDC/IDC/ICDC/CDS/CTDC) reusing the filter/ - join core. -- **Phase 3:** GEO/SRA accession-list builders; on merge with the importer branch, - fold shared HTTP utilities into the `gacdi` package. - ## License See [LICENSE](LICENSE). diff --git a/cli_tools/mcdi/mcdi/__init__.py b/cli_tools/mcdi/mcdi/__init__.py index e0192e9..24c15fa 100644 --- a/cli_tools/mcdi/mcdi/__init__.py +++ b/cli_tools/mcdi/mcdi/__init__.py @@ -1,25 +1,14 @@ """MCDI — Multi-Commons Data Importer. -One command, two subcommands: - -- ``mcdi manifest`` (:mod:`mcdi.manifest`): filter-driven generation of - download manifests (and enriched metadata tables) for NIH/NCI cancer data - repositories, starting with the NCI Genomic Data Commons (GDC). The builder - emits a deliberate *two-file split*: a lean, CLI/importer-ready manifest - (``id/filename/md5/size/state``) and a rich metadata table joining - clinical/molecular annotations by barcode, plus a match report so selections - and joins are never silently wrong. -- ``mcdi download`` (:mod:`mcdi.download`): downloads the files listed - in a GDC or PDC manifest, auto-detecting which commons it came from. +``mcdi manifest``: build a filtered GDC manifest + enriched metadata table. +``mcdi download``: download files from a GDC, PDC, or IDC manifest. """ import os -__version__ = "0.4.0" +__version__ = "0.5.0" -# Build identifier baked into the container image at build time (e.g. the git -# commit SHA). Lets you confirm the exact code a run used, even when the version -# number hasn't changed. Empty for local/editable installs. +# Container image build id (e.g. git SHA); empty for local/editable installs. BUILD = os.environ.get("MCDI_BUILD", "").strip() diff --git a/cli_tools/mcdi/mcdi/download/cli.py b/cli_tools/mcdi/mcdi/download/cli.py index 795567b..caed83c 100644 --- a/cli_tools/mcdi/mcdi/download/cli.py +++ b/cli_tools/mcdi/mcdi/download/cli.py @@ -1,4 +1,4 @@ -"""``mcdi download`` — download the files listed in a GDC or PDC manifest.""" +"""``mcdi download`` — download the files listed in a GDC, PDC, or IDC manifest.""" from __future__ import annotations @@ -17,7 +17,7 @@ def add_arguments(subparsers: argparse._SubParsersAction) -> argparse.ArgumentPa """Attach the ``download`` subcommand to ``subparsers``.""" parser = subparsers.add_parser( "download", - help="Download files listed in a GDC or PDC manifest.", + help="Download files listed in a GDC, PDC, or IDC manifest.", ) parser.add_argument("--manifest", required=True, type=Path, help="Path to the exported manifest file") parser.add_argument("--output-dir", required=True, type=Path, help="Directory to download files into") diff --git a/cli_tools/mcdi/mcdi/download/engine.py b/cli_tools/mcdi/mcdi/download/engine.py index 188b5d7..fe9b279 100644 --- a/cli_tools/mcdi/mcdi/download/engine.py +++ b/cli_tools/mcdi/mcdi/download/engine.py @@ -71,9 +71,8 @@ def _dest_path(output_dir: Path, entry: FileEntry) -> Path: def _already_present(dest_path: Path, entry: FileEntry, verify: bool) -> tuple[bool, str]: """Return ``(satisfied, detail)``: whether ``dest_path`` already correctly holds ``entry``. - Shared by the download step and the pre-flight access check, so a file - that doesn't need (re-)downloading also doesn't need its remote - accessibility re-verified on every rerun. + Shared by the download step and the pre-flight access check, so a rerun + skips both re-downloading and re-verifying access for files already present. """ if not (dest_path.exists() and dest_path.stat().st_size > 0): return False, "" @@ -94,19 +93,10 @@ class DownloadResult: def _archived_path(output_dir: Path, entry: FileEntry) -> Path: - """Where ``entry``'s archive is relocated to once successfully extracted. - - Mirrors ``entry.rel_dir``/filename under a sibling of ``output_dir`` - (``.mcdi-archives/...``). Moving the archive there - instead - of leaving it in ``output_dir`` next to what it was extracted into - - means (a) a tool that recursively collects everything under - ``output_dir`` (e.g. a Galaxy ``discover_datasets`` with ``recurse``) - only ever sees the extracted contents, not a redundant copy of the - packed archive, and (b) the archive's presence here doubles as the - idempotency marker: extraction is already done for this entry iff a file - exists here. The archive isn't deleted, just moved aside - if extraction - fails, it's left in ``output_dir`` instead, so there's still something - to show for the download. + """Where ``entry``'s archive moves to after successful extraction: a sibling + ``.mcdi-archives/`` tree, mirroring ``entry.rel_dir``. + + Its presence there doubles as the idempotency marker for "already extracted". """ archive_root = output_dir.with_name(output_dir.name + ".mcdi-archives") return archive_root / entry.rel_dir / entry.filename @@ -187,15 +177,12 @@ def _download_one( return result -# HTTP statuses worth retrying at the batch level: rate limiting and -# server-side/transient failures. Anything else (401/403/404/...) is a -# permanent-looking failure that a delayed retry won't fix. +# Transient: rate limiting + server errors. Not 401/403/404 (permanent). _RETRYABLE_STATUSES = {429, 500, 502, 503, 504} def _is_retryable(result: DownloadResult) -> bool: - # A checksum mismatch could be transient transfer corruption, not - # necessarily a bad source file, so it's worth one more attempt too. + # A checksum mismatch could be transient transfer corruption, not a bad source file. if result.status == "checksum_mismatch": return True if result.status != "error": @@ -206,8 +193,7 @@ def _is_retryable(result: DownloadResult) -> bool: return int(detail.split()[1]) in _RETRYABLE_STATUSES except (IndexError, ValueError): return False - # Any other "error" detail came from a requests.RequestException (timeout, - # connection reset, DNS hiccup, ...) - inherently transient. + # Any other "error" detail is a requests.RequestException (timeout, connection reset, ...) - transient. return True @@ -224,8 +210,7 @@ def run( session = build_session() rate_limit = source.rate_limit() pacer = Pacer(rate_limit) if rate_limit else None - # PDC's per-IP rate limit applies regardless of thread count, so cap - # effective concurrency to 1 when a pacer is active to keep pacing honest. + # Cap concurrency to 1 when paced (e.g. PDC's per-IP limit), regardless of --workers. effective_workers = 1 if pacer else workers def _pass(batch: list[FileEntry]) -> list[DownloadResult]: @@ -248,10 +233,7 @@ def _pass(batch: list[FileEntry]) -> list[DownloadResult]: results_by_id = {r.entry.file_id: r for r in _pass(entries)} - # Batch-level retry, on top of the per-request transport retry already - # inside `build_session()`. This matters most for non-interactive runs - # (e.g. a Galaxy job) where nothing will manually rerun the command on - # the same output directory if a few files fail transiently. + # Batch-level retry, on top of build_session()'s own transport-level retries. attempt = 0 while attempt < retries: retry_entries = [r.entry for r in results_by_id.values() if _is_retryable(r)] @@ -302,17 +284,8 @@ def check_access( ) -> list[AccessFailure]: """Probe every entry with a 1-byte ranged request; return the ones that fail. - Meant to run before any real download, so a manifest with one - inaccessible file (e.g. controlled-access without a valid token) is - caught in seconds instead of after downloading everything else first. - Two things narrow what actually needs a round trip: entries already - correctly present in ``output_dir`` (same check the download step uses - - including, if ``extract`` is set, entries already extracted and relocated - to ``output_dir``'s sibling archive directory - so a rerun over - mostly-complete output doesn't re-verify remote access for files it isn't - going to touch anyway) skip the check outright, and entries the source - can already vouch for as open (see ``Source.known_open``) skip the - per-file probe specifically. + Skips entries already present in ``output_dir`` (or already extracted, if + ``extract``) and ones ``Source.known_open`` vouches for. """ session = build_session() rate_limit = source.rate_limit() diff --git a/cli_tools/mcdi/mcdi/download/sources/__init__.py b/cli_tools/mcdi/mcdi/download/sources/__init__.py index 2ff1de1..7711056 100644 --- a/cli_tools/mcdi/mcdi/download/sources/__init__.py +++ b/cli_tools/mcdi/mcdi/download/sources/__init__.py @@ -3,24 +3,28 @@ from pathlib import Path from ...errors import InputError -from .base import CANDIDATE_DELIMITERS, FileEntry, RateLimit, Source, read_header +from .base import CANDIDATE_DELIMITERS, FileEntry, RateLimit, Source, read_header, read_lines from .gdc import GDCSource +from .idc import IDCSource from .pdc import PDCSource SOURCES: dict[str, type[Source]] = { GDCSource.name: GDCSource, PDCSource.name: PDCSource, + IDCSource.name: IDCSource, } def detect_source(path: Path) -> str: """Identify which data commons a manifest came from by matching its header row, under each candidate delimiter, against every registered source's - schema (``Source.sniff``).""" + schema (``Source.sniff``), or - for headerless formats - its leading raw + lines (``Source.sniff_lines``).""" + lines = read_lines(path) matches = [ name for name, cls in SOURCES.items() - if any(cls.sniff(read_header(path, d)) for d in CANDIDATE_DELIMITERS) + if any(cls.sniff(read_header(path, d)) for d in CANDIDATE_DELIMITERS) or cls.sniff_lines(lines) ] if len(matches) == 1: return matches[0] @@ -34,4 +38,6 @@ def detect_source(path: Path) -> str: ) -__all__ = ["FileEntry", "RateLimit", "Source", "GDCSource", "PDCSource", "SOURCES", "detect_source"] +__all__ = [ + "FileEntry", "RateLimit", "Source", "GDCSource", "PDCSource", "IDCSource", "SOURCES", "detect_source", +] diff --git a/cli_tools/mcdi/mcdi/download/sources/base.py b/cli_tools/mcdi/mcdi/download/sources/base.py index eb133b9..92bb39f 100644 --- a/cli_tools/mcdi/mcdi/download/sources/base.py +++ b/cli_tools/mcdi/mcdi/download/sources/base.py @@ -12,6 +12,9 @@ # source's schema. CANDIDATE_DELIMITERS = ("\t", ",") +# Leading raw lines read for Source.sniff_lines (past a manifest's leading comments). +CONTENT_SNIFF_LINES = 20 + def read_header(path: Path, delimiter: str) -> list[str]: """Read and split a manifest's header row using the given delimiter.""" @@ -19,6 +22,12 @@ def read_header(path: Path, delimiter: str) -> list[str]: return next(csv.reader(f, delimiter=delimiter), []) +def read_lines(path: Path, limit: int = CONTENT_SNIFF_LINES) -> list[str]: + """Read up to ``limit`` leading raw lines, for ``Source.sniff_lines``.""" + with open(path) as f: + return [line for line, _ in zip(f, range(limit))] + + @dataclass class FileEntry: """A single file to download, normalized across manifest formats.""" @@ -41,6 +50,11 @@ class Source(ABC): def sniff(header_fields: list[str]) -> bool: """Return True if a manifest with these header fields belongs to this source.""" + @staticmethod + def sniff_lines(lines: list[str]) -> bool: + """Hook for headerless formats ``sniff`` can't identify. Default: no match.""" + return False + @abstractmethod def parse_manifest(self, path: Path) -> list[FileEntry]: """Read a manifest file and return the files it lists.""" @@ -54,12 +68,8 @@ def rate_limit(self) -> Optional["RateLimit"]: return None def known_open(self, entries: list[FileEntry], session: requests.Session) -> set[str]: - """Best-effort ``file_id``s known accessible without a per-file network probe. - - Lets a source short-circuit the pre-flight access check (see - ``download.engine.check_access``) for entries it can already vouch - for cheaply, e.g. via a single bulk metadata query. Default: none, - so every entry gets individually probed. + """``file_id``s known accessible without a per-file probe (see + ``download.engine.check_access``). Default: none. """ return set() diff --git a/cli_tools/mcdi/mcdi/download/sources/gdc.py b/cli_tools/mcdi/mcdi/download/sources/gdc.py index cad5438..af90850 100644 --- a/cli_tools/mcdi/mcdi/download/sources/gdc.py +++ b/cli_tools/mcdi/mcdi/download/sources/gdc.py @@ -54,11 +54,7 @@ def request_kwargs(self, entry: FileEntry) -> dict: return {"headers": headers} def known_open(self, entries: list[FileEntry], session: requests.Session) -> set[str]: - """Bulk-query GDC's ``access`` field for every id; no token needed for this. - - This just narrows which entries the per-file pre-flight probe has to - touch — open files never need the individual round trip. - """ + """Bulk-query GDC's ``access`` field for every id (no token needed).""" ids = [e.file_id for e in entries] if not ids: return set() diff --git a/cli_tools/mcdi/mcdi/download/sources/idc.py b/cli_tools/mcdi/mcdi/download/sources/idc.py new file mode 100644 index 0000000..324ebfd --- /dev/null +++ b/cli_tools/mcdi/mcdi/download/sources/idc.py @@ -0,0 +1,226 @@ +from __future__ import annotations + +import csv +import re +import xml.etree.ElementTree as ET +from concurrent.futures import ThreadPoolExecutor, as_completed +from dataclasses import dataclass +from pathlib import Path +from typing import Optional + +import requests +from idc_index import IDCClient + +from ...errors import ApiError, InputError +from ...net import build_session +from .base import CANDIDATE_DELIMITERS, FileEntry, Source, read_header + +# Columns in an IDC cohort manifest (CSV/TSV): series-level (one row per +# series) or BigQuery-export (one row per SOPInstanceUID, same columns). +REQUIRED_HEADERS = {"SeriesInstanceUID", "collection_id", "PatientID", "gcs_url"} + +# Matches ``cp :////*``. Scheme-agnostic: the +# bucket always comes from the index (see IdcIndexLookup), never the manifest +# text. Non-matching lines (comments, blanks, stray headers) are skipped, not errors. +_S5CMD_CP_RE = re.compile(r"^cp\s+\S+://[^\s/]+/([^\s/]+)/\*(?:\s|$)") + +# Only bucket confirmed anonymously listable/gettable; GCS is out of scope for now. +_SUPPORTED_BUCKET = "idc-open-data" + +# Fixed concurrency for the per-series S3 listing fan-out (not plumbed through --workers). +_LISTING_WORKERS = 8 + +_S3_NS = {"s3": "http://s3.amazonaws.com/doc/2006-03-01/"} + + +@dataclass +class SeriesInfo: + """One series' location, resolved via idc-index's metadata index.""" + + series_uid: str + crdc_series_uuid: str + aws_bucket: str + collection_id: str + patient_id: str + study_uid: str + modality: str + + +class IdcIndexLookup: + """Wraps ``idc_index.IDCClient``'s bundled metadata index. + + Constructor-injection seam for unit testing without a real ``IDCClient`` + (needs ``s5cmd`` on ``PATH``, loads a ~77MB index). + """ + + def __init__(self) -> None: + self._client: Optional[IDCClient] = None + + def _get_client(self) -> IDCClient: + if self._client is None: + self._client = IDCClient() + return self._client + + def resolve(self, column: str, values: list[str]) -> tuple[dict[str, SeriesInfo], list[str]]: + """Batch-resolve values against the current index, then prior releases for misses. + + Returns ``(resolved_by_value, still_unresolved)``. Vectorized, not per-value + (the index can hold upwards of a million rows). + """ + client = self._get_client() + resolved: dict[str, SeriesInfo] = {} + remaining = list(values) + for df in (client.index, client.prior_versions_index): + if not remaining: + break + matches = df[df[column].isin(remaining)].drop_duplicates(subset=[column]) + for _, row in matches.iterrows(): + resolved[row[column]] = SeriesInfo( + series_uid=row["SeriesInstanceUID"], + crdc_series_uuid=row["crdc_series_uuid"], + aws_bucket=row["aws_bucket"], + collection_id=row["collection_id"], + patient_id=row["PatientID"], + study_uid=row["StudyInstanceUID"], + modality=row["Modality"], + ) + remaining = [v for v in remaining if v not in resolved] + return resolved, remaining + + +def _list_series_objects(session: requests.Session, bucket: str, crdc_series_uuid: str) -> list[tuple[str, int]]: + """Anonymous, paginated ListObjectsV2 over one series' object prefix + (public-read bucket, no signing needed).""" + prefix = f"{crdc_series_uuid}/" + keys: list[tuple[str, int]] = [] + token: Optional[str] = None + while True: + params = {"list-type": "2", "prefix": prefix} + if token: + params["continuation-token"] = token + resp = session.get(f"https://{bucket}.s3.amazonaws.com", params=params, timeout=60) + if resp.status_code != 200: + raise ApiError(f"Could not list IDC series {crdc_series_uuid!r} objects: HTTP {resp.status_code}") + root = ET.fromstring(resp.content) + for c in root.findall("s3:Contents", _S3_NS): + key = c.findtext("s3:Key", namespaces=_S3_NS) + size = int(c.findtext("s3:Size", default="0", namespaces=_S3_NS)) + keys.append((key, size)) + token = root.findtext("s3:NextContinuationToken", namespaces=_S3_NS) + is_truncated = root.findtext("s3:IsTruncated", namespaces=_S3_NS) + if is_truncated != "true" or not token: + break + if not keys: + raise ApiError( + f"IDC series {crdc_series_uuid!r} resolved via idc-index but has no objects " + f"at s3://{bucket}/{prefix}" + ) + return keys + + +def _csv_delimiter(path: Path) -> Optional[str]: + """Delimiter the manifest's header matches ``REQUIRED_HEADERS`` under, or None.""" + for d in CANDIDATE_DELIMITERS: + if REQUIRED_HEADERS.issubset(set(read_header(path, d))): + return d + return None + + +def _extract_csv_series_uids(path: Path, delimiter: str) -> list[str]: + seen: list[str] = [] + seen_set: set[str] = set() + with open(path, newline="") as f: + for row in csv.DictReader(f, delimiter=delimiter): + series_uid = (row.get("SeriesInstanceUID") or "").strip() + if series_uid and series_uid not in seen_set: + seen_set.add(series_uid) + seen.append(series_uid) + return seen + + +def _extract_s5cmd_series_uuids(path: Path) -> list[str]: + seen: list[str] = [] + seen_set: set[str] = set() + with open(path) as f: + for raw in f: + m = _S5CMD_CP_RE.match(raw.strip()) + if not m: + continue # comment, blank, or an unrecognized/header line - skip, don't error + crdc_series_uuid = m.group(1) + if crdc_series_uuid not in seen_set: + seen_set.add(crdc_series_uuid) + seen.append(crdc_series_uuid) + return seen + + +def _extract_series_refs(path: Path) -> tuple[str, list[str]]: + """(index_column, ordered_deduplicated_ids) a manifest references.""" + delimiter = _csv_delimiter(path) + if delimiter is not None: + return "SeriesInstanceUID", _extract_csv_series_uids(path, delimiter) + return "crdc_series_uuid", _extract_s5cmd_series_uuids(path) + + +class IDCSource(Source): + name = "idc" + + def __init__(self, lookup: Optional[IdcIndexLookup] = None): + self._lookup = lookup or IdcIndexLookup() + + @staticmethod + def sniff(header_fields: list[str]) -> bool: + return REQUIRED_HEADERS.issubset(set(header_fields)) + + @staticmethod + def sniff_lines(lines: list[str]) -> bool: + return any(_S5CMD_CP_RE.match(line.strip()) for line in lines) + + def parse_manifest(self, path: Path) -> list[FileEntry]: + column, ids = _extract_series_refs(path) + if not ids: + raise InputError(f"IDC manifest {path} contained no resolvable series references.") + + resolved, unresolved = self._lookup.resolve(column, ids) + if unresolved: + raise InputError( + f"{len(unresolved)} IDC series not found in idc-index metadata (current or prior " + f"releases); manifest may be stale: {unresolved}" + ) + + session = build_session() + entries: list[FileEntry] = [] + with ThreadPoolExecutor(max_workers=min(_LISTING_WORKERS, len(resolved))) as pool: + futures = [pool.submit(self._entries_for_series, session, info) for info in resolved.values()] + for future in as_completed(futures): + entries.extend(future.result()) + return entries + + def _entries_for_series(self, session: requests.Session, info: SeriesInfo) -> list[FileEntry]: + if info.aws_bucket != _SUPPORTED_BUCKET: + raise ApiError( + f"IDC series {info.series_uid!r} resolves to unsupported bucket {info.aws_bucket!r} " + f"(only {_SUPPORTED_BUCKET!r} is currently supported)." + ) + rel_dir = ( + Path("idc") / info.collection_id / info.patient_id / info.study_uid + / f"{info.modality}_{info.series_uid}" + ) + return [ + FileEntry( + file_id=key, + filename=key.rsplit("/", 1)[-1], + rel_dir=rel_dir, + url=f"https://{info.aws_bucket}.s3.amazonaws.com/{key}", + size=size, + md5=None, # no reliable per-object checksum (S3 ETag unsound for multipart uploads) + ) + for key, size in _list_series_objects(session, info.aws_bucket, info.crdc_series_uuid) + ] + + def request_kwargs(self, entry: FileEntry) -> dict: + return {} + + def known_open(self, entries: list[FileEntry], session: requests.Session) -> set[str]: + """Already confirmed to exist via the listing in ``parse_manifest`` - skip the + redundant per-file probe in ``engine.check_access``.""" + return {e.file_id for e in entries} diff --git a/cli_tools/mcdi/mcdi/download/sources/pdc.py b/cli_tools/mcdi/mcdi/download/sources/pdc.py index 7f01c11..20f7e41 100644 --- a/cli_tools/mcdi/mcdi/download/sources/pdc.py +++ b/cli_tools/mcdi/mcdi/download/sources/pdc.py @@ -60,6 +60,5 @@ def request_kwargs(self, entry: FileEntry) -> dict: return {} def rate_limit(self) -> RateLimit: - # Mirrors NCI's reference download script: pace requests to avoid - # tripping PDC's 24h per-IP restriction on repeated file downloads. + # Avoids tripping PDC's 24h per-IP restriction on repeated file downloads. return RateLimit(max_per_window=10, window_seconds=600, per_file_sleep_seconds=2) diff --git a/cli_tools/mcdi/mcdi/manifest/cbioportal.py b/cli_tools/mcdi/mcdi/manifest/cbioportal.py index 44a4ba2..a825bae 100644 --- a/cli_tools/mcdi/mcdi/manifest/cbioportal.py +++ b/cli_tools/mcdi/mcdi/manifest/cbioportal.py @@ -1,16 +1,8 @@ """cBioPortal client: list clinical attributes and fetch per-sample values. -Attribute ids (e.g. ``SUBTYPE`` for PAM50, ``ER_STATUS_BY_IHC``) are -study-specific, so we do not hard-map them: the caller supplies the study id and -an optional attribute list (or ``all``). Sample ids look like ``TCGA-XX-XXXX-01``. - -Clinical attributes live at two levels in cBioPortal: - -- SAMPLE level (e.g. ANEUPLOIDY_SCORE, CANCER_TYPE, MSI scores), and -- PATIENT level (e.g. SUBTYPE/PAM50, ER_STATUS_BY_IHC, PR_STATUS_BY_IHC, HER2). - -We fetch **both** and merge the patient-level values onto every sample of that -patient, so subtype/receptor-status columns are included. +Attribute ids (e.g. ``SUBTYPE``, ``ER_STATUS_BY_IHC``) are study-specific, not +hard-mapped. Clinical data exists at SAMPLE and PATIENT level; both are fetched +and patient-level values are merged onto each of that patient's samples. """ from __future__ import annotations @@ -57,9 +49,7 @@ def fetch_clinical( ) -> tuple[dict[str, dict], list[str]]: """Return ``{sample_id: {attr: value}}`` and the ordered attribute columns. - Fetches SAMPLE- and PATIENT-level clinical data and merges the patient values - onto each of that patient's samples. If *attribute_ids* is given, only those - attributes are kept (order preserved). + If *attribute_ids* is given, only those attributes are kept (order preserved). """ sample_records = _get_clinical(session, study_id, "SAMPLE", base) patient_records = _get_clinical(session, study_id, "PATIENT", base) diff --git a/cli_tools/mcdi/mcdi/manifest/cli.py b/cli_tools/mcdi/mcdi/manifest/cli.py index b73d37c..ad58b1b 100644 --- a/cli_tools/mcdi/mcdi/manifest/cli.py +++ b/cli_tools/mcdi/mcdi/manifest/cli.py @@ -85,9 +85,7 @@ def add_arguments(subparsers: argparse._SubParsersAction) -> argparse.ArgumentPa def _run_gdc(args: argparse.Namespace) -> int: session = build_session() - # Normalise and sanity-check the cBioPortal study id(s) early, before the - # (expensive) GDC query, so a bad id fails fast with a clear message. - # One or more comma-separated ids are allowed (they get merged). + # Validate cBioPortal study id(s) before the (expensive) GDC query, so bad input fails fast. if args.cbioportal_study: studies = enrich.split_studies(args.cbioportal_study) bad = [s for s in studies if any(c.isspace() for c in s)] @@ -144,17 +142,14 @@ def _run_gdc(args: argparse.Namespace) -> int: total_matching = gdc.count(session, filters) file_rows = gdc.query_files(session, filters, max_files=args.max_files, total=total_matching) - # Drop rows without a file id: the GaCDI GDC importer skips empty-id manifest - # rows, so excluding them keeps the manifest and metadata table aligned. + # Drop rows without a file id so the manifest and metadata stay row-aligned. dropped = [r for r in file_rows if not r.file_id] if dropped: log.warning("Dropped %d file(s) with no file id.", len(dropped)) file_rows = [r for r in file_rows if r.file_id] - # Deterministic order so the manifest is reproducible across runs (workflows). + # Deterministic order. file_rows.sort(key=lambda r: r.file_id) - # No files matched: write empty, self-explanatory outputs and skip the - # (now pointless) annotation fetch/join. if not file_rows: io.write_manifest(args.manifest_out, []) io.write_metadata(args.metadata_out, [], []) diff --git a/cli_tools/mcdi/mcdi/manifest/enrich.py b/cli_tools/mcdi/mcdi/manifest/enrich.py index 84aaef0..49ffa3b 100644 --- a/cli_tools/mcdi/mcdi/manifest/enrich.py +++ b/cli_tools/mcdi/mcdi/manifest/enrich.py @@ -1,8 +1,6 @@ """Collect optional annotation sources into a single ``{sample_id: {attr: val}}``. -Sources (any combination): cBioPortal sample clinical data and a user-uploaded -annotation TSV. GDC-native fields (barcode, sample_type, project) are always -carried through by the join and need no collection here. +Sources: cBioPortal sample clinical data and a user-uploaded annotation TSV. """ from __future__ import annotations @@ -54,10 +52,7 @@ def collect( ) -> tuple[dict[str, dict], list[str]]: """Gather and merge all enrichment sources into one annotation table. - ``cbioportal_study`` may name several studies (comma-separated); they are - fetched and merged in order. Columns are the union across sources; for a given - sample/attribute the first source with a non-empty value wins (so list the - highest-priority study first). + Merged in order; first source with a non-empty value wins per attribute. """ annotations: dict[str, dict] = {} columns: list[str] = [] diff --git a/cli_tools/mcdi/mcdi/manifest/filters.py b/cli_tools/mcdi/mcdi/manifest/filters.py index 95bcbcc..eb47b6a 100644 --- a/cli_tools/mcdi/mcdi/manifest/filters.py +++ b/cli_tools/mcdi/mcdi/manifest/filters.py @@ -1,15 +1,11 @@ -"""Build a GDC ``filters`` object from guided flags, custom facets and raw JSON. - -Owning filter construction in Python (rather than the Galaxy Cheetah template) -keeps the wrapper trivial and lets us unit-test the exact query we send. -""" +"""Build a GDC ``filters`` object from guided flags, custom facets and raw JSON.""" from __future__ import annotations from ..errors import InputError -# Guided flag -> GDC field. File-level fields are unprefixed on the files -# endpoint; case-level fields use the ``cases.`` path. +# Guided flag -> GDC field: unprefixed (file-level), ``cases.`` (case-level), or +# ``analysis.`` (pipeline-level, e.g. workflow_type). GUIDED_FIELDS = { "project": "cases.project.project_id", "primary_site": "cases.primary_site", diff --git a/cli_tools/mcdi/mcdi/manifest/io.py b/cli_tools/mcdi/mcdi/manifest/io.py index 6c0d419..5fc9d21 100644 --- a/cli_tools/mcdi/mcdi/manifest/io.py +++ b/cli_tools/mcdi/mcdi/manifest/io.py @@ -34,7 +34,8 @@ def _top_counts(rows: list[dict], column: str, top: int = 15) -> tuple[list[tupl counter = Counter((r.get(column) or "(blank)") for r in rows) return counter.most_common(top), len(counter) -# gdc-client / gacdi_gdc require exactly these columns, in this order. +# gdc-client requires exactly these columns, in this order; mcdi download only +# needs id/filename/md5/size present (any order, extra columns OK). MANIFEST_COLUMNS = ["id", "filename", "md5", "size", "state"] BASE_METADATA_COLUMNS = [ @@ -90,12 +91,8 @@ def write_report( max_examples: int = 25, extra: list[tuple[str, str, str]] | None = None, ) -> None: - """Write a meaningful tabular report: ``category key value`` rows. - - Sections: a **summary** (match counts, download size, annotation coverage), a - **composition** breakdown of what the manifest contains, per-attribute - annotation **coverage**, GDC **facets** (preview mode), and a short capped list - of **unmatched** samples — instead of dumping every file UUID. + """Write a tabular report (``category key value`` rows): summary, composition, + annotation coverage, GDC facets (preview mode), and capped unmatched examples. """ rows: list[tuple[str, str, str]] = [] @@ -125,7 +122,7 @@ def add(category: str, key: str, value) -> None: if report.collisions: add("summary", "annotation_key_collisions", len(report.collisions)) - # --- composition: what am I about to download? --------------------- + # --- composition ----------------------------------------------------- if merged_rows: for column, label in _COMPOSITION_COLUMNS: counts, distinct = _top_counts(merged_rows, column) diff --git a/cli_tools/mcdi/mcdi/manifest/join.py b/cli_tools/mcdi/mcdi/manifest/join.py index f26be3e..f78c3b4 100644 --- a/cli_tools/mcdi/mcdi/manifest/join.py +++ b/cli_tools/mcdi/mcdi/manifest/join.py @@ -1,8 +1,6 @@ """Barcode normalization and the manifest<->annotation join (with QC report). -The join is *left* on the manifest: every selected file is kept, matched or not, -so downloads are never silently dropped. A report records match counts, unmatched -files, unused annotations and colliding keys. +Left join on the manifest: every selected file is kept, matched or not. """ from __future__ import annotations diff --git a/cli_tools/mcdi/mcdi/manifest/model.py b/cli_tools/mcdi/mcdi/manifest/model.py index 54b0a1f..fb5aa90 100644 --- a/cli_tools/mcdi/mcdi/manifest/model.py +++ b/cli_tools/mcdi/mcdi/manifest/model.py @@ -52,9 +52,7 @@ def disease_type(row: dict) -> str | None: return _first_match(row, _DISEASE_TYPE) or row.get("cases.disease_type") or None -# Best-effort mapping of a downloaded file to a Galaxy datatype extension, so the -# manifest/metadata can tell downstream workflow tools how to interpret each file. -# Filename suffix wins (most reliable); GDC data_format is the fallback. +# Best-effort file -> Galaxy datatype extension. Filename suffix wins; data_format is the fallback. _EXT_BY_SUFFIX = [ (".bam", "bam"), (".bai", "bai"), diff --git a/cli_tools/mcdi/mcdi/net.py b/cli_tools/mcdi/mcdi/net.py index b4c3fd6..c69cff8 100644 --- a/cli_tools/mcdi/mcdi/net.py +++ b/cli_tools/mcdi/mcdi/net.py @@ -1,9 +1,4 @@ -"""Minimal retrying HTTP session. - -TEMPORARY: mirrors ``gacdi.net`` from the NIH_commons branch. When the branches -merge, replace this module with an import of ``gacdi.net.build_session`` to avoid -duplication. -""" +"""Minimal retrying HTTP session.""" from __future__ import annotations diff --git a/cli_tools/mcdi/pyproject.toml b/cli_tools/mcdi/pyproject.toml index c74fd3e..ce4eb0c 100644 --- a/cli_tools/mcdi/pyproject.toml +++ b/cli_tools/mcdi/pyproject.toml @@ -5,14 +5,15 @@ build-backend = "hatchling.build" [project] name = "mcdi" dynamic = ["version"] -description = "MCDI (Multi-Commons Data Importer) — filter-driven manifest + enriched metadata generation for NIH/NCI cancer data repositories, plus a GDC/PDC manifest downloader, via `mcdi manifest` and `mcdi download`." +description = "MCDI (Multi-Commons Data Importer) — filter-driven manifest + enriched metadata generation for NIH/NCI cancer data repositories, plus a GDC/PDC/IDC manifest downloader, via `mcdi manifest` and `mcdi download`." readme = "README.md" -requires-python = ">=3.9" +requires-python = ">=3.10" license = { file = "LICENSE" } authors = [{ name = "GaCDI contributors" }] -keywords = ["galaxy", "cancer", "gdc", "manifest", "cbioportal", "tcga", "nci"] +keywords = ["galaxy", "cancer", "gdc", "manifest", "cbioportal", "tcga", "nci", "idc"] dependencies = [ "requests>=2.28", + "idc-index>=0.12.0,<0.13.0", ] [project.optional-dependencies] diff --git a/cli_tools/mcdi/tests/test_download_idc.py b/cli_tools/mcdi/tests/test_download_idc.py new file mode 100644 index 0000000..dc67589 --- /dev/null +++ b/cli_tools/mcdi/tests/test_download_idc.py @@ -0,0 +1,254 @@ +from pathlib import Path + +import pytest + +from mcdi.download.sources.idc import ( + IDCSource, + SeriesInfo, + _extract_series_refs, +) +from mcdi.errors import ApiError, InputError + +SERIES_A = SeriesInfo( + series_uid="1.2.series.A", + crdc_series_uuid="aaaaaaaa-0000-0000-0000-000000000000", + aws_bucket="idc-open-data", + collection_id="coll", + patient_id="pat1", + study_uid="1.2.study.A", + modality="MR", +) +SERIES_B = SeriesInfo( + series_uid="1.2.series.B", + crdc_series_uuid="bbbbbbbb-1111-1111-1111-111111111111", + aws_bucket="idc-open-data", + collection_id="coll", + patient_id="pat2", + study_uid="1.2.study.B", + modality="CT", +) + +CSV_HEADER = "PatientID,collection_id,StudyInstanceUID,SeriesInstanceUID,source_DOI,study_uuid,series_uuid,gcs_url\n" + + +class FakeLookup: + """Duck-compatible stand-in for IdcIndexLookup.resolve(), keyed by id.""" + + def __init__(self, table: dict): + self._table = table + + def resolve(self, column, values): + resolved = {v: self._table[v] for v in values if v in self._table} + unresolved = [v for v in values if v not in self._table] + return resolved, unresolved + + +def _xml(keys_sizes, next_token=None, truncated=False): + contents = "".join(f"{k}{s}" for k, s in keys_sizes) + trunc = "true" if truncated else "false" + token = f"{next_token}" if next_token else "" + return ( + '' + '' + f"{contents}{trunc}{token}" + ) + + +def _mock_listing(requests_mock, bucket, series_pages): + """series_pages: {crdc_series_uuid: [(keys_sizes, is_truncated), ...]}.""" + + def callback(request, context): + prefix = request.qs["prefix"][0] + crdc_series_uuid = prefix.rstrip("/") + pages = series_pages[crdc_series_uuid] + token = request.qs.get("continuation-token", [None])[0] + idx = int(token) if token else 0 + keys_sizes, truncated = pages[idx] + next_token = str(idx + 1) if truncated else None + context.status_code = 200 + return _xml(keys_sizes, next_token=next_token, truncated=truncated) + + requests_mock.get(f"https://{bucket}.s3.amazonaws.com/", text=callback) + + +# --- Phase 1: pure reference extraction ------------------------------------------------- + + +def test_extract_series_refs_csv_dedups_across_instance_rows(tmp_path): + manifest = tmp_path / "cohort.csv" + manifest.write_text( + CSV_HEADER + + "pat1,coll,1.2.study.A,1.2.series.A,,,,gs://x\n" + + "pat1,coll,1.2.study.A,1.2.series.A,,,,gs://x\n" # BigQuery/instance-level: repeated series row + + "pat2,coll,1.2.study.B,1.2.series.B,,,,gs://y\n" + ) + column, ids = _extract_series_refs(manifest) + assert column == "SeriesInstanceUID" + assert ids == ["1.2.series.A", "1.2.series.B"] + + +def test_extract_series_refs_s5cmd_tolerates_comments_and_stray_header(tmp_path): + manifest = tmp_path / "m.s5cmd" + manifest.write_text( + "# generated by the IDC portal\n" + "study_manifest_cp_command\n" # a stray non-cp header row some exports include + "\n" + "cp s3://idc-open-data/aaaaaaaa-0000-0000-0000-000000000000/* .\n" + "cp s3://idc-open-data/aaaaaaaa-0000-0000-0000-000000000000/* .\n" # dup, should collapse + "cp gs://idc-open-data/bbbbbbbb-1111-1111-1111-111111111111/* .\n" # gs:// scheme still matches + ) + column, ids = _extract_series_refs(manifest) + assert column == "crdc_series_uuid" + assert ids == [ + "aaaaaaaa-0000-0000-0000-000000000000", + "bbbbbbbb-1111-1111-1111-111111111111", + ] + + +# --- sniff / sniff_lines ----------------------------------------------------------------- + + +def test_sniff_matches_csv_cohort_header(): + header = CSV_HEADER.strip().split(",") + assert IDCSource.sniff(header) + assert not IDCSource.sniff(["id", "filename", "md5", "size", "state"]) + + +def test_sniff_lines_matches_s5cmd_and_ignores_other_content(): + assert IDCSource.sniff_lines(["# comment\n", "cp s3://idc-open-data/uuid1/* .\n"]) + assert not IDCSource.sniff_lines(["id\tfilename\tmd5\tsize\tstate\n"]) + + +# --- parse_manifest: resolution errors ---------------------------------------------------- + + +def test_parse_manifest_raises_one_error_naming_all_unresolved(tmp_path): + manifest = tmp_path / "m.s5cmd" + manifest.write_text( + "cp s3://idc-open-data/aaaaaaaa-0000-0000-0000-000000000000/* .\n" + "cp s3://idc-open-data/cccccccc-2222-2222-2222-222222222222/* .\n" + "cp s3://idc-open-data/dddddddd-3333-3333-3333-333333333333/* .\n" + ) + lookup = FakeLookup({SERIES_A.crdc_series_uuid: SERIES_A}) + with pytest.raises(InputError) as exc_info: + IDCSource(lookup=lookup).parse_manifest(manifest) + message = str(exc_info.value) + assert "cccccccc-2222-2222-2222-222222222222" in message + assert "dddddddd-3333-3333-3333-333333333333" in message + + +def test_parse_manifest_empty_manifest_raises_input_error(tmp_path): + manifest = tmp_path / "m.s5cmd" + manifest.write_text("# nothing here but comments\n\n") + with pytest.raises(InputError): + IDCSource(lookup=FakeLookup({})).parse_manifest(manifest) + + +# --- parse_manifest: successful resolution + listing -------------------------------------- + + +def test_parse_manifest_builds_entries_from_listing(tmp_path, requests_mock): + manifest = tmp_path / "m.s5cmd" + manifest.write_text(f"cp s3://idc-open-data/{SERIES_A.crdc_series_uuid}/* .\n") + _mock_listing( + requests_mock, + "idc-open-data", + {SERIES_A.crdc_series_uuid: [([(f"{SERIES_A.crdc_series_uuid}/a.dcm", 100)], False)]}, + ) + + entries = IDCSource(lookup=FakeLookup({SERIES_A.crdc_series_uuid: SERIES_A})).parse_manifest(manifest) + + assert len(entries) == 1 + entry = entries[0] + assert entry.filename == "a.dcm" + assert entry.size == 100 + assert entry.md5 is None + assert entry.url == f"https://idc-open-data.s3.amazonaws.com/{SERIES_A.crdc_series_uuid}/a.dcm" + assert entry.rel_dir == Path("idc/coll/pat1/1.2.study.A/MR_1.2.series.A") + + +def test_parse_manifest_handles_pagination(tmp_path, requests_mock): + manifest = tmp_path / "m.s5cmd" + manifest.write_text(f"cp s3://idc-open-data/{SERIES_A.crdc_series_uuid}/* .\n") + uuid = SERIES_A.crdc_series_uuid + _mock_listing( + requests_mock, + "idc-open-data", + { + uuid: [ + ([(f"{uuid}/a.dcm", 10)], True), + ([(f"{uuid}/b.dcm", 20)], False), + ] + }, + ) + + entries = IDCSource(lookup=FakeLookup({uuid: SERIES_A})).parse_manifest(manifest) + + assert {e.filename for e in entries} == {"a.dcm", "b.dcm"} + + +def test_parse_manifest_fans_out_across_multiple_series_concurrently(tmp_path, requests_mock): + manifest = tmp_path / "m.s5cmd" + manifest.write_text( + f"cp s3://idc-open-data/{SERIES_A.crdc_series_uuid}/* .\n" + f"cp s3://idc-open-data/{SERIES_B.crdc_series_uuid}/* .\n" + ) + _mock_listing( + requests_mock, + "idc-open-data", + { + SERIES_A.crdc_series_uuid: [([(f"{SERIES_A.crdc_series_uuid}/a.dcm", 10)], False)], + SERIES_B.crdc_series_uuid: [([(f"{SERIES_B.crdc_series_uuid}/b.dcm", 20)], False)], + }, + ) + lookup = FakeLookup({SERIES_A.crdc_series_uuid: SERIES_A, SERIES_B.crdc_series_uuid: SERIES_B}) + + entries = IDCSource(lookup=lookup).parse_manifest(manifest) + + assert {e.filename for e in entries} == {"a.dcm", "b.dcm"} + assert {e.rel_dir for e in entries} == { + Path("idc/coll/pat1/1.2.study.A/MR_1.2.series.A"), + Path("idc/coll/pat2/1.2.study.B/CT_1.2.series.B"), + } + + +# --- parse_manifest: listing-side errors --------------------------------------------------- + + +def test_parse_manifest_unsupported_bucket_raises_api_error(tmp_path): + manifest = tmp_path / "m.s5cmd" + manifest.write_text(f"cp s3://idc-open-data/{SERIES_A.crdc_series_uuid}/* .\n") + weird_bucket_series = SeriesInfo(**{**SERIES_A.__dict__, "aws_bucket": "some-other-bucket"}) + lookup = FakeLookup({SERIES_A.crdc_series_uuid: weird_bucket_series}) + with pytest.raises(ApiError): + IDCSource(lookup=lookup).parse_manifest(manifest) + + +def test_parse_manifest_empty_listing_raises_api_error(tmp_path, requests_mock): + manifest = tmp_path / "m.s5cmd" + manifest.write_text(f"cp s3://idc-open-data/{SERIES_A.crdc_series_uuid}/* .\n") + _mock_listing(requests_mock, "idc-open-data", {SERIES_A.crdc_series_uuid: [([], False)]}) + lookup = FakeLookup({SERIES_A.crdc_series_uuid: SERIES_A}) + with pytest.raises(ApiError): + IDCSource(lookup=lookup).parse_manifest(manifest) + + +# --- known_open ------------------------------------------------------------------------ + + +def test_known_open_vouches_for_every_entry(tmp_path, requests_mock): + manifest = tmp_path / "m.s5cmd" + manifest.write_text(f"cp s3://idc-open-data/{SERIES_A.crdc_series_uuid}/* .\n") + _mock_listing( + requests_mock, + "idc-open-data", + {SERIES_A.crdc_series_uuid: [([(f"{SERIES_A.crdc_series_uuid}/a.dcm", 10)], False)]}, + ) + source = IDCSource(lookup=FakeLookup({SERIES_A.crdc_series_uuid: SERIES_A})) + entries = source.parse_manifest(manifest) + + assert source.known_open(entries, session=None) == {e.file_id for e in entries} + + +def test_request_kwargs_is_empty(): + assert IDCSource().request_kwargs(entry=None) == {} diff --git a/cli_tools/mcdi/tests/test_download_idc_network.py b/cli_tools/mcdi/tests/test_download_idc_network.py new file mode 100644 index 0000000..7bbf0d7 --- /dev/null +++ b/cli_tools/mcdi/tests/test_download_idc_network.py @@ -0,0 +1,67 @@ +from __future__ import annotations + +from pathlib import Path + +import pytest + +from mcdi.download import engine +from mcdi.download.sources import detect_source +from mcdi.download.sources.idc import IDCSource + +# A small (3 objects, ~105KB total), stable IDC series - itself one of +# idc-index's own maintained test fixtures - used purely to exercise the +# download pipeline against the real idc-index metadata index and the real +# public idc-open-data S3 bucket. +CRDC_SERIES_UUID = "28621ba9-1aca-4aab-a2a1-f6d2c3e2ab19" +SERIES_INSTANCE_UID = "1.3.6.1.4.1.14519.5.2.1.7695.1700.153974929648969296590126728101" +EXPECTED_FILE_COUNT = 3 +EXPECTED_TOTAL_SIZE = 104904 + + +def _write_s5cmd_manifest(path: Path) -> None: + path.write_text( + "# generated by the IDC portal\n" + f"cp s3://idc-open-data/{CRDC_SERIES_UUID}/* .\n" + ) + + +def _write_csv_manifest(path: Path) -> None: + path.write_text( + "PatientID,collection_id,StudyInstanceUID,SeriesInstanceUID,source_DOI,study_uuid,series_uuid,gcs_url\n" + f"ISPY1_1181,ispy1,1.3.6.1.4.1.14519.5.2.1.7695.1700.251181729149149196962664129666," + f"{SERIES_INSTANCE_UID},,,,gs://placeholder\n" + ) + + +@pytest.mark.network +@pytest.mark.parametrize("write_manifest,expected_source", [ + (_write_s5cmd_manifest, "idc"), + (_write_csv_manifest, "idc"), +]) +def test_idc_manifest_round_trip(tmp_path, write_manifest, expected_source): + manifest_path = tmp_path / "idc_manifest.txt" + write_manifest(manifest_path) + + assert detect_source(manifest_path) == expected_source + + source = IDCSource() + entries = source.parse_manifest(manifest_path) + assert len(entries) == EXPECTED_FILE_COUNT + assert sum(e.size for e in entries) == EXPECTED_TOTAL_SIZE + + output_dir = tmp_path / "downloads" + failures = engine.check_access(entries, source, output_dir=output_dir, workers=4) + assert failures == [] + + results = engine.run(entries, source, output_dir, workers=4, verify=True) + assert {r.status for r in results} == {"downloaded"} + for entry in entries: + dest = output_dir / entry.rel_dir / entry.filename + assert dest.is_file() + assert dest.stat().st_size == entry.size + + # Re-running against the same output dir should skip files already + # downloaded rather than re-fetching them (no checksum to verify against + # for IDC entries, but presence + size is still enough to skip). + results_again = engine.run(entries, source, output_dir, workers=4, verify=True) + assert {r.status for r in results_again} == {"skipped"} diff --git a/cli_tools/mcdi/tests/test_download_sources.py b/cli_tools/mcdi/tests/test_download_sources.py index deb580b..e679ab8 100644 --- a/cli_tools/mcdi/tests/test_download_sources.py +++ b/cli_tools/mcdi/tests/test_download_sources.py @@ -84,6 +84,20 @@ def test_detect_source_ignores_filename_extension(tmp_path): assert detect_source(pdc_manifest) == "pdc" +def test_gdc_pdc_sniff_lines_default_false_and_unaffected_by_content_hook(): + # GDC/PDC don't override sniff_lines - the default no-match behavior + # must hold regardless of what content is passed, so adding the hook for + # IDC can't change their detection. + assert not GDCSource.sniff_lines(["cp s3://idc-open-data/uuid1/* .\n"]) + assert not PDCSource.sniff_lines(["cp s3://idc-open-data/uuid1/* .\n"]) + + +def test_detect_source_matches_headerless_manifest_via_sniff_lines(tmp_path): + manifest = tmp_path / "m.s5cmd" + manifest.write_text("# IDC portal export\ncp s3://idc-open-data/uuid1/* .\n") + assert detect_source(manifest) == "idc" + + def test_pdc_parse_manifest_ignores_filename_extension(tmp_path): manifest = tmp_path / "dataset_3.dat" manifest.write_text( diff --git a/cli_tools/mcdi/tests/test_importer_contract.py b/cli_tools/mcdi/tests/test_importer_contract.py index 10dd684..40b2670 100644 --- a/cli_tools/mcdi/tests/test_importer_contract.py +++ b/cli_tools/mcdi/tests/test_importer_contract.py @@ -1,9 +1,8 @@ -"""Compatibility contract with the GaCDI GDC importer (NIH_commons branch). +"""Compatibility contract between ``mcdi manifest gdc`` and ``mcdi download``. -These tests encode what the importer's ``parse_gdc_manifest`` requires and how its -history ``summary`` is keyed, so a change here that would break the download -hand-off fails loudly. Kept self-contained (the importer package is not on this -branch) by replicating the importer's exact expectations. +These tests encode what ``mcdi download``'s GDC parsing (``parse_gdc_manifest``) +requires and how a downloaded collection's datasets are keyed, so a change to the +manifest/metadata builder that would break the download hand-off fails loudly. """ from __future__ import annotations @@ -13,14 +12,14 @@ from mcdi.cli import main from mcdi.manifest.io import BASE_METADATA_COLUMNS, MANIFEST_COLUMNS -# --- mirrored from NIH_commons: gacdi/manifest.py and gacdi/history.py --- +# --- mirrors mcdi.download.sources.gdc.GDCSource's manifest expectations --- IMPORTER_REQUIRED_MANIFEST_COLUMNS = {"id", "filename", "md5", "size"} IMPORTER_SUMMARY_KEYS = {"file_id", "filename"} # how downloaded files are identified IMPORTER_MANIFEST_INPUT_FORMATS = {"tabular", "txt"} def _simulate_importer_parse(path): - """Replica of gacdi.manifest.parse_gdc_manifest's contract.""" + """Replica of mcdi.download.sources.gdc.GDCSource.parse_manifest's contract.""" with open(path, newline="") as fh: reader = csv.DictReader(fh, delimiter="\t") header = {c.strip().lower() for c in (reader.fieldnames or [])} diff --git a/tools/manifest_downloader/macros.xml b/tools/manifest_downloader/macros.xml index d63750b..929f619 100644 --- a/tools/manifest_downloader/macros.xml +++ b/tools/manifest_downloader/macros.xml @@ -1,13 +1,8 @@ - 0.4.0 + 0.5.0 0 25.1 - quay.io/goeckslab/mcdi:@TOOL_VERSION@ diff --git a/tools/manifest_downloader/manifest_downloader.xml b/tools/manifest_downloader/manifest_downloader.xml index 1f00158..cc27ede 100644 --- a/tools/manifest_downloader/manifest_downloader.xml +++ b/tools/manifest_downloader/manifest_downloader.xml @@ -6,7 +6,7 @@ hidden="false" profile="@PROFILE@" > - Download files listed in a GDC or PDC manifest + Download files listed in a GDC, PDC, or IDC manifest macros.xml @@ -22,7 +22,7 @@ name="manifest_file" type="data" format="txt" - label="Manifest file (GDC or PDC — the commons is auto-detected)" + label="Manifest file (GDC, PDC, or IDC — the commons is auto-detected)" /> - @@ -74,15 +71,56 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + `_ or the -`Proteomic Data Commons (PDC) `_, using -``mcdi download``. The manifest's data commons is auto-detected from its -header row, so the same tool handles either kind of manifest without any +Downloads the files listed in a manifest exported from the NCI +`Genomic Data Commons (GDC) `_, the +`Proteomic Data Commons (PDC) `_, or the +`Imaging Data Commons (IDC) `_, +using ``mcdi download``. The manifest's data commons is auto-detected from its +content, so the same tool handles any of these manifest kinds without any extra configuration. **Manifest** @@ -94,6 +132,14 @@ extra configuration. to the files you want and use "Export File Manifest" (CSV or TSV). PDC manifests embed signed download links that expire after 7 days — re-export if a run starts failing on files that previously worked. +- **IDC**: build a cohort in the `IDC portal `_ + and export either manifest format it offers — the **s5cmd manifest** + (a ``.s5cmd`` script of ``cp s3://...`` lines, one per series) or the + **CSV/TSV cohort manifest** (a table keyed by ``SeriesInstanceUID``) — both + are accepted. A manifest entry is a whole DICOM series (many files), so the + actual file list is resolved at run time against the public + `idc-open-data` bucket; only AWS-hosted series are currently supported. No + token is needed — all IDC data is open access. **Authorization token (optional, GDC only)** diff --git a/tools/manifest_downloader/test-data/idc_manifest.csv b/tools/manifest_downloader/test-data/idc_manifest.csv new file mode 100644 index 0000000..9d5ac93 --- /dev/null +++ b/tools/manifest_downloader/test-data/idc_manifest.csv @@ -0,0 +1,2 @@ +PatientID,collection_id,StudyInstanceUID,SeriesInstanceUID,source_DOI,study_uuid,series_uuid,gcs_url +ISPY1_1181,ispy1,1.3.6.1.4.1.14519.5.2.1.7695.1700.251181729149149196962664129666,1.3.6.1.4.1.14519.5.2.1.7695.1700.153974929648969296590126728101,,,,gs://idc-open-data/28621ba9-1aca-4aab-a2a1-f6d2c3e2ab19 diff --git a/tools/manifest_downloader/test-data/idc_manifest.s5cmd b/tools/manifest_downloader/test-data/idc_manifest.s5cmd new file mode 100644 index 0000000..0ce8a40 --- /dev/null +++ b/tools/manifest_downloader/test-data/idc_manifest.s5cmd @@ -0,0 +1,4 @@ +# To download the files in this manifest, install s5cmd +# (https://github.com/peak/s5cmd) and run: +# s5cmd --no-sign-request run idc_manifest.s5cmd +cp s3://idc-open-data/28621ba9-1aca-4aab-a2a1-f6d2c3e2ab19/* . diff --git a/tools/manifest_gdc/gacdi_manifest_gdc.xml b/tools/manifest_gdc/gacdi_manifest_gdc.xml index 00d88e5..42feb0c 100644 --- a/tools/manifest_gdc/gacdi_manifest_gdc.xml +++ b/tools/manifest_gdc/gacdi_manifest_gdc.xml @@ -240,7 +240,6 @@ -
diff --git a/tools/manifest_gdc/macros.xml b/tools/manifest_gdc/macros.xml index fb6752b..82d0627 100644 --- a/tools/manifest_gdc/macros.xml +++ b/tools/manifest_gdc/macros.xml @@ -3,11 +3,6 @@ 0 22.05 - quay.io/goeckslab/mcdi:@TOOL_VERSION@