Skip to content

Adding a New Source

[!IMPORTANT] Repository Access & Collaboration: The gathering codebase is maintained in a private repository and users do not have direct access to push or open merge requests. If you wish to collaborate on integrating a new atmospheric network, satellite mission, or custom sensor, please contact the DIVA platform administrators. The engineering team will coordinate development, provide sandbox access, or integrate the source directly.

Adding a new data source to gathering requires only one file. Drop a Python module into gathering/sources/, subclass Gatherer, and define a Pydantic parameter model.


1. Quickstart Skeleton

Create a new file, e.g. gathering/sources/my_network.py:

from pathlib import Path
from typing import Any
from pydantic import BaseModel, Field
from gathering.base import Gatherer
from gathering.grid import Grid, parse_grid
from gathering.integrity import verify_file_integrity
from gathering.logging import logger


class MyNetworkParameters(BaseModel):
    station: str = Field(description="Station identifier code")
    product: list[str] = Field(default_factory=lambda: ["default_prod"], description="Products to fetch")
    start_date: str | None = Field(default=None, description="Start date YYYY-MM-DD")
    end_date: str | None = Field(default=None, description="End date YYYY-MM-DD")
    grid: Grid | dict[str, Any] | None = Field(default=None, description="Spatial Region of Interest")
    max_workers: int = Field(default=1, ge=1, description="Concurrent download threads")


class MyNetworkGatherer(Gatherer):
    NAME = "my_network"
    ALIASES = ["mynet", "my_source"]
    PARAMS_MODEL = MyNetworkParameters

    def __init__(self, *, caching_sub_location: Path | str = Path("my_network"), **kwargs: Any):
        super().__init__(caching_sub_location=caching_sub_location, **kwargs)

    def _validate_parameters(self, **kwargs: Any) -> dict[str, Any]:
        """Validate input arguments against the Pydantic model."""
        kwargs.pop("overwrite_cache", None)
        return MyNetworkParameters(**kwargs).model_dump()

    def _set_caching_name(self, **kwargs: Any) -> None:
        """Define the local caching directory or filename."""
        self.caching_name = f"{kwargs['station']}_{kwargs.get('start_date', 'all')}"

    def _call_source_fetch(self, **kwargs: Any) -> list[str]:
        """Search, authenticate, download, and verify files. Return file paths."""
        params = MyNetworkParameters(**kwargs)
        # Download logic here (see sections below for Grid and Integrity details)...
        return []

2. Using Spatial Grid Functions (gathering.grid)

Gathering provides standard spatial Region of Interest (ROI) adapters in gathering.grid. These allow your gatherer to support AERONET site name resolution, bounding box transformations, and polygon geometries.

Declaring Grid in Parameter Models

Include grid and optional bounding_box in your Pydantic model:

from gathering.grid import Grid

class MySatelliteParameters(BaseModel):
    grid: Grid | dict[str, Any] | list[float] | None = Field(
        default=None,
        description="Spatial ROI defined as a Grid instance, site dictionary, or [min_lon, min_lat, max_lon, max_lat].",
    )
    bounding_box: list[float] | None = Field(
        default=None,
        description="Bounding box [min_lon, min_lat, max_lon, max_lat].",
    )

Parsing & Adapting Grid for Remote APIs

Inside _call_source_fetch, use parse_grid to convert user inputs into a structured Grid object:

from gathering.grid import parse_grid, get_aeronet_site_coords

def _call_source_fetch(self, **kwargs: Any) -> list[str]:
    params = MySatelliteParameters(**kwargs)

    # 1. Parse Grid object (resolves AERONET sites, dicts, or bbox lists)
    roi_grid: Grid | None = parse_grid(params.grid) if params.grid else None

    # 2. Convert to STAC / OGC Bounding Box: [min_lon, min_lat, max_lon, max_lat]
    if roi_grid:
        stac_bbox = roi_grid.to_stac_bbox()
        # e.g. [-3.705, 37.064, -3.505, 37.264]

    # 3. Convert to NASA CMR Bounding Box: (min_lon, min_lat, max_lon, max_lat)
    if roi_grid:
        cmr_bbox = roi_grid.to_cmr_bounding_box()
        # Pass directly to cmr_base or search query: search_cmr_granules(..., bounding_box=cmr_bbox)

    # 4. Convert to WKT Polygon (for OData / Copernicus APIs)
    if roi_grid:
        wkt_polygon = roi_grid.to_wkt_polygon()
        # e.g. "POLYGON((-3.705 37.064, -3.505 37.064, ...))"

    # 5. Direct AERONET Site Coordinate Lookup (if needed)
    lat, lon, elev = get_aeronet_site_coords("Granada")
    # Returns (37.164, -3.605, 680.0)

3. File Integrity Verification & Safety (gathering.integrity)

To prevent corrupted downloads, incomplete transfers, and server error pages from polluting the local cache, use the verification utilities in gathering.integrity.

Key Verification Utilities

Function Signature Purpose
verify_file_integrity (path, min_size_bytes=100, check_archive=True, expected_checksum=None, checksum_algo='md5') -> tuple[bool, str | None] Full integrity check: size, magic bytes, archive headers, HTML error payloads, and checksum.
verify_checksum (path, expected_hash, algorithm='md5') -> bool Validates file MD5 or SHA-256 against an expected remote hash.
compute_md5 / compute_sha256 (path) -> str Computes hex digest in memory-efficient chunks.
validate_and_clean_cache (cache_dir, min_size_bytes=100, remove_invalid=True) -> tuple[list[Path], list[Path]] Scans a cache directory, validating existing files and purging corrupt ones.

Robust Atomic Download Pattern

Always write remote downloads to a .tmp file, verify integrity, and atomically replace the destination path:

import requests
from gathering.integrity import verify_file_integrity
from gathering.logging import logger

def _download_and_verify(self, url: str, dest_path: Path, expected_md5: str | None = None) -> Path | None:
    dest_path.parent.mkdir(parents=True, exist_ok=True)

    # 1. Check existing cache
    if dest_path.exists() and not self.overwrite_cache:
        valid, reason = verify_file_integrity(dest_path, expected_checksum=expected_md5)
        if valid:
            logger.debug(f"Cache hit: {dest_path.name}")
            return dest_path
        logger.warning(f"Purging corrupted cached file ({reason}): {dest_path.name}")
        dest_path.unlink(missing_ok=True)

    # 2. Download to temporary file
    temp_path = dest_path.with_suffix(dest_path.suffix + ".tmp")
    try:
        with requests.get(url, stream=True, timeout=60) as resp:
            resp.raise_for_status()
            with open(temp_path, "wb") as f:
                for chunk in resp.iter_content(chunk_size=65536):
                    if chunk:
                        f.write(chunk)

        # 3. Verify downloaded file integrity (magic bytes, size, HTML errors, checksum)
        is_valid, error_msg = verify_file_integrity(
            temp_path,
            min_size_bytes=256,
            check_archive=True,
            expected_checksum=expected_md5,
            checksum_algo="md5",
        )

        if not is_valid:
            logger.error(f"Download verification failed for {dest_path.name}: {error_msg}")
            temp_path.unlink(missing_ok=True)
            return None

        # 4. Atomic move to final cache location
        temp_path.replace(dest_path)
        return dest_path

    except Exception as exc:
        logger.error(f"Download failed for {dest_path.name}: {exc}")
        if temp_path.exists():
            temp_path.unlink(missing_ok=True)
        return None

What verify_file_integrity Automatically Detects:

  1. Empty / Incomplete Files: Rejects files below min_size_bytes (default 100 bytes).
  2. HTML Error Disguises: Server 401/403/404/500/502 error pages returned with HTTP 200 (e.g. OAuth login redirects).
  3. Magic Byte Verification: Checks valid file signatures for NetCDF-3 (CDF\x01), NetCDF-4/HDF5 (\x89HDF\r\n\x1a\n), HDF4, GZIP (\x1f\x8b), BZIP2 (BZh), and ZIP (PK\x03\x04).
  4. Archive Integrity: Tests gzip/zip central directory headers without full decompression.
  5. Cryptographic Checksums: Verifies against server-provided MD5 or SHA-256 digests.

4. Multi-Granule Concurrency (max_workers)

If your data source produces multiple granules or assets per query, support concurrent downloads using ThreadPoolExecutor:

from concurrent.futures import ThreadPoolExecutor, as_completed

def _call_source_fetch(self, **kwargs: Any) -> list[str]:
    params = MyNetworkParameters(**kwargs)
    items_to_download = self._search_remote_items(params)

    if params.max_workers <= 1 or len(items_to_download) == 1:
        # Sequential execution
        return [self._download_item(it, params) for it in items_to_download]

    # Parallel execution
    workers = min(params.max_workers, len(items_to_download))
    downloaded: list[str] = []

    with ThreadPoolExecutor(max_workers=workers) as executor:
        futures = {executor.submit(self._download_item, item, params): item for item in items_to_download}
        for future in as_completed(futures):
            res = future.result()
            if res:
                downloaded.append(str(res))

    return downloaded

5. Automatic Framework Integration

Once your subclass is placed in gathering/sources/, gathering automatically wires it into:

  1. CLI Discovery & Autogenerated Subcommands:

    gathering list                                          # Lists your gatherer
    gathering my_network --station Granada --max-workers 4  # Generated CLI command
    

  2. Python Task Builder:

    from gathering import Task
    
    task = Task.my_network(station="Granada", max_workers=4)
    result = task.run()
    

  3. YAML Configuration Engine:

    tasks:
      - source: my_network
        station: "Granada"
        grid:
          site: "Granada"
          size_yx: [100, 100]
          resolution_m: 1000.0
        max_workers: 4
    

  4. Testing Checklist:

  5. Create unit tests in tests/test_my_network_gatherer.py.
  6. Test parameter validation, spatial grid parsing, mock download flow, and verify_file_integrity cache behavior.