diff --git a/CITATION.cff b/CITATION.cff index 2fd96de..f4926f9 100644 --- a/CITATION.cff +++ b/CITATION.cff @@ -5,7 +5,9 @@ authors: given-names: "Sérgio Souza" orcid: "https://orcid.org/0000-0002-0232-4549" title: "DisSCube: Declarative Spatial Layer for Dynamic Models" -version: 0.2.0 +abstract: "An open-source, declarative engine for constructing cellular spatial data cubes for dynamic modeling, spatial simulation, and environmental analysis." +version: 0.3.0 +date-released: "2026-09-29" url: "https://github.com/DisSModel/disscube" repository-code: "https://github.com/DisSModel/disscube" license: MIT @@ -14,4 +16,6 @@ keywords: - "geospatial" - "land use and land cover change" - "brazil data cube" + - "cellular automata" - "reproducibility" + - "fair-software" diff --git a/MANIFEST.in b/MANIFEST.in index c1a7121..dcd190b 100644 --- a/MANIFEST.in +++ b/MANIFEST.in @@ -1,2 +1,3 @@ include LICENSE include README.md +recursive-include disscube/data * diff --git a/README.md b/README.md index 56bf79d..beb0738 100644 --- a/README.md +++ b/README.md @@ -234,19 +234,17 @@ data/derived/{grid_id}/{partition}/{spec_hash}/{variable_name}.zarr ``` disscube/ -├── client/ CubeClient — public entry point -├── models/ GridSpec, SpatialSource, SpatialDerivation, Variable… -├── derivation.py Declarative Derivation (front end over SpatialDerivation) +├── client.py CubeClient — public entry point +├── models/ GridSpec, SpatialSource, SpatialDerivation, Variable, Derivation… ├── operators/ Operators as classes (self-registered via __init_subclass__) │ ├── base.py Operator ABC + OPERATOR_REGISTRY │ ├── zonal.py mean, sum, majority, percentage, attribute, presence… │ └── proximity.py distance, min_distance, count -├── pipeline/ Stages: Normalizer → GridAligner → Aggregator → Writer +├── pipeline/ Pipeline execution & planning (schema, runner) + internal stages ├── catalog/ CatalogStore (Protocol) + SQLite and JSON implementations -├── storage/ AssetStore (fsspec — local and S3) +├── storage.py AssetStore (fsspec — local and S3) ├── api/ Experimental HTTP API (optional `api` extra) -├── config/ Pipeline files (TOML): schema, planning, running -├── cli.py `disscube validate` / `disscube run` +├── cli.py `disscube validate` / `disscube run` / `disscube export` ├── sources/ Adapters that bring external data in as SpatialSources, │ │ each with a checksum and a provenance.json sidecar │ ├── _raster.py Window2D, windowed reads, composites, mosaics, register_raster @@ -255,8 +253,7 @@ disscube/ │ ├── mapbiomas.py MapBiomas annual land-cover maps (Collection 11, 10 m series) │ ├── prodes.py PRODES deforestation (download + cache, legend from the .qml) │ └── classified.py any classified map with its legend, e.g. from SITS -└── utils/ Grids (grids.py), BDC tile geometry (bdc_importer.py), - checksums (files.py) +└── utils.py Checksums (sha256_file) and BDC tile importer (import_bdc_grids) ``` ## Adding a new operator @@ -319,6 +316,20 @@ The `purity_threshold` field on `Derivation` is included in `spec_hash` but is n **STAC: reading only** `disscube.sources.bdc` reads Brazil Data Cube cubes through their STAC catalog (search, windowed reads, per-tile composites, mosaics) and writes local GeoTIFFs that are registered as ordinary sources. Derived variables are not published back as STAC, and the `valid_from`/`valid_until` and `bbox` fields on `Derivation` only follow STAC naming conventions. +## Citation + +If you use DisSCube in your research, dynamic modeling, or spatial data pipelines, please cite it using the metadata from [`CITATION.cff`](CITATION.cff) or the following BibTeX entry: + +```bibtex +@software{costa_disscube_2026, + author = {Costa, S{\'e}rgio Souza}, + title = {{DisSCube: Declarative Spatial Layer for Dynamic Models}}, + year = {2026}, + version = {0.3.0}, + url = {https://github.com/DisSModel/disscube} +} +``` + ## License DisSCube is part of the DisSModel ecosystem and is released under the MIT License. See [LICENSE](LICENSE) for details. diff --git a/disscube/__init__.py b/disscube/__init__.py index 0ea5f20..54092e5 100644 --- a/disscube/__init__.py +++ b/disscube/__init__.py @@ -1,5 +1,19 @@ from disscube.client import CubeClient -from disscube.derivation import Derivation -from disscube.models import DerivedVariable, GridSpec, SpatialDerivation, SpatialSource, Variable +from disscube.models import ( + Derivation, + DerivedVariable, + GridSpec, + SpatialDerivation, + SpatialSource, + Variable, +) -__all__ = ["CubeClient", "Derivation", "DerivedVariable", "GridSpec", "SpatialDerivation", "SpatialSource", "Variable"] +__all__ = [ + "CubeClient", + "Derivation", + "DerivedVariable", + "GridSpec", + "SpatialDerivation", + "SpatialSource", + "Variable", +] diff --git a/disscube/api/app.py b/disscube/api/app.py index ce92de6..c3b9b52 100644 --- a/disscube/api/app.py +++ b/disscube/api/app.py @@ -30,10 +30,13 @@ "The DisSCube HTTP API requires the optional 'api' extra: pip install \"disscube[api]\"" ) from exc +import os + from disscube.client import CubeClient from disscube.models import DerivedVariable, GridSpec, SpatialDerivation, SpatialSource -from . import config +DEFAULT_CATALOG_PATH = os.getenv("DISSCUBE_CATALOG", "./catalog.db") +DEFAULT_STORE_PATH = os.getenv("DISSCUBE_STORE", "./data/") def get_cube(request: Request) -> CubeClient: @@ -48,7 +51,7 @@ def create_app(catalog_path: str | None = None, store_path: str | None = None) - """ Build the API application. - Parameters default to ``config.CATALOG_PATH`` / ``config.STORE_PATH``. + Parameters default to ``DEFAULT_CATALOG_PATH`` / ``DEFAULT_STORE_PATH``. The ``CubeClient`` is created in the application lifespan, so building the app has no side effects on disk. """ @@ -56,8 +59,8 @@ def create_app(catalog_path: str | None = None, store_path: str | None = None) - @asynccontextmanager async def lifespan(app: FastAPI) -> AsyncIterator[None]: app.state.cube = CubeClient( - catalog_path or config.CATALOG_PATH, - store_path or config.STORE_PATH, + catalog_path or DEFAULT_CATALOG_PATH, + store_path or DEFAULT_STORE_PATH, ) yield diff --git a/disscube/api/config.py b/disscube/api/config.py deleted file mode 100644 index c51459e..0000000 --- a/disscube/api/config.py +++ /dev/null @@ -1,4 +0,0 @@ -import os - -CATALOG_PATH = os.getenv("DISSCUBE_CATALOG", "./catalog.db") -STORE_PATH = os.getenv("DISSCUBE_STORE", "./data/") diff --git a/disscube/cli.py b/disscube/cli.py index 03e83f2..c0fa654 100644 --- a/disscube/cli.py +++ b/disscube/cli.py @@ -1,10 +1,10 @@ """ Command line: run DisSCube pipeline files. - disscube validate pipeline.toml + disscube validate pipeline.toml [--json] disscube fetch pipeline.toml - disscube run pipeline.toml [--workspace DIR] [--output OUT.tif] [-v] - disscube export pipeline.toml --output OUT.tif [--workspace DIR] [--variables ...] [-v] + disscube run pipeline.toml [--workspace DIR] [--output OUT.tif] [--dry-run] [--json] [-v] + disscube export pipeline.toml --output OUT.tif [--workspace DIR] [--variables ...] [--json] [-v] """ from __future__ import annotations @@ -20,6 +20,7 @@ def main(argv: list[str] | None = None) -> int: p_val = sub.add_parser("validate", help="check a pipeline file without fetching anything") p_val.add_argument("file") + p_val.add_argument("--json", action="store_true", help="output validation result as JSON") p_fetch = sub.add_parser("fetch", help="download and verify remote file sources declared in the pipeline") p_fetch.add_argument("file") @@ -29,6 +30,8 @@ def main(argv: list[str] | None = None) -> int: p_run.add_argument("--workspace", help="output folder (default: the file's 'workspace', " "else a folder named after the file)") p_run.add_argument("--output", "-o", help="export derived variables to a multi-band GeoTIFF") + p_run.add_argument("--dry-run", action="store_true", help="simulate plan execution without downloading or computing") + p_run.add_argument("--json", action="store_true", help="output execution report as JSON") p_run.add_argument("-v", "--verbose", action="store_true", help="log each step") p_exp = sub.add_parser("export", help="export derived variables from an existing data cube to GeoTIFF") @@ -36,20 +39,22 @@ def main(argv: list[str] | None = None) -> int: p_exp.add_argument("--output", "-o", required=True, help="output GeoTIFF file path (e.g. data/cellspace.tif)") p_exp.add_argument("--workspace", help="workspace folder (default: data/cube or from pipeline)") p_exp.add_argument("--variables", nargs="*", help="specific variables to export (default: all derived)") + p_exp.add_argument("--json", action="store_true", help="output export result as JSON") p_exp.add_argument("-v", "--verbose", action="store_true", help="log each step") args = parser.parse_args(argv) - from disscube.config import plan, run - from disscube.config.runner import PipelineError + from disscube.pipeline import PipelineError, plan, run logging.basicConfig(level=logging.INFO if getattr(args, "verbose", False) else logging.WARNING, format="%(levelname)s %(name)s: %(message)s") + is_json = getattr(args, "json", False) try: p = plan(args.file) - print(p.summary()) + if not is_json: + print(p.summary()) if args.command == "fetch": - from disscube.config.runner import _fetch_file_source, _local_file, _raw_cache_dir, _resolve - from disscube.config.schema import FileSource + from disscube.pipeline.runner import _fetch_file_source, _local_file, _raw_cache_dir, _resolve + from disscube.pipeline.schema import FileSource raw = _raw_cache_dir() raw.mkdir(parents=True, exist_ok=True) fetched = 0 @@ -63,18 +68,69 @@ def main(argv: list[str] | None = None) -> int: print(f"Fetch completed: {fetched} remote sources verified.") return 0 if args.command == "validate": - print("OK") + if is_json: + import json + print(json.dumps({ + "status": "ok", + "file": str(p.file.path), + "name": p.file.config.name, + "grid": p.grid.name if p.grid else None, + "sources": [s.id for s in p.sources], + "derives": [d.target for d in p.derives], + }, indent=2)) + else: + print("OK") + return 0 + if getattr(args, "dry_run", False): + if is_json: + import json + print(json.dumps({ + "status": "ok", + "dry_run": True, + "file": str(p.file.path), + "grid": p.grid.name if p.grid else None, + "sources": [{"id": s.id, "type": s.config.type} for s in p.sources], + "derives": [{"target": d.target, "source": d.source, "operator": d.operator} for d in p.derives], + }, indent=2)) + else: + print(f"[dry-run] Plan is valid. Would process {len(p.sources)} sources and derive {len(p.derives)} variables.") return 0 if args.command == "export": - from disscube.config.runner import export_cube + from disscube.pipeline import export_cube exp = export_cube(p, output=args.output, workspace=args.workspace, variables=args.variables) - print(f"workspace : {exp.workspace}") - print(f"exported : {len(exp.variables)} variables on {exp.grid_id}") - print(f"output : {exp.output}") + if is_json: + import json + print(json.dumps({ + "status": "ok", + "workspace": str(exp.workspace), + "grid_id": exp.grid_id, + "variables": exp.variables, + "output": str(exp.output), + }, indent=2)) + else: + print(f"workspace : {exp.workspace}") + print(f"exported : {len(exp.variables)} variables on {exp.grid_id}") + print(f"output : {exp.output}") return 0 report = run(p, workspace=args.workspace, export_geotiff=getattr(args, "output", None)) + if is_json: + import json + print(json.dumps({ + "status": "ok", + "workspace": str(report.workspace), + "grid_id": report.grid_id, + "sources": report.sources, + "derived": report.derived, + "exported": str(report.exported) if report.exported else None, + "record": str(report.record), + }, indent=2, default=str)) + return 0 except PipelineError as exc: - print(f"error: {exc}", file=sys.stderr) + if is_json: + import json + print(json.dumps({"status": "error", "error": str(exc)}, indent=2), file=sys.stderr) + else: + print(f"error: {exc}", file=sys.stderr) return 2 print(f"workspace : {report.workspace}") print(f"derived : {len(report.derived)} products on {report.grid_id}") diff --git a/disscube/client/cube_client.py b/disscube/client.py similarity index 97% rename from disscube/client/cube_client.py rename to disscube/client.py index 84422b8..cc1dfb6 100644 --- a/disscube/client/cube_client.py +++ b/disscube/client.py @@ -2,6 +2,7 @@ import logging import os +import sys from typing import TYPE_CHECKING import numpy as np @@ -21,9 +22,9 @@ # circular import; dissmodel is only needed by to_lucc_data(). from dissmodel.geo.raster.backend import RasterBackend - from disscube.derivation import Derivation + from disscube.models import Derivation -log = logging.getLogger(__name__) +log = logging.getLogger("disscube.client.cube_client") class CubeClient: @@ -198,7 +199,7 @@ def _exists(d: DerivedVariable) -> bool: return ok temporal = [d for d in matches if d.times and _exists(d)] - static = [d for d in matches if not d.times and _exists(d)] + static = [d for d in matches if not d.times and _exists(d)] if temporal: # Stack temporal slices along time axis sorted by first time value @@ -363,3 +364,9 @@ def to_data_source(self, derived_id: str) -> dict: "checksum": derived.content_hash, "type": "local" if derived.asset_url.startswith("/") else "s3" } + + +# Backward-compatibility alias for legacy code importing disscube.client.cube_client +sys.modules[f"{__name__}.cube_client"] = sys.modules[__name__] + +__all__ = ["CubeClient"] diff --git a/disscube/client/__init__.py b/disscube/client/__init__.py deleted file mode 100644 index 794f39a..0000000 --- a/disscube/client/__init__.py +++ /dev/null @@ -1,3 +0,0 @@ -from .cube_client import CubeClient - -__all__ = ["CubeClient"] diff --git a/disscube/config/__init__.py b/disscube/config/__init__.py deleted file mode 100644 index 2b38bbe..0000000 --- a/disscube/config/__init__.py +++ /dev/null @@ -1,15 +0,0 @@ -""" -Pipeline files: declare a whole data preparation — grid, sources, derived -variables — in one TOML file, the DisSCube counterpart of a TerraME -``fillCellularSpace`` script. - - disscube validate pipeline.toml # check the file, no downloads - disscube run pipeline.toml # fetch the sources and derive the variables - -See ``docs/guides/pipeline_files.md`` for the format. -""" - -from disscube.config.runner import Plan, RunReport, load, plan, run -from disscube.config.schema import PipelineConfig - -__all__ = ["PipelineConfig", "Plan", "RunReport", "load", "plan", "run"] diff --git a/disscube/models/__init__.py b/disscube/models/__init__.py index a0cf0a9..eb78ead 100644 --- a/disscube/models/__init__.py +++ b/disscube/models/__init__.py @@ -1,12 +1,28 @@ -from .grid import GridAnchor, GridSpec, SpatialRelation +from .derivation import Derivation +from .grid import ( + BDC_CRS, + BRAZIL_BBOX, + SIMULATION_GRIDS, + GridAnchor, + GridSpec, + SpatialRelation, + register_local_grid, + register_simulation_grids, +) from .variable import DerivedVariable, SpatialDerivation, SpatialSource, Variable __all__ = [ + "BDC_CRS", + "BRAZIL_BBOX", + "SIMULATION_GRIDS", + "Derivation", "DerivedVariable", "GridAnchor", "GridSpec", "SpatialDerivation", "SpatialRelation", "SpatialSource", - "Variable" + "Variable", + "register_local_grid", + "register_simulation_grids", ] diff --git a/disscube/derivation.py b/disscube/models/derivation.py similarity index 89% rename from disscube/derivation.py rename to disscube/models/derivation.py index 379b633..b469925 100644 --- a/disscube/derivation.py +++ b/disscube/models/derivation.py @@ -7,12 +7,14 @@ ``end_datetime`` (and ``datetime`` for the static/instant case). No STAC logic is implemented here; the alignment is naming-only. - ``bbox`` (optional, reserved) follows the STAC bounding-box field order: - [xmin, ymin, xmax, ymax] in EPSG:4326. The field is not populated + [xmin, ymin, xmax, ymax] in EPSG:4326. The field is not populated from data and is not used in execution. No STAC code, catalog, API, or export is implemented in this module. """ +from __future__ import annotations + import hashlib import json from typing import Any, Literal @@ -28,7 +30,7 @@ class Derivation(BaseModel): Declarative description of a single derivation intent. Acts as a thin, additive front-end over the existing - ``SpatialDerivation`` / ``Variable`` machinery. Instantiation validates + ``SpatialDerivation`` / ``Variable`` machinery. Instantiation validates the operator name and operator-specific field requirements so that errors surface before any I/O (fail-fast). @@ -49,18 +51,18 @@ class Derivation(BaseModel): Defaults to ``"driver"``. valid_from : str | None Start of the temporal validity window (ISO 8601 or year string). - Aligns with STAC ``start_datetime``. ``None`` means no lower bound. + Aligns with STAC ``start_datetime``. ``None`` means no lower bound. valid_until : str | None End of the temporal validity window (ISO 8601 or year string). - Aligns with STAC ``end_datetime``. ``None`` means no upper bound. + Aligns with STAC ``end_datetime``. ``None`` means no upper bound. Both ``None`` → static variable (aligns with STAC ``datetime``). purity_threshold : float | None - Reserved for future purity-masking logic. Included in + Reserved for future purity-masking logic. Included in ``spec_hash()`` so that two derivations with different thresholds - are always distinct products. Currently unused in execution. + are always distinct products. Currently unused in execution. bbox : list[float] | None Optional bounding box ``[xmin, ymin, xmax, ymax]`` in EPSG:4326. - Aligns with the STAC ``bbox`` field. Reserved — not used in + Aligns with the STAC ``bbox`` field. Reserved — not used in execution and excluded from ``spec_hash()``. params : dict Operator options, e.g. ``{"crs": "EPSG:5880"}`` for ``distance`` or @@ -84,7 +86,7 @@ class Derivation(BaseModel): fill: Literal["nearest"] | None = None @model_validator(mode="after") - def _validate_operator(self) -> "Derivation": + def _validate_operator(self) -> Derivation: available = sorted(OPERATOR_REGISTRY) if self.operator not in OPERATOR_REGISTRY: raise ValueError( @@ -157,7 +159,7 @@ def spec_hash(self, grid_id: str = "__global__") -> str: Deterministic SHA-256 hash of the derivation spec. Delegates to ``SpatialDerivation.spec_hash()`` for the base fields, - then folds in ``purity_threshold`` when it is set. ``bbox`` is + then folds in ``purity_threshold`` when it is set. ``bbox`` is excluded because it is descriptive metadata and does not affect what is computed. @@ -165,7 +167,7 @@ def spec_hash(self, grid_id: str = "__global__") -> str: ---------- grid_id : str Grid identifier used for the underlying - ``SpatialDerivation.spec_hash()``. Defaults to ``"__global__"`` + ``SpatialDerivation.spec_hash()``. Defaults to ``"__global__"`` for grid-agnostic comparisons. Returns diff --git a/disscube/models/grid.py b/disscube/models/grid.py index b5f13ec..36fca79 100644 --- a/disscube/models/grid.py +++ b/disscube/models/grid.py @@ -1,10 +1,15 @@ +import logging +import math import re import warnings -from typing import Literal +from typing import Any, Literal import numpy as np from affine import Affine from pydantic import BaseModel +from pyproj import Transformer + +log = logging.getLogger(__name__) class GridAnchor(BaseModel): @@ -51,7 +56,7 @@ def cols(self) -> int: @property def transform(self) -> Affine: # North-up transform: origin at (minx, maxy), negative y-scale - return Affine.translation(self.bbox[0], self.bbox[3]) * Affine.scale(self.resolution, -self.resolution) + return Affine.translation(self.bbox[0], self.bbox[3]) @ Affine.scale(self.resolution, -self.resolution) @property def xs(self) -> np.ndarray: @@ -122,3 +127,114 @@ def parse_cell_id(cell_id: str) -> tuple[str, int, int]: if isinstance(e, ValueError) and "Invalid coordinate format" in str(e): raise raise ValueError(f"Invalid cell_id format: {cell_id}") from e + + +# --------------------------------------------------------------------------- +# Master Grid Constants (Brazil Data Cube Standard) +# --------------------------------------------------------------------------- + +BDC_CRS = ( + "+proj=aea +lat_0=-12 +lon_0=-54 +lat_1=-2 +lat_2=-22" + " +x_0=5000000 +y_0=10000000 +ellps=GRS80 +units=m +no_defs" +) + +# Full Brazil bbox in BDC Albers, snapped to 5 km mesh. +BRAZIL_BBOX: list[float] = [2_720_000, 7_500_000, 7_870_000, 11_830_000] + +# National reference resolutions +SIMULATION_GRIDS = [ + ("BR/5km", 5_000.0, "National LUCC grid — 5 km pixels, BDC Albers"), + ("BR/1km", 1_000.0, "National LUCC grid — 1 km pixels, BDC Albers"), +] + +# --------------------------------------------------------------------------- +# Grid Registration Utilities +# --------------------------------------------------------------------------- + +def register_simulation_grids(cube: Any) -> None: + """Register national simulation grids (pixel resolution of derived output).""" + for grid_id, resolution, description in SIMULATION_GRIDS: + grid = GridSpec( + id=grid_id, + type="reference", + crs=BDC_CRS, + resolution=resolution, + bbox=BRAZIL_BBOX, + description=description, + ) + cube.register_grid(grid) + log.info("[grid] registered %r (%d m pixels, %d rows × %d cols)", + grid_id, int(resolution), grid.rows, grid.cols) + + +def register_local_grid( + cube: Any, + name: str | None = None, + state: str | None = None, + bbox_geo: tuple[float, float, float, float] | None = None, + resolution: float = 5_000.0, + snap: bool = True, +) -> GridSpec: + """ + Register a local simulation grid snapped to the national mesh. + + This ensures that any Area of Interest (AOI) has pixels that align + perfectly with the national master grids, enabling interoperability + without resampling. + """ + name = name or state + if not name: + raise ValueError("Either 'name' or 'state' must be provided.") + + if bbox_geo is None: + raise ValueError("'bbox_geo' must be provided.") + + transformer = Transformer.from_crs("EPSG:4326", BDC_CRS, always_xy=True) + + corners = [ + (bbox_geo[0], bbox_geo[1]), # SW + (bbox_geo[0], bbox_geo[3]), # NW + (bbox_geo[2], bbox_geo[3]), # NE + (bbox_geo[2], bbox_geo[1]), # SE + ] + xs, ys = zip(*(transformer.transform(lon, lat) for lon, lat in corners)) + + minx, miny, maxx, maxy = min(xs), min(ys), max(xs), max(ys) + + if snap: + minx = math.floor(minx / resolution) * resolution + miny = math.floor(miny / resolution) * resolution + maxx = math.ceil(maxx / resolution) * resolution + maxy = math.ceil(maxy / resolution) * resolution + + is_km = resolution >= 1000 and math.isclose(resolution % 1000, 0, abs_tol=1e-5) + if is_km: + res_str = f"{int(resolution // 1000)}km" + else: + res_str = f"{int(resolution)}m" + + grid_id = f"{name}/{res_str}" + grid = GridSpec( + id=grid_id, + type="reference", + crs=BDC_CRS, + resolution=resolution, + bbox=[minx, miny, maxx, maxy], + description=f"{name} simulation grid — {resolution:.0f} m pixels, BDC Albers", + ) + cube.register_grid(grid) + + # Back-project for human-readable summary + back = Transformer.from_crs(BDC_CRS, "EPSG:4326", always_xy=True) + lon_min, lat_min = back.transform(minx, miny) + lon_max, lat_max = back.transform(maxx, maxy) + + log.info( + "[grid] registered %r bbox(BDC Albers)=[%.0f, %.0f, %.0f, %.0f]" + " bbox(geo)=lon[%.2f,%.2f] lat[%.2f,%.2f]" + " size=%d×%d cells (%.0f m pixels)", + grid_id, minx, miny, maxx, maxy, + lon_min, lon_max, lat_min, lat_max, + grid.rows, grid.cols, resolution, + ) + return grid diff --git a/disscube/pipeline/__init__.py b/disscube/pipeline/__init__.py index c7ef57d..221882b 100644 --- a/disscube/pipeline/__init__.py +++ b/disscube/pipeline/__init__.py @@ -1,3 +1,41 @@ -from .context import PipelineContext, PipelineStage +""" +DisSCube pipeline package: declarative pipeline files, planning, execution and raster alignment. +""" -__all__ = ["PipelineContext", "PipelineStage"] +from __future__ import annotations + +import sys + +# Backward-compatibility alias for legacy code importing disscube.config +from disscube.pipeline import runner as _runner +from disscube.pipeline import schema as _schema +from disscube.pipeline.context import PipelineContext, PipelineStage +from disscube.pipeline.runner import ( + ExportReport, + PipelineError, + Plan, + RunReport, + export_cube, + load, + plan, + run, +) +from disscube.pipeline.schema import PipelineConfig + +sys.modules["disscube.config"] = sys.modules[__name__] +sys.modules["disscube.config.runner"] = _runner +sys.modules["disscube.config.schema"] = _schema + +__all__ = [ + "ExportReport", + "PipelineConfig", + "PipelineContext", + "PipelineError", + "PipelineStage", + "Plan", + "RunReport", + "export_cube", + "load", + "plan", + "run", +] diff --git a/disscube/pipeline/aligner.py b/disscube/pipeline/aligner.py index df5521f..b1fb06f 100644 --- a/disscube/pipeline/aligner.py +++ b/disscube/pipeline/aligner.py @@ -336,7 +336,7 @@ def _align_fine( # Fine transform shares the target grid origin (north-up). fine_transform = ( Affine.translation(grid.bbox[0], grid.bbox[3]) - * Affine.scale(fine_res, -fine_res) + @ Affine.scale(fine_res, -fine_res) ) nodata = _source_nodata(band) diff --git a/disscube/config/runner.py b/disscube/pipeline/runner.py similarity index 98% rename from disscube/config/runner.py rename to disscube/pipeline/runner.py index 19adc28..aece1df 100644 --- a/disscube/config/runner.py +++ b/disscube/pipeline/runner.py @@ -15,7 +15,7 @@ import numpy as np from pydantic import ValidationError -from disscube.config.schema import ( +from disscube.pipeline.schema import ( BdcSource, ClassifiedSource, DeriveConfig, @@ -27,7 +27,7 @@ ProdesSource, UnionSource, ) -from disscube.utils.files import sha256_file +from disscube.utils import sha256_file log = logging.getLogger(__name__) @@ -71,7 +71,7 @@ class PlannedDerive: fill: str | None = None def derivation(self): - from disscube.derivation import Derivation + from disscube.models import Derivation return Derivation(target=self.target, source_id=self.source, operator=self.operator, class_code=self.class_code, role=self.role, params=self.params, fill=self.fill) @@ -298,6 +298,11 @@ def _save_geotiff_from_backend(backend, variables: list[str], grid: GridConfig, nodata=np.nan, compress="deflate", ) as dst: + dst.update_tags( + TIFFTAG_SOFTWARE="DisSCube 0.3.0", + GRID_ID=grid.name, + CONVENTIONS="CF-1.8", + ) for idx, (var, arr) in enumerate(zip(variables, arrays), start=1): dst.write(arr, idx) dst.set_band_description(idx, var) @@ -433,7 +438,7 @@ def _register_grid(cube, g: GridConfig) -> tuple[str, list[float]]: from pyproj import Transformer from disscube.models import GridSpec - from disscube.utils.grids import register_local_grid + from disscube.utils import register_local_grid min_x, min_y, max_x, max_y = g.bbox if g.crs is None: diff --git a/disscube/config/schema.py b/disscube/pipeline/schema.py similarity index 100% rename from disscube/config/schema.py rename to disscube/pipeline/schema.py diff --git a/disscube/pipeline/writer.py b/disscube/pipeline/writer.py index 11a88ba..9f9e7bf 100644 --- a/disscube/pipeline/writer.py +++ b/disscube/pipeline/writer.py @@ -38,10 +38,31 @@ def execute(self, ctx: PipelineContext) -> PipelineContext: da.attrs["grid_id"] = grid.id da.attrs["role"] = derivation.role da.attrs["spec_hash"] = spec_hash + da.attrs["conventions"] = "CF-1.8" + op_name = getattr(derivation, "operator", None) + if not op_name and hasattr(derivation, "variables"): + for v in derivation.variables: + if getattr(v, "name", None) == var_name: + op_name = getattr(v, "operator", None) + break + if op_name: + da.attrs["operator"] = str(op_name) + if getattr(derivation, "source_id", None): + da.attrs["source_id"] = derivation.source_id if tile_id: da.attrs["tile_id"] = tile_id if "spatial_ref" in da.coords: da.attrs["crs"] = grid.crs + + # Coordinate metadata (CF conventions) + if "x" in da.coords and not da.coords["x"].attrs.get("standard_name"): + is_geo = "4326" in str(grid.crs).lower() or "longlat" in str(grid.crs).lower() + da.coords["x"].attrs["standard_name"] = "longitude" if is_geo else "projection_x_coordinate" + da.coords["x"].attrs["units"] = "degrees_east" if is_geo else "m" + if "y" in da.coords and not da.coords["y"].attrs.get("standard_name"): + is_geo = "4326" in str(grid.crs).lower() or "longlat" in str(grid.crs).lower() + da.coords["y"].attrs["standard_name"] = "latitude" if is_geo else "projection_y_coordinate" + da.coords["y"].attrs["units"] = "degrees_north" if is_geo else "m" # Storage path: derived/{grid_id}/{tile_id or 'global'}/{spec_hash}/{var_name}.zarr partition = tile_id if tile_id else "global" diff --git a/disscube/sources/_raster.py b/disscube/sources/_raster.py index 7128c9a..63f1bce 100644 --- a/disscube/sources/_raster.py +++ b/disscube/sources/_raster.py @@ -25,8 +25,7 @@ from rasterio.warp import transform_bounds from rasterio.windows import Window, from_bounds -from disscube.utils.files import sha256_file -from disscube.utils.grids import BDC_CRS +from disscube.utils import BDC_CRS, sha256_file log = logging.getLogger(__name__) diff --git a/disscube/sources/classified.py b/disscube/sources/classified.py index 7383f38..be48ae0 100644 --- a/disscube/sources/classified.py +++ b/disscube/sources/classified.py @@ -26,7 +26,7 @@ from disscube.sources._categorical import read_legend from disscube.sources._raster import read_window, register_raster -from disscube.utils.files import sha256_file +from disscube.utils import sha256_file def register_classified_map( diff --git a/disscube/sources/prodes.py b/disscube/sources/prodes.py index 07c85cd..3f1e859 100644 --- a/disscube/sources/prodes.py +++ b/disscube/sources/prodes.py @@ -39,7 +39,7 @@ from disscube.sources._categorical import read_qml_legend, reclassify, strip_code from disscube.sources._raster import Window2D, read_window, register_raster -from disscube.utils.files import sha256_file +from disscube.utils import sha256_file log = logging.getLogger(__name__) diff --git a/disscube/storage/local.py b/disscube/storage.py similarity index 68% rename from disscube/storage/local.py rename to disscube/storage.py index 6f769c1..023a6bb 100644 --- a/disscube/storage/local.py +++ b/disscube/storage.py @@ -1,3 +1,11 @@ +""" +Storage layer for DisSCube assets (local and remote via fsspec). +""" + +from __future__ import annotations + +import sys + import fsspec @@ -16,3 +24,9 @@ def exists(self, relative_path: str) -> bool: def open(self, relative_path: str, mode: str = "rb"): return self.fs.open(f"{self.path}/{relative_path}", mode=mode) + + +# Backward-compatibility alias for legacy code importing disscube.storage.local +sys.modules[f"{__name__}.local"] = sys.modules[__name__] + +__all__ = ["AssetStore"] diff --git a/disscube/storage/__init__.py b/disscube/storage/__init__.py deleted file mode 100644 index bd4b20d..0000000 --- a/disscube/storage/__init__.py +++ /dev/null @@ -1,3 +0,0 @@ -from .local import AssetStore - -__all__ = ["AssetStore"] diff --git a/disscube/utils/bdc_importer.py b/disscube/utils.py similarity index 62% rename from disscube/utils/bdc_importer.py rename to disscube/utils.py index 10b16c5..b357a35 100644 --- a/disscube/utils/bdc_importer.py +++ b/disscube/utils.py @@ -1,24 +1,59 @@ +""" +Utilities and helper functions for DisSCube. + +Consolidates file checksums, BDC coordinate systems, grid registration, +and BDC shapefile importer helpers into a single module. +""" + +from __future__ import annotations + +import hashlib import logging +import sys from importlib.resources import files +from pathlib import Path +from typing import Any from shapely.geometry import shape -from disscube.client import CubeClient -from disscube.models import SpatialSource - -from .grids import BDC_CRS, register_simulation_grids +from disscube.models.grid import ( + BDC_CRS, + BRAZIL_BBOX, + SIMULATION_GRIDS, + register_local_grid, + register_simulation_grids, +) +from disscube.models.variable import SpatialSource log = logging.getLogger(__name__) + +# --------------------------------------------------------------------------- +# File utilities +# --------------------------------------------------------------------------- + +def sha256_file(path: str | Path, chunk_size: int = 1 << 20) -> str: + """SHA-256 of a file's bytes, as ``"sha256:"``. + + Use it as ``SpatialSource.checksum``: when the file changes, derivations + from that source get a new ``spec_hash`` instead of a stale cache hit. + """ + digest = hashlib.sha256() + with open(path, "rb") as fh: + for block in iter(lambda: fh.read(chunk_size), b""): + digest.update(block) + return f"sha256:{digest.hexdigest()}" + + # --------------------------------------------------------------------------- -# BDC Specific Constants +# BDC Specific Constants & Importer # --------------------------------------------------------------------------- # Tile sizes of BDC Grid V2 (Albers equal-area, metres); ~1°, ~2° and ~4°. BDC_TILE_LEVELS = [ - ("SM", "sm_path", "BDC Small tile grid (105.6 km × 105.6 km, ~1°)"), - ("MD", "md_path", "BDC Medium tile grid (211.2 km × 211.2 km, ~2°)"), - ("LG", "lg_path", "BDC Large tile grid (422.4 km × 422.4 km, ~4°)"), + ("SM", "sm_path", "BDC Small tile grid (105.6 km × 105.6 km, ~1°)"), + ("MD", "md_path", "BDC Medium tile grid (211.2 km × 211.2 km, ~2°)"), + ("LG", "lg_path", "BDC Large tile grid (422.4 km × 422.4 km, ~4°)"), ] @@ -29,12 +64,6 @@ def bundled_bdc_grid(level: str) -> str: ``level`` is one of ``"SM"``, ``"MD"`` or ``"LG"``. Provenance, checksums and licensing notes are in ``disscube/data/bdc_grids/README.md``. - - Note: the bundled ``.prj`` files carry a non-existent authority code - (``EPSG:200000``, emitted by the BDC GeoServer), so readers that resolve - the CRS from the file — e.g. ``geopandas.read_file`` — fail on them. The - importer does not read the file CRS; it uses :data:`BDC_CRS`, which is the - same projection. """ level = level.upper() if level not in {lvl for lvl, _, _ in BDC_TILE_LEVELS}: @@ -42,16 +71,13 @@ def bundled_bdc_grid(level: str) -> str: path = files("disscube") / "data" / "bdc_grids" / f"BDC_{level}_V2.zip" return f"zip://{path}" -# --------------------------------------------------------------------------- -# Public API -# --------------------------------------------------------------------------- def import_bdc_grids( - cube: CubeClient, + cube: Any, sm_path: str | None = None, md_path: str | None = None, lg_path: str | None = None, -): +) -> None: """ Import BDC tiles and national simulation grids into the catalog. @@ -67,11 +93,7 @@ def import_bdc_grids( _register_tile_sources(cube, paths) -# --------------------------------------------------------------------------- -# Internal helpers -# --------------------------------------------------------------------------- - -def _register_tile_sources(cube: CubeClient, paths: dict[str, str]) -> None: +def _register_tile_sources(cube: Any, paths: dict[str, str]) -> None: """Register BDC tile envelopes as SpatialSources.""" try: import fiona @@ -106,3 +128,23 @@ def _register_tile_sources(cube: CubeClient, paths: dict[str, str]) -> None: count += 1 log.info("[tiles] registered %d BDC_%s tiles", count, label) + + +# --------------------------------------------------------------------------- +# Backward-compatibility aliases for legacy submodules +# --------------------------------------------------------------------------- +sys.modules[f"{__name__}.files"] = sys.modules[__name__] +sys.modules[f"{__name__}.grids"] = sys.modules[__name__] +sys.modules[f"{__name__}.bdc_importer"] = sys.modules[__name__] + +__all__ = [ + "BDC_CRS", + "BDC_TILE_LEVELS", + "BRAZIL_BBOX", + "SIMULATION_GRIDS", + "bundled_bdc_grid", + "import_bdc_grids", + "register_local_grid", + "register_simulation_grids", + "sha256_file", +] diff --git a/disscube/utils/files.py b/disscube/utils/files.py deleted file mode 100644 index 4efaf84..0000000 --- a/disscube/utils/files.py +++ /dev/null @@ -1,19 +0,0 @@ -"""Small file helpers.""" - -from __future__ import annotations - -import hashlib -from pathlib import Path - - -def sha256_file(path: str | Path, chunk_size: int = 1 << 20) -> str: - """SHA-256 of a file's bytes, as ``"sha256:"``. - - Use it as ``SpatialSource.checksum``: when the file changes, derivations - from that source get a new ``spec_hash`` instead of a stale cache hit. - """ - digest = hashlib.sha256() - with open(path, "rb") as fh: - for block in iter(lambda: fh.read(chunk_size), b""): - digest.update(block) - return f"sha256:{digest.hexdigest()}" diff --git a/disscube/utils/grids.py b/disscube/utils/grids.py deleted file mode 100644 index 4827a77..0000000 --- a/disscube/utils/grids.py +++ /dev/null @@ -1,124 +0,0 @@ -import logging -import math - -from pyproj import Transformer - -from disscube.client import CubeClient -from disscube.models import GridSpec - -log = logging.getLogger(__name__) - -# --------------------------------------------------------------------------- -# Master Grid Constants (Brazil Data Cube Standard) -# --------------------------------------------------------------------------- - -# EPSG:10857 is the official code for BDC Albers (SIRGAS 2000 / Brazil Albers -# Equal Area Conic), registered in 2023. Using it instead of a raw proj4 string -# ensures QGIS and other tools recognise the CRS automatically without requiring -# manual configuration or a custom CRS entry. -#BDC_CRS = "EPSG:10857" -BDC_CRS = ( - "+proj=aea +lat_0=-12 +lon_0=-54 +lat_1=-2 +lat_2=-22" - " +x_0=5000000 +y_0=10000000 +ellps=GRS80 +units=m +no_defs" -) - -# Full Brazil bbox in BDC Albers, snapped to 5 km mesh. -BRAZIL_BBOX: list[float] = [2_720_000, 7_500_000, 7_870_000, 11_830_000] - -# National reference resolutions -SIMULATION_GRIDS = [ - ("BR/5km", 5_000.0, "National LUCC grid — 5 km pixels, BDC Albers"), - ("BR/1km", 1_000.0, "National LUCC grid — 1 km pixels, BDC Albers"), -] - -# --------------------------------------------------------------------------- -# Grid Registration Utilities -# --------------------------------------------------------------------------- - -def register_simulation_grids(cube: CubeClient) -> None: - """Register national simulation grids (pixel resolution of derived output).""" - for grid_id, resolution, description in SIMULATION_GRIDS: - grid = GridSpec( - id=grid_id, - type="reference", - crs=BDC_CRS, - resolution=resolution, - bbox=BRAZIL_BBOX, - description=description, - ) - cube.register_grid(grid) - log.info("[grid] registered %r (%d m pixels, %d rows × %d cols)", - grid_id, int(resolution), grid.rows, grid.cols) - - -def register_local_grid( - cube: CubeClient, - name: str | None = None, - state: str | None = None, - bbox_geo: tuple[float, float, float, float] | None = None, - resolution: float = 5_000.0, - snap: bool = True, -) -> GridSpec: - """ - Register a local simulation grid snapped to the national mesh. - - This ensures that any Area of Interest (AOI) has pixels that align - perfectly with the national master grids, enabling interoperability - without resampling. - """ - name = name or state - if not name: - raise ValueError("Either 'name' or 'state' must be provided.") - - if bbox_geo is None: - raise ValueError("'bbox_geo' must be provided.") - - transformer = Transformer.from_crs("EPSG:4326", BDC_CRS, always_xy=True) - - corners = [ - (bbox_geo[0], bbox_geo[1]), # SW - (bbox_geo[0], bbox_geo[3]), # NW - (bbox_geo[2], bbox_geo[3]), # NE - (bbox_geo[2], bbox_geo[1]), # SE - ] - xs, ys = zip(*(transformer.transform(lon, lat) for lon, lat in corners)) - - minx, miny, maxx, maxy = min(xs), min(ys), max(xs), max(ys) - - if snap: - minx = math.floor(minx / resolution) * resolution - miny = math.floor(miny / resolution) * resolution - maxx = math.ceil(maxx / resolution) * resolution - maxy = math.ceil(maxy / resolution) * resolution - - is_km = resolution >= 1000 and math.isclose(resolution % 1000, 0, abs_tol=1e-5) - if is_km: - res_str = f"{int(resolution // 1000)}km" - else: - res_str = f"{int(resolution)}m" - - grid_id = f"{name}/{res_str}" - grid = GridSpec( - id=grid_id, - type="reference", - crs=BDC_CRS, - resolution=resolution, - bbox=[minx, miny, maxx, maxy], - description=f"{name} simulation grid — {resolution:.0f} m pixels, BDC Albers", - ) - cube.register_grid(grid) - - # Back-project for human-readable summary - back = Transformer.from_crs(BDC_CRS, "EPSG:4326", always_xy=True) - lon_min, lat_min = back.transform(minx, miny) - lon_max, lat_max = back.transform(maxx, maxy) - - log.info( - "[grid] registered %r bbox(BDC Albers)=[%.0f, %.0f, %.0f, %.0f]" - " bbox(geo)=lon[%.2f,%.2f] lat[%.2f,%.2f]" - " size=%d×%d cells (%.0f m pixels)", - grid_id, minx, miny, maxx, maxy, - lon_min, lon_max, lat_min, lat_max, - grid.rows, grid.cols, resolution, - ) - return grid diff --git a/docs/architecture/overview.md b/docs/architecture/overview.md index dcfea92..ea59bcf 100644 --- a/docs/architecture/overview.md +++ b/docs/architecture/overview.md @@ -28,7 +28,7 @@ SpatialSource ──► SpatialDerivation ──► Variable ──► DerivedVa **Declarative (recommended):** ```python -from disscube.derivation import Derivation +from disscube.models import Derivation d = Derivation( target="forest_pct", diff --git a/docs/index.md b/docs/index.md index e660a6f..4c88224 100644 --- a/docs/index.md +++ b/docs/index.md @@ -23,7 +23,7 @@ pip install -e . ```python from disscube.client import CubeClient -from disscube.derivation import Derivation +from disscube.models import Derivation cube = CubeClient(catalog="catalog.db", store="./data/") diff --git a/examples/07_bdc_cube.py b/examples/07_bdc_cube.py index c6fc2b6..2000c5c 100644 --- a/examples/07_bdc_cube.py +++ b/examples/07_bdc_cube.py @@ -39,8 +39,7 @@ from disscube import CubeClient, Derivation, SpatialSource from disscube.sources import Window2D, register_raster -from disscube.utils.files import sha256_file -from disscube.utils.grids import BDC_CRS, register_local_grid +from disscube.utils import BDC_CRS, register_local_grid, sha256_file COLLECTION = "LANDSAT-16D-1" PERIOD = "2020-07-01/2020-09-30" # dry season diff --git a/examples/08_mapbiomas_land_use.py b/examples/08_mapbiomas_land_use.py index 95a66b0..712f8ed 100644 --- a/examples/08_mapbiomas_land_use.py +++ b/examples/08_mapbiomas_land_use.py @@ -39,7 +39,7 @@ from disscube import CubeClient, Derivation from disscube.sources import Window2D, register_raster from disscube.sources.mapbiomas import CLASSES, MAPBIOMAS_NODATA, register_mapbiomas_source -from disscube.utils.grids import register_local_grid +from disscube.utils import register_local_grid YEARS = (2000, 2020) BBOX = (-44.35, -2.62, -44.20, -2.47) # São Luís, Ilha do Maranhão (WGS84) diff --git a/examples/09_prodes_deforestation.py b/examples/09_prodes_deforestation.py index fa781b5..bb1fba9 100644 --- a/examples/09_prodes_deforestation.py +++ b/examples/09_prodes_deforestation.py @@ -35,7 +35,7 @@ from disscube import CubeClient, Derivation from disscube.sources import prodes -from disscube.utils.grids import register_local_grid +from disscube.utils import register_local_grid YEARS = (2008, 2016, 2024) BBOX = (-54.842, -3.587, -54.459, -3.168) # LuccME Lab15 cellular space (cs_moju) diff --git a/examples/README.md b/examples/README.md index 5fe406c..7fe4037 100644 --- a/examples/README.md +++ b/examples/README.md @@ -47,9 +47,16 @@ DisSCube prepares data for models; it stops at `CubeClient.to_lucc_data()`. Examples that run simulations with the prepared data (BR-MANGUE, LUCC) belong to the model repositories, where those dependencies live. -## Utilities (`tools/`) +## Exporting to GeoTIFF -| Script | Purpose | -|---|---| -| `tools/zarr_to_tif.py` | Converts a derived Zarr to GeoTIFF | -| `tools/import_bdc_tiles.py` | Registers the BDC tile grids in a catalog | +Derived data cubes can be exported to multi-band GeoTIFFs using the CLI: + +```bash +disscube export examples/pipelines/itaituba_fill.toml --output data/cellspace.tif +``` + +Or directly during execution: + +```bash +disscube run examples/pipelines/itaituba_fill.toml --output data/cellspace.tif +``` diff --git a/pyproject.toml b/pyproject.toml index c79cd19..be630c6 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -103,9 +103,13 @@ line-length = 120 target-version = "py311" exclude = ["docs", "site", "build"] +[tool.ruff.lint.isort] +known-first-party = ["disscube"] + + [tool.mypy] python_version = "3.12" -exclude = ["tests/", "docs/", "examples/", "tools/"] +exclude = ["tests/", "docs/", "examples/"] [[tool.mypy.overrides]] # Third-party (or optional-extra) dependencies with no bundled type stubs. diff --git a/requirements.txt b/requirements.txt index 01544cb..e018412 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,18 +10,32 @@ numpy pandas pyproj fsspec +toml +rioxarray +affine +pooch>=1.8.0 +dissmodel>=0.6.0,<0.7.0 + +# Optional: BDC +fiona +pystac-client + +# Optional: S3 s3fs + +# Optional: API fastapi uvicorn python-multipart -toml -rioxarray -affine -# Development & Documentation -mkdocs-material -pymdown-extensions +# Development & Testing pytest -pytest-mock -black -flake8 +pytest-cov +mypy>=1.13 +ruff>=0.16,<0.17 +httpx +ipykernel + +# Documentation +mkdocs>=1.6,<2 +mkdocs-material>=9.5 diff --git a/tests/test_anchoring.py b/tests/test_anchoring.py index 37d4e06..a89fdc8 100644 --- a/tests/test_anchoring.py +++ b/tests/test_anchoring.py @@ -1,81 +1,76 @@ -import os - from disscube.client import CubeClient from disscube.models import GridSpec, SpatialRelation -# Initialize client -if os.path.exists("test_catalog.json"): - os.remove("test_catalog.json") -cube = CubeClient(catalog="test_catalog.json", store="./data/") -# 1. Register BDC Reference Grid (Parent 0) -bdc_tile = GridSpec( - id="BDC_SM_089094", - type="reference", - crs="EPSG:200000", - resolution=10.0, - bbox=[5000000, 10000000, 5105600, 10105600], - description="BDC SM Tile 089094" -) -cube.register_grid(bdc_tile) -print(f"Registered: {bdc_tile.id}") +def test_grid_anchoring_and_lineage(tmp_path): + catalog_path = tmp_path / "catalog.json" + store_path = tmp_path / "store" + cube = CubeClient(catalog=str(catalog_path), store=str(store_path)) + + # 1. Register BDC Reference Grid (Parent 0) + bdc_tile = GridSpec( + id="BDC_SM_089094", + type="reference", + crs="EPSG:200000", + resolution=10.0, + bbox=[5000000, 10000000, 5105600, 10105600], + description="BDC SM Tile 089094", + ) + cube.register_grid(bdc_tile) -# 2. Register 5km Grid -grid_5km = GridSpec( - id="REGIONAL_5KM", - type="local", - crs="EPSG:200000", - resolution=5000.0, - bbox=[5010000, 10010000, 5060000, 10060000], - description="Regional 5km model" -) -cube.register_grid(grid_5km) -cube.register_relation(SpatialRelation( - source_grid_id="REGIONAL_5KM", - target_grid_id="BDC_SM_089094", - strategy="simple" -)) -print(f"Registered: {grid_5km.id}") + # 2. Register 5km Grid + grid_5km = GridSpec( + id="REGIONAL_5KM", + type="local", + crs="EPSG:200000", + resolution=5000.0, + bbox=[5010000, 10010000, 5060000, 10060000], + description="Regional 5km model", + ) + cube.register_grid(grid_5km) + cube.register_relation( + SpatialRelation( + source_grid_id="REGIONAL_5KM", + target_grid_id="BDC_SM_089094", + strategy="simple", + ) + ) -# 3. Register 1km Grid -grid_1km = GridSpec( - id="LOCAL_1KM", - type="local", - crs="EPSG:200000", - resolution=1000.0, - bbox=[5020000, 10020000, 5030000, 10030000], - description="Local 1km model" -) -cube.register_grid(grid_1km) -cube.register_relation(SpatialRelation( - source_grid_id="LOCAL_1KM", - target_grid_id="REGIONAL_5KM", - strategy="simple" -)) -print(f"Registered: {grid_1km.id}") + # 3. Register 1km Grid + grid_1km = GridSpec( + id="LOCAL_1KM", + type="local", + crs="EPSG:200000", + resolution=1000.0, + bbox=[5020000, 10020000, 5030000, 10030000], + description="Local 1km model", + ) + cube.register_grid(grid_1km) + cube.register_relation( + SpatialRelation( + source_grid_id="LOCAL_1KM", + target_grid_id="REGIONAL_5KM", + strategy="simple", + ) + ) -# 4. Verify Catalog Persistence -print("\nVerifying catalog content...") -all_grids = cube.catalog.list_grids() -for g in all_grids: - relations = cube.get_relations(g.id) - rel_info = "" - for r in relations: - if r.source_grid_id == g.id: - rel_info += f" -> Relates to: {r.target_grid_id} ({r.strategy})" - print(f"- {g.id} [{g.type}]{rel_info}") + # 4. Verify Catalog Persistence + all_grids = cube.catalog.list_grids() + assert len(all_grids) == 3 + grid_ids = {g.id for g in all_grids} + assert grid_ids == {"BDC_SM_089094", "REGIONAL_5KM", "LOCAL_1KM"} -# 5. Simple lineage check -def get_lineage(grid_id, cube): - lineage = [grid_id] - relations = cube.get_relations(grid_id) - # Find relation where current grid is source - source_rel = next((r for r in relations if r.source_grid_id == grid_id), None) - while source_rel: - lineage.append(source_rel.target_grid_id) - grid_id = source_rel.target_grid_id - relations = cube.get_relations(grid_id) + # 5. Lineage check + def get_lineage(grid_id, c): + lineage = [grid_id] + relations = c.get_relations(grid_id) source_rel = next((r for r in relations if r.source_grid_id == grid_id), None) - return lineage + while source_rel: + lineage.append(source_rel.target_grid_id) + grid_id = source_rel.target_grid_id + relations = c.get_relations(grid_id) + source_rel = next((r for r in relations if r.source_grid_id == grid_id), None) + return lineage -print(f"\nLineage for LOCAL_1KM: {' -> '.join(get_lineage('LOCAL_1KM', cube))}") + lineage = get_lineage("LOCAL_1KM", cube) + assert lineage == ["LOCAL_1KM", "REGIONAL_5KM", "BDC_SM_089094"] diff --git a/tests/test_derivation.py b/tests/test_derivation.py index d60a6e0..7e7565b 100644 --- a/tests/test_derivation.py +++ b/tests/test_derivation.py @@ -1,10 +1,10 @@ """ -Tests for the declarative Derivation model (disscube/derivation.py). +Tests for the declarative Derivation model (disscube/models/derivation.py). """ import pytest -from disscube.derivation import Derivation +from disscube.models import Derivation from disscube.models.variable import SpatialDerivation, Variable # ── Construction-time validation ────────────────────────────────────────────── diff --git a/tests/test_fetch.py b/tests/test_fetch.py index 730d53a..6635f96 100644 --- a/tests/test_fetch.py +++ b/tests/test_fetch.py @@ -1,7 +1,7 @@ from unittest.mock import patch -from disscube.config.runner import _fetch_file_source -from disscube.config.schema import FileSource +from disscube.pipeline.runner import _fetch_file_source +from disscube.pipeline.schema import FileSource def test_fetch_file_source_calls_pooch(tmp_path): diff --git a/tests/test_geometry.py b/tests/test_geometry.py index 14fbb05..b6d7b4d 100644 --- a/tests/test_geometry.py +++ b/tests/test_geometry.py @@ -18,7 +18,7 @@ def test_gridspec_properties(): assert grid.cols == 10 # North-up transform: origin at (0, 1000), dx=100, dy=-100 - expected_transform = Affine.translation(0, 1000) * Affine.scale(100, -100) + expected_transform = Affine.translation(0, 1000) @ Affine.scale(100, -100) assert grid.transform == expected_transform # xs should be [50, 150, ..., 950] diff --git a/tests/test_pipeline_file.py b/tests/test_pipeline_file.py index 79d7959..bcb2332 100644 --- a/tests/test_pipeline_file.py +++ b/tests/test_pipeline_file.py @@ -18,8 +18,7 @@ from disscube import CubeClient, GridSpec, SpatialDerivation, SpatialSource, Variable from disscube.cli import main as cli -from disscube.config import load, plan, run -from disscube.config.runner import PipelineError +from disscube.pipeline import PipelineError, load, plan, run from disscube.sources import Window2D ROOT = Path(__file__).resolve().parents[1] @@ -270,6 +269,71 @@ def test_cli_reports_errors(tmp_path, capsys): assert "unknown source" in capsys.readouterr().err +def test_cli_validate_json(capsys): + assert cli(["validate", str(ITAITUBA), "--json"]) == 0 + out = capsys.readouterr().out + data = json.loads(out) + assert data["status"] == "ok" + assert data["grid"] == "itaituba/5km" + assert "elevation" in data["derives"] + assert "elevation" in data["sources"] + + +def test_cli_validate_json_error(tmp_path, capsys): + path = _toml(tmp_path, '[[derive]]\ntarget="x"\nsource="nope"\noperator="mean"\n') + assert cli(["validate", str(path), "--json"]) == 2 + err = capsys.readouterr().err + data = json.loads(err) + assert data["status"] == "error" + assert "unknown source" in data["error"] + + +def test_cli_run_dry_run(capsys): + assert cli(["run", str(ITAITUBA), "--dry-run"]) == 0 + out = capsys.readouterr().out + assert "[dry-run]" in out + assert "Plan is valid" in out + + +def test_cli_run_dry_run_json(capsys): + assert cli(["run", str(ITAITUBA), "--dry-run", "--json"]) == 0 + out = capsys.readouterr().out + data = json.loads(out) + assert data["status"] == "ok" + assert data["dry_run"] is True + assert data["grid"] == "itaituba/5km" + assert any(s["id"] == "elevation" for s in data["sources"]) + assert any(d["target"] == "elevation" for d in data["derives"]) + + +def test_cli_run_json(tmp_path, capsys): + ws = tmp_path / "ws_json" + assert cli(["run", str(ITAITUBA), "--workspace", str(ws), "--json"]) == 0 + out = capsys.readouterr().out + data = json.loads(out) + assert data["status"] == "ok" + assert data["workspace"] == str(ws) + assert data["grid_id"] == "itaituba/5km" + assert len(data["derived"]) == 5 + assert (ws / "run.json").exists() + + +def test_cli_export_json(tmp_path, capsys): + ws = tmp_path / "ws_export" + assert cli(["run", str(ITAITUBA), "--workspace", str(ws)]) == 0 + capsys.readouterr() + out_tif = tmp_path / "exported.tif" + assert cli(["export", str(ITAITUBA), "--workspace", str(ws), "--output", str(out_tif), "--json"]) == 0 + out = capsys.readouterr().out + data = json.loads(out) + assert data["status"] == "ok" + assert data["workspace"] == str(ws) + assert data["grid_id"] == "itaituba/5km" + assert "elevation" in data["variables"] + assert data["output"] == str(out_tif) + assert out_tif.exists() + + def test_sources_only_pipeline_without_grid_or_extent(tmp_path): tif = _raster(tmp_path / "raw.tif") path = tmp_path / "sources_only.toml" diff --git a/tests/test_pipeline_options.py b/tests/test_pipeline_options.py index 6f2d83d..8d0344f 100644 --- a/tests/test_pipeline_options.py +++ b/tests/test_pipeline_options.py @@ -20,8 +20,7 @@ from shapely.geometry import LineString, Point, box from disscube import CubeClient, Derivation, GridSpec, SpatialDerivation, SpatialSource, Variable -from disscube.config import plan, run -from disscube.config.runner import PipelineError +from disscube.pipeline import PipelineError, plan, run GEO = GridSpec(id="geo", type="local", crs="EPSG:4326", resolution=0.1, bbox=[-50.0, -10.0, -49.0, -9.0]) UTM = GridSpec(id="utm", type="local", crs="EPSG:31983", resolution=300, bbox=[0, 0, 3000, 3000]) diff --git a/tests/test_roundtrip.py b/tests/test_roundtrip.py index b8e81e2..3d140b2 100644 --- a/tests/test_roundtrip.py +++ b/tests/test_roundtrip.py @@ -194,3 +194,20 @@ def test_temporal_variable_roundtrip(tmp_path, cube): ) assert "time" in da.dims assert list(da.coords["time"].values) == [2020] + + +def test_cf_conventions_metadata_roundtrip(tmp_path, cube): + """Zarr storage enriches derived variables with CF-1.8 metadata and coordinates.""" + src = np.ones((2, 2), dtype=np.float32) + tif = _write_tif(tmp_path / "cf_test.tif", src) + + _derive(cube, tif, operator="mean", var_name="cf_var") + + da = cube.load("cf_var", grid_id="G") + assert da.attrs.get("conventions") == "CF-1.8" + assert da.attrs.get("operator") == "mean" + assert da.attrs.get("source_id") == "S" + assert da.attrs.get("grid_mapping") == "spatial_ref" + assert da.coords["x"].attrs.get("standard_name") == "projection_x_coordinate" + assert da.coords["y"].attrs.get("standard_name") == "projection_y_coordinate" + diff --git a/tools/import_bdc_tiles.py b/tools/import_bdc_tiles.py deleted file mode 100644 index f1373fe..0000000 --- a/tools/import_bdc_tiles.py +++ /dev/null @@ -1,25 +0,0 @@ -""" -tools/import_bdc_tiles.py - -Imports the BDC tiles (SM / MD / LG) as SpatialSources in the catalog. -One-time operation; may take a few minutes depending on the size of the -shapefiles. - -Usage: - python tools/import_bdc_tiles.py -""" - -from disscube.client import CubeClient -from disscube.utils.bdc_importer import import_bdc_grids - - -def main(): - cube = CubeClient(catalog="catalog.db", store="./data/") - - print("=== Importing BDC tiles (one-time, may be slow) ===") - import_bdc_grids(cube) # uses the BDC Grid V2 files bundled with DisSCube - print("=== Tiles BDC registrados ===") - - -if __name__ == "__main__": - main() diff --git a/tools/zarr_to_tif.py b/tools/zarr_to_tif.py deleted file mode 100644 index 4d0750b..0000000 --- a/tools/zarr_to_tif.py +++ /dev/null @@ -1,47 +0,0 @@ -import os -import sys - -import rioxarray # noqa: F401 (registers the .rio accessor on xarray objects) -import xarray as xr - - -def convert_zarr_to_tif(zarr_path, output_tif): - if not os.path.exists(zarr_path): - print(f"Error: Path {zarr_path} does not exist.") - return - - print(f"Opening Zarr: {zarr_path}") - try: - # Abre o dataset Zarr - ds = xr.open_zarr(zarr_path) - - # The Zarr written by DisSCube is a Dataset; take its first data variable. - var_names = list(ds.data_vars) - if not var_names: - print("Error: No data variables found in Zarr.") - return - - da = ds[var_names[0]] - - # Make sure the spatial dims are set correctly for rioxarray - if 'x' not in da.coords or 'y' not in da.coords: - print("Error: Spatial coordinates (x, y) not found in DataArray.") - return - - print(f"Converting variable '{var_names[0]}' to GeoTIFF...") - - # Escreve o TIFF - da.rio.to_raster(output_tif) - print(f"Success! Saved to {output_tif}") - - except (OSError, ValueError, KeyError) as e: - print(f"An error occurred: {e}") - -if __name__ == "__main__": - if len(sys.argv) < 3: - print("Usage: python tools/zarr_to_tif.py ") - sys.exit(1) - - zarr_path = sys.argv[1] - output_tif = sys.argv[2] - convert_zarr_to_tif(zarr_path, output_tif)