Adding a New Source¶
[!IMPORTANT] Repository Access & Collaboration: The
gatheringcodebase 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:¶
- Empty / Incomplete Files: Rejects files below
min_size_bytes(default 100 bytes). - HTML Error Disguises: Server 401/403/404/500/502 error pages returned with HTTP 200 (e.g. OAuth login redirects).
- 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). - Archive Integrity: Tests gzip/zip central directory headers without full decompression.
- 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:
-
CLI Discovery & Autogenerated Subcommands:
-
Python Task Builder:
-
YAML Configuration Engine:
-
Testing Checklist:
- Create unit tests in
tests/test_my_network_gatherer.py. - Test parameter validation, spatial grid parsing, mock download flow, and
verify_file_integritycache behavior.