from __future__ import annotations import json import pathlib from datetime import datetime, timezone from pathlib import Path from typing import Any from uuid import UUID import uuid from fastapi import UploadFile from sqlalchemy.orm import Session from app.core.errors import AppError from app.models import Area, Dataset, DatasetVersion, Project from app.schemas.dataset import ( DatasetCreateResponse, DatasetStorageResponse, DatasetTemporalUpdate, DatasetVectorSummary, DatasetVersionRead, ) from app.services.geojson_service import parse_geojson_payload, load_dataset_text from app.services.raster_service import extract_raster_metadata from app.services.storage_service import StorageService from app.services.vector_feature_service import VectorFeatureService class DatasetService: VECTOR_EXTENSIONS = {".geojson", ".json"} RASTER_EXTENSIONS = {".tif", ".tiff", ".geotiff"} VECTOR_TYPES = {"vector", "geojson"} RASTER_TYPES = {"raster", "tif", "tiff", "geotiff"} VALID_DATASET_ROLES = {"source", "derived", "reference"} VALID_TEMPORAL_GRANULARITIES = {"snapshot", "day", "month", "year", "period"} @staticmethod def _normalize_datetime(value: datetime | None) -> datetime | None: if value is None: return None if value.tzinfo is None: return value.replace(tzinfo=timezone.utc) return value.astimezone(timezone.utc) @staticmethod def _validate_temporal_metadata( *, temporal_series_key: str | None, observed_at: datetime | None, valid_from: datetime | None, valid_to: datetime | None, temporal_granularity: str | None, source_version: str | None, ) -> dict[str, Any]: normalized_key = (temporal_series_key or "").strip() or None normalized_observed_at = DatasetService._normalize_datetime(observed_at) normalized_valid_from = DatasetService._normalize_datetime(valid_from) normalized_valid_to = DatasetService._normalize_datetime(valid_to) normalized_granularity = (temporal_granularity or "").strip().lower() or None normalized_source_version = (source_version or "").strip() or None if normalized_key and len(normalized_key) > 255: raise AppError(code="INVALID_TEMPORAL_METADATA", message="temporal_series_key is too long", status_code=400) if normalized_granularity and normalized_granularity not in DatasetService.VALID_TEMPORAL_GRANULARITIES: raise AppError( code="INVALID_TEMPORAL_METADATA", message="temporal_granularity must be snapshot, day, month, year or period", status_code=400, ) if normalized_valid_from and normalized_valid_to and normalized_valid_to < normalized_valid_from: raise AppError( code="INVALID_TEMPORAL_METADATA", message="valid_to must be on or after valid_from", status_code=400, ) if normalized_key and normalized_observed_at is None: raise AppError( code="INVALID_TEMPORAL_METADATA", message="observed_at is required when temporal_series_key is provided", status_code=400, ) if normalized_observed_at and normalized_key is None: raise AppError( code="INVALID_TEMPORAL_METADATA", message="temporal_series_key is required when observed_at is provided", status_code=400, ) return { "temporal_series_key": normalized_key, "observed_at": normalized_observed_at, "valid_from": normalized_valid_from, "valid_to": normalized_valid_to, "temporal_granularity": normalized_granularity, "source_version": normalized_source_version, } @staticmethod def _to_response(dataset: Dataset) -> DatasetCreateResponse: metadata_json = dataset.metadata_json if isinstance(dataset.metadata_json, dict) else {} return DatasetCreateResponse( id=dataset.id, name=dataset.name, dataset_type=dataset.dataset_type, source=dataset.source, dataset_role=dataset.dataset_role, source_name=dataset.source_name, reference_layer_name=dataset.reference_layer_name, source_metadata=dataset.source_metadata, provenance_metadata=dataset.provenance_metadata, imported_at=dataset.imported_at, temporal_series_key=dataset.temporal_series_key, observed_at=dataset.observed_at, valid_from=dataset.valid_from, valid_to=dataset.valid_to, temporal_granularity=dataset.temporal_granularity, source_version=dataset.source_version, project_id=dataset.project_id, area_id=dataset.area_id, storage_path=dataset.storage_path, original_filename=dataset.original_filename, stored_filename=dataset.stored_filename, content_type=dataset.content_type, size_bytes=dataset.size_bytes, checksum_sha256=dataset.checksum_sha256, crs=dataset.crs, bounds_json=dataset.bounds_json, metadata_json=dataset.metadata_json, vector_summary=DatasetService._extract_vector_summary(dataset.dataset_type, metadata_json), status=dataset.status, derived_from_dataset_id=dataset.derived_from_dataset_id, created_at=dataset.created_at, feature_count=metadata_json.get("feature_count"), ) @staticmethod def _canonical_dataset_type(dataset_type: str) -> str: normalized = (dataset_type or "").strip().lower() if normalized in DatasetService.VECTOR_TYPES: return "vector" if normalized in DatasetService.RASTER_TYPES: return "raster" raise AppError( code="INVALID_DATASET_TYPE", message="dataset_type must be 'vector' or 'raster' (or legacy 'geojson')", status_code=400, ) @staticmethod def _normalize_stored_dataset_type(dataset_type: str) -> str: normalized = (dataset_type or "").strip().lower() if normalized in DatasetService.VECTOR_TYPES: return "vector" if normalized in DatasetService.RASTER_TYPES: return "raster" return normalized @staticmethod def _is_vector_type(dataset_type: str) -> bool: return DatasetService._normalize_stored_dataset_type(dataset_type) == "vector" @staticmethod def _is_raster_type(dataset_type: str) -> bool: return DatasetService._normalize_stored_dataset_type(dataset_type) == "raster" @staticmethod def _normalize_dataset_role(dataset_role: str | None) -> str: normalized = (dataset_role or "").strip().lower() or "source" if normalized not in DatasetService.VALID_DATASET_ROLES: raise AppError( code="INVALID_DATASET_ROLE", message="dataset_role must be one of: source, derived, reference", status_code=400, ) return normalized @staticmethod def _extension_for_path(filename: str) -> str: return Path(filename).suffix.lower() @staticmethod def _validate_upload_filename(filename: str | None) -> str: if not filename: raise AppError(code="INVALID_UPLOAD", message="Missing file name", status_code=400) return filename @staticmethod def list_datasets(db: Session, project_id: UUID, limit: int = 50, offset: int = 0) -> tuple[list[DatasetCreateResponse], int]: total = db.query(Dataset).filter(Dataset.project_id == project_id).count() rows = ( db.query(Dataset) .filter(Dataset.project_id == project_id) .order_by(Dataset.created_at.desc()) .offset(offset) .limit(limit) .all() ) return [DatasetService._to_response(row) for row in rows], total @staticmethod def _extract_vector_summary(dataset_type: str, metadata_json: dict) -> DatasetVectorSummary | None: if not DatasetService._is_vector_type(dataset_type): return None if not isinstance(metadata_json, dict): return None return DatasetVectorSummary( feature_count=metadata_json.get("feature_count"), geometry_types=metadata_json.get("geometry_types"), bounds_json=metadata_json.get("bounds_json"), approximate_area_m2=metadata_json.get("approximate_area_m2"), crs=metadata_json.get("crs"), feature_geometry_count=metadata_json.get("feature_geometry_count"), invalid_features=metadata_json.get("invalid_features"), crs_assumed=metadata_json.get("crs_assumed"), ) @staticmethod def _extract_raster_bounds_json(metadata_json: dict[str, Any]) -> dict[str, float] | None: existing = metadata_json.get("bounds_json") if isinstance(existing, dict): return existing bounds = metadata_json.get("bounds") if isinstance(bounds, (list, tuple)) and len(bounds) == 4: return { "minx": float(bounds[0]), "miny": float(bounds[1]), "maxx": float(bounds[2]), "maxy": float(bounds[3]), } return None @staticmethod def _extract_raster_resolution_json(metadata_json: dict[str, Any]) -> dict[str, float] | None: existing = metadata_json.get("resolution_json") if isinstance(existing, dict): return existing resolution = metadata_json.get("resolution") if isinstance(resolution, (list, tuple)) and len(resolution) >= 2: return {"x": float(resolution[0]), "y": float(resolution[1])} return None @staticmethod def _extract_raster_bands_json(metadata_json: dict[str, Any]) -> dict[str, Any] | None: existing = metadata_json.get("bands_json") if isinstance(existing, dict): return existing bands_json: dict[str, Any] = {} if metadata_json.get("band_count") is not None: bands_json["band_count"] = int(metadata_json["band_count"]) if metadata_json.get("dtype") is not None: bands_json["dtype"] = metadata_json["dtype"] return bands_json or None @staticmethod async def upload_dataset( db: Session, project_id: UUID, file: UploadFile, dataset_type: str, source: str, dataset_role: str = "source", source_name: str | None = None, reference_layer_name: str | None = None, source_metadata: dict | None = None, provenance_metadata: dict | None = None, area_id: UUID | None = None, temporal_series_key: str | None = None, observed_at: datetime | None = None, valid_from: datetime | None = None, valid_to: datetime | None = None, temporal_granularity: str | None = None, source_version: str | None = None, ) -> DatasetCreateResponse: if not db.get(Project, project_id): raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404) filename = DatasetService._validate_upload_filename(file.filename) canonical_type = DatasetService._canonical_dataset_type(dataset_type) normalized_role = DatasetService._normalize_dataset_role(dataset_role) temporal = DatasetService._validate_temporal_metadata( temporal_series_key=temporal_series_key, observed_at=observed_at, valid_from=valid_from, valid_to=valid_to, temporal_granularity=temporal_granularity, source_version=source_version, ) normalized_source_name = source_name if normalized_role == "reference" and not normalized_source_name: normalized_source_name = "manual" if normalized_role == "reference" and canonical_type == "raster": raise AppError(code="INVALID_DATASET_ROLE", message="reference role requires vector dataset type", status_code=400) extension = DatasetService._extension_for_path(filename) if canonical_type == "vector" and extension not in DatasetService.VECTOR_EXTENSIONS: raise AppError(code="INVALID_UPLOAD", message="Vector uploads require .geojson or .json files", status_code=415) if canonical_type == "raster" and extension not in DatasetService.RASTER_EXTENSIONS: raise AppError( code="INVALID_UPLOAD", message="Raster uploads require .tif, .tiff or .geotiff files", status_code=415, ) raw = await file.read() storage_info = StorageService.persist_dataset_file( project_id=str(project_id), dataset_id=str(dataset_id := uuid.uuid4()), dataset_type=canonical_type, original_filename=filename, content=raw, content_type=file.content_type, ) metadata: dict[str, Any] = {} vector_payload: dict[str, Any] | None = None status = "uploaded" try: status = "validating" if canonical_type == "vector": try: text = raw.decode("utf-8") except UnicodeDecodeError as exc: raise AppError(code="INVALID_UPLOAD", message="Upload must be UTF-8 encoded", status_code=400) from exc metadata = parse_geojson_payload(text) vector_payload = json.loads(text) status = "ready" else: metadata = extract_raster_metadata(storage_info["storage_path"]) status = "ready" except ValueError as exc: status = "failed" StorageService.remove_dataset_file(storage_info["storage_path"]) raise AppError(code="INVALID_GEOJSON", message=str(exc), status_code=400) from exc except AppError as exc: if canonical_type == "raster" and exc.code == "RASTER_PROCESSING_UNAVAILABLE": status = "failed" metadata = { "processing_error": exc.message, "processing_code": exc.code, } else: StorageService.remove_dataset_file(storage_info["storage_path"]) raise bounds_json = metadata.get("bounds_json") if isinstance(metadata, dict) else None resolution_json = metadata.get("resolution_json") if isinstance(metadata, dict) else None bands_json = metadata.get("bands_json") if isinstance(metadata, dict) else None if canonical_type == "raster" and isinstance(metadata, dict): bounds_json = DatasetService._extract_raster_bounds_json(metadata) resolution_json = DatasetService._extract_raster_resolution_json(metadata) bands_json = DatasetService._extract_raster_bands_json(metadata) dataset = Dataset( id=dataset_id, project_id=project_id, area_id=area_id, name=filename, dataset_type=canonical_type, source=source, dataset_role=normalized_role, source_name=normalized_source_name, reference_layer_name=reference_layer_name if normalized_role == "reference" else None, source_metadata=source_metadata, provenance_metadata=provenance_metadata, imported_at=datetime.now(timezone.utc), **temporal, storage_path=storage_info["storage_path"], original_filename=storage_info["original_filename"], stored_filename=storage_info["stored_filename"], content_type=storage_info["content_type"], size_bytes=storage_info["size_bytes"], checksum_sha256=storage_info["checksum_sha256"], crs=metadata.get("crs") if isinstance(metadata, dict) else None, bounds_json=bounds_json, resolution_json=resolution_json, bands_json=bands_json, metadata_json=metadata, status=status, ) try: db.add(dataset) db.add( DatasetVersion( dataset_id=dataset.id, version=1, storage_path=dataset.storage_path, source_version=dataset.source_version, observed_at=dataset.observed_at, valid_from=dataset.valid_from, valid_to=dataset.valid_to, checksum_sha256=dataset.checksum_sha256, source_metadata=dataset.source_metadata, provenance_metadata=dataset.provenance_metadata, ) ) if canonical_type == "vector" and vector_payload is not None and status == "ready": feature_class = reference_layer_name if normalized_role == "reference" else None VectorFeatureService.persist_geojson_features( db=db, dataset_id=dataset.id, payload=vector_payload, feature_class=feature_class, commit=False, ) db.commit() db.refresh(dataset) except Exception: db.rollback() StorageService.remove_dataset_file(storage_info["storage_path"]) raise return DatasetService._to_response(dataset) @staticmethod def import_vector_bytes( db: Session, *, project_id: UUID, filename: str, content: bytes, source: str, source_name: str, dataset_role: str, reference_layer_name: str | None, source_metadata: dict[str, Any], provenance_metadata: dict[str, Any], area_id: UUID | None = None, temporal_series_key: str | None = None, observed_at: datetime | None = None, valid_from: datetime | None = None, valid_to: datetime | None = None, temporal_granularity: str | None = None, source_version: str | None = None, content_type: str = "application/geo+json", ) -> DatasetCreateResponse: if not db.get(Project, project_id): raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404) if area_id is not None: area = db.get(Area, area_id) if not area: raise AppError(code="AREA_NOT_FOUND", message="Area not found", status_code=404) if area.project_id != project_id: raise AppError(code="INVALID_DATASET_SCOPE", message="Area does not belong to this project", status_code=400) if not content: raise AppError(code="INVALID_UPLOAD", message="Vector artifact is empty", status_code=400) safe_filename = DatasetService._validate_upload_filename(filename) if DatasetService._extension_for_path(safe_filename) not in DatasetService.VECTOR_EXTENSIONS: raise AppError(code="INVALID_UPLOAD", message="Vector artifacts require .geojson or .json files", status_code=415) normalized_role = DatasetService._normalize_dataset_role(dataset_role) normalized_source_name = (source_name or "").strip() or ("manual" if normalized_role == "reference" else None) temporal = DatasetService._validate_temporal_metadata( temporal_series_key=temporal_series_key, observed_at=observed_at, valid_from=valid_from, valid_to=valid_to, temporal_granularity=temporal_granularity, source_version=source_version, ) try: text = content.decode("utf-8") except UnicodeDecodeError as exc: raise AppError(code="INVALID_UPLOAD", message="Vector artifact must be UTF-8 encoded", status_code=400) from exc try: metadata = parse_geojson_payload(text) vector_payload = json.loads(text) except (ValueError, json.JSONDecodeError) as exc: raise AppError(code="INVALID_GEOJSON", message=str(exc), status_code=400) from exc dataset_id = uuid.uuid4() storage_info = StorageService.persist_dataset_file( project_id=str(project_id), dataset_id=str(dataset_id), dataset_type="vector", original_filename=safe_filename, content=content, content_type=content_type, ) dataset = Dataset( id=dataset_id, project_id=project_id, area_id=area_id, name=safe_filename, dataset_type="vector", source=source, dataset_role=normalized_role, source_name=normalized_source_name, reference_layer_name=reference_layer_name if normalized_role == "reference" else None, source_metadata=source_metadata, provenance_metadata=provenance_metadata, imported_at=datetime.now(timezone.utc), **temporal, storage_path=storage_info["storage_path"], original_filename=storage_info["original_filename"], stored_filename=storage_info["stored_filename"], content_type=storage_info["content_type"], size_bytes=storage_info["size_bytes"], checksum_sha256=storage_info["checksum_sha256"], crs=metadata.get("crs"), bounds_json=metadata.get("bounds_json"), metadata_json=metadata, status="ready", ) try: db.add(dataset) db.add( DatasetVersion( dataset_id=dataset.id, version=1, storage_path=dataset.storage_path, source_version=dataset.source_version, observed_at=dataset.observed_at, valid_from=dataset.valid_from, valid_to=dataset.valid_to, checksum_sha256=dataset.checksum_sha256, source_metadata=dataset.source_metadata, provenance_metadata=dataset.provenance_metadata, ) ) VectorFeatureService.persist_geojson_features( db=db, dataset_id=dataset.id, payload=vector_payload, feature_class=reference_layer_name if normalized_role == "reference" else None, commit=False, ) db.commit() db.refresh(dataset) except Exception: db.rollback() StorageService.remove_dataset_file(storage_info["storage_path"]) raise return DatasetService._to_response(dataset) @staticmethod def import_raster_bytes( db: Session, *, project_id: UUID, filename: str, content: bytes, source: str, source_name: str, source_metadata: dict[str, Any], provenance_metadata: dict[str, Any], area_id: UUID | None = None, temporal_series_key: str | None = None, observed_at: datetime | None = None, valid_from: datetime | None = None, valid_to: datetime | None = None, temporal_granularity: str | None = None, source_version: str | None = None, content_type: str = "image/tiff", ) -> DatasetCreateResponse: if not db.get(Project, project_id): raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404) if area_id is not None: area = db.get(Area, area_id) if not area: raise AppError(code="AREA_NOT_FOUND", message="Area not found", status_code=404) if area.project_id != project_id: raise AppError(code="INVALID_DATASET_SCOPE", message="Area does not belong to this project", status_code=400) if not content: raise AppError(code="INVALID_UPLOAD", message="Raster artifact is empty", status_code=400) safe_filename = DatasetService._validate_upload_filename(filename) if DatasetService._extension_for_path(safe_filename) not in DatasetService.RASTER_EXTENSIONS: raise AppError(code="INVALID_UPLOAD", message="Raster artifacts require a GeoTIFF filename", status_code=415) temporal = DatasetService._validate_temporal_metadata( temporal_series_key=temporal_series_key, observed_at=observed_at, valid_from=valid_from, valid_to=valid_to, temporal_granularity=temporal_granularity, source_version=source_version, ) dataset_id = uuid.uuid4() storage_info = StorageService.persist_dataset_file( project_id=str(project_id), dataset_id=str(dataset_id), dataset_type="raster", original_filename=safe_filename, content=content, content_type=content_type, ) try: metadata = extract_raster_metadata(storage_info["storage_path"]) dataset = Dataset( id=dataset_id, project_id=project_id, area_id=area_id, name=safe_filename, dataset_type="raster", source=source, dataset_role="source", source_name=source_name, source_metadata=source_metadata, provenance_metadata=provenance_metadata, imported_at=datetime.now(timezone.utc), **temporal, storage_path=storage_info["storage_path"], original_filename=storage_info["original_filename"], stored_filename=storage_info["stored_filename"], content_type=storage_info["content_type"], size_bytes=storage_info["size_bytes"], checksum_sha256=storage_info["checksum_sha256"], crs=metadata.get("crs"), bounds_json=DatasetService._extract_raster_bounds_json(metadata), resolution_json=DatasetService._extract_raster_resolution_json(metadata), bands_json=DatasetService._extract_raster_bands_json(metadata), metadata_json=metadata, status="ready", ) db.add(dataset) db.add( DatasetVersion( dataset_id=dataset.id, version=1, storage_path=dataset.storage_path, source_version=dataset.source_version, observed_at=dataset.observed_at, valid_from=dataset.valid_from, valid_to=dataset.valid_to, checksum_sha256=dataset.checksum_sha256, source_metadata=dataset.source_metadata, provenance_metadata=dataset.provenance_metadata, ) ) db.commit() db.refresh(dataset) return DatasetService._to_response(dataset) except Exception: db.rollback() StorageService.remove_dataset_file(storage_info["storage_path"]) raise @staticmethod def import_partitioned_vector_artifact( db: Session, *, project_id: UUID, area_id: UUID, artifact_path: str | Path, partition_paths: list[str | Path], original_filename: str, source: str, dataset_role: str, source_name: str, reference_layer_name: str | None, metadata_json: dict[str, Any], source_metadata: dict[str, Any], provenance_metadata: dict[str, Any], temporal_series_key: str, observed_at: datetime, temporal_granularity: str = "snapshot", source_version: str | None = None, batch_size: int = 1000, ) -> DatasetCreateResponse: if not db.get(Project, project_id): raise AppError(code="PROJECT_NOT_FOUND", message="Project not found", status_code=404) area = db.get(Area, area_id) if not area: raise AppError(code="AREA_NOT_FOUND", message="Area not found", status_code=404) if area.project_id != project_id: raise AppError(code="INVALID_DATASET_SCOPE", message="Area does not belong to this project", status_code=400) if not partition_paths: raise AppError( code="INVALID_GEOJSON_PARTITIONS", message="At least one GeoJSON partition is required", status_code=400, ) filename = DatasetService._validate_upload_filename(original_filename) if DatasetService._extension_for_path(filename) not in DatasetService.VECTOR_EXTENSIONS: raise AppError(code="INVALID_UPLOAD", message="Vector artifacts require .geojson or .json files", status_code=415) normalized_role = DatasetService._normalize_dataset_role(dataset_role) temporal = DatasetService._validate_temporal_metadata( temporal_series_key=temporal_series_key, observed_at=observed_at, valid_from=observed_at, valid_to=None, temporal_granularity=temporal_granularity, source_version=source_version, ) metadata = dict(metadata_json) expected_feature_count = int(metadata.get("feature_count") or 0) if expected_feature_count <= 0: raise AppError( code="INVALID_GEOJSON_PARTITIONS", message="Partition metadata must declare a positive feature_count", status_code=400, ) dataset_id = uuid.uuid4() storage_info = StorageService.persist_dataset_file_from_path( project_id=str(project_id), dataset_id=str(dataset_id), dataset_type="vector", original_filename=filename, source_path=artifact_path, content_type="application/geo+json", ) dataset = Dataset( id=dataset_id, project_id=project_id, area_id=area_id, name=filename, dataset_type="vector", source=source, dataset_role=normalized_role, source_name=source_name, reference_layer_name=reference_layer_name if normalized_role == "reference" else None, source_metadata=source_metadata, provenance_metadata=provenance_metadata, imported_at=datetime.now(timezone.utc), **temporal, storage_path=storage_info["storage_path"], original_filename=storage_info["original_filename"], stored_filename=storage_info["stored_filename"], content_type=storage_info["content_type"], size_bytes=storage_info["size_bytes"], checksum_sha256=storage_info["checksum_sha256"], crs=str(metadata.get("crs") or "EPSG:4326"), bounds_json=metadata.get("bounds_json"), metadata_json=metadata, status="ready", ) try: db.add(dataset) db.add( DatasetVersion( dataset_id=dataset.id, version=1, storage_path=dataset.storage_path, source_version=dataset.source_version, observed_at=dataset.observed_at, valid_from=dataset.valid_from, checksum_sha256=dataset.checksum_sha256, source_metadata=dataset.source_metadata, provenance_metadata=dataset.provenance_metadata, ) ) persisted_count = VectorFeatureService.persist_geojson_partitions( db, dataset.id, partition_paths, feature_class=reference_layer_name if normalized_role == "reference" else None, batch_size=batch_size, ) if persisted_count != expected_feature_count: raise AppError( code="PARTITION_FEATURE_COUNT_MISMATCH", message=( f"Regional artifact declares {expected_feature_count} features but " f"{persisted_count} queryable features were indexed" ), status_code=400, ) db.commit() db.refresh(dataset) except Exception: db.rollback() StorageService.remove_dataset_file(storage_info["storage_path"]) raise return DatasetService._to_response(dataset) @staticmethod def refresh_metadata(db: Session, dataset_id: UUID) -> DatasetCreateResponse: dataset = DatasetService._get_dataset(db, dataset_id) if not dataset.storage_path: raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) if not Path(dataset.storage_path).exists(): raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) try: if DatasetService._is_vector_type(dataset.dataset_type): metadata = parse_geojson_payload(load_dataset_text(dataset.storage_path)) elif DatasetService._is_raster_type(dataset.dataset_type): metadata = extract_raster_metadata(dataset.storage_path) else: raise AppError(code="INVALID_DATASET_TYPE", message="Cannot refresh metadata for this dataset type", status_code=400) dataset.status = "ready" except ValueError as exc: dataset.status = "failed" raise AppError(code="INVALID_GEOJSON", message=str(exc), status_code=400) from exc except AppError as exc: if DatasetService._is_raster_type(dataset.dataset_type) and exc.code == "RASTER_PROCESSING_UNAVAILABLE": dataset.status = "failed" metadata = {"processing_error": exc.message, "processing_code": exc.code} else: dataset.status = "failed" raise bounds_json = metadata.get("bounds_json") if isinstance(metadata, dict) else dataset.bounds_json resolution_json = metadata.get("resolution_json") if isinstance(metadata, dict) else dataset.resolution_json bands_json = metadata.get("bands_json") if isinstance(metadata, dict) else dataset.bands_json if DatasetService._is_raster_type(dataset.dataset_type) and isinstance(metadata, dict): bounds_json = DatasetService._extract_raster_bounds_json(metadata) resolution_json = DatasetService._extract_raster_resolution_json(metadata) bands_json = DatasetService._extract_raster_bands_json(metadata) dataset.crs = metadata.get("crs") if isinstance(metadata, dict) else dataset.crs dataset.bounds_json = bounds_json dataset.metadata_json = metadata dataset.resolution_json = resolution_json dataset.bands_json = bands_json db.add(dataset) db.commit() db.refresh(dataset) return DatasetService._to_response(dataset) @staticmethod def update_temporal_metadata(db: Session, dataset_id: UUID, payload: DatasetTemporalUpdate) -> DatasetCreateResponse: dataset = DatasetService._get_dataset(db, dataset_id) temporal = DatasetService._validate_temporal_metadata(**payload.model_dump()) if all(getattr(dataset, field) == value for field, value in temporal.items()): return DatasetService._to_response(dataset) for field, value in temporal.items(): setattr(dataset, field, value) latest_version = ( db.query(DatasetVersion) .filter(DatasetVersion.dataset_id == dataset.id) .order_by(DatasetVersion.version.desc()) .first() ) db.add(dataset) db.add( DatasetVersion( dataset_id=dataset.id, version=(latest_version.version + 1) if latest_version else 1, storage_path=dataset.storage_path, source_version=dataset.source_version, observed_at=dataset.observed_at, valid_from=dataset.valid_from, valid_to=dataset.valid_to, checksum_sha256=dataset.checksum_sha256, source_metadata=dataset.source_metadata, provenance_metadata=dataset.provenance_metadata, ) ) db.commit() db.refresh(dataset) return DatasetService._to_response(dataset) @staticmethod def list_versions(db: Session, dataset_id: UUID) -> list[DatasetVersionRead]: DatasetService._get_dataset(db, dataset_id) rows = ( db.query(DatasetVersion) .filter(DatasetVersion.dataset_id == dataset_id) .order_by(DatasetVersion.version.desc()) .all() ) return [DatasetVersionRead.model_validate(row) for row in rows] @staticmethod def get_dataset(db: Session, dataset_id: UUID) -> Dataset: dataset = db.get(Dataset, dataset_id) if not dataset: raise AppError(code="DATASET_NOT_FOUND", message="Dataset not found", status_code=404) return dataset @staticmethod def _get_dataset(db: Session, dataset_id: UUID) -> Dataset: return DatasetService.get_dataset(db, dataset_id) @staticmethod def get_dataset_geojson(db: Session, dataset_id: UUID) -> dict: dataset = DatasetService._get_dataset(db, dataset_id) if not DatasetService._is_vector_type(dataset.dataset_type): raise AppError(code="INVALID_DATASET_TYPE", message="Dataset is not a vector dataset", status_code=400) if not dataset.storage_path: raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) if not pathlib.Path(dataset.storage_path).exists(): raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) raw = load_dataset_text(dataset.storage_path) try: return json.loads(raw) except Exception as exc: raise AppError(code="INVALID_GEOJSON", message="Stored dataset is not valid JSON", status_code=500) from exc @staticmethod def inspect_vector_dataset(db: Session, dataset_id: UUID) -> dict[str, Any]: dataset = DatasetService._get_dataset(db, dataset_id) if not DatasetService._is_vector_type(dataset.dataset_type): raise AppError(code="INVALID_DATASET_TYPE", message="Dataset is not a vector dataset", status_code=400) if not dataset.storage_path or not Path(dataset.storage_path).exists(): raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) metadata = dataset.metadata_json or {} if not isinstance(metadata, dict): metadata = {} summary = DatasetService._extract_vector_summary(dataset.dataset_type, metadata) return { "dataset": { "id": str(dataset.id), "name": dataset.name, "dataset_type": dataset.dataset_type, "status": dataset.status, "source": dataset.source, "storage": DatasetStorageResponse( original_filename=dataset.original_filename, stored_filename=dataset.stored_filename, content_type=dataset.content_type, size_bytes=dataset.size_bytes, checksum_sha256=dataset.checksum_sha256, ).model_dump(), "feature_count": metadata.get("feature_count"), "crs": metadata.get("crs"), }, "summary": summary.model_dump() if summary else None, "metadata": metadata, } @staticmethod def vector_summary(db: Session, dataset_id: UUID) -> dict[str, Any]: dataset = DatasetService._get_dataset(db, dataset_id) if not DatasetService._is_vector_type(dataset.dataset_type): raise AppError(code="INVALID_DATASET_TYPE", message="Dataset is not a vector dataset", status_code=400) metadata = dataset.metadata_json or {} if not isinstance(metadata, dict): metadata = {} summary = DatasetService._extract_vector_summary(dataset.dataset_type, metadata) if not summary: raise AppError(code="INVALID_GEOJSON", message="Vector summary unavailable", status_code=422) return summary.model_dump() @staticmethod def raster_metadata(db: Session, dataset_id: UUID) -> dict[str, Any]: dataset = DatasetService._get_dataset(db, dataset_id) if not DatasetService._is_raster_type(dataset.dataset_type): raise AppError(code="INVALID_DATASET_TYPE", message="Dataset is not a raster dataset", status_code=400) if not dataset.storage_path: raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) if not Path(dataset.storage_path).exists(): raise AppError(code="DATASET_FILE_MISSING", message="Stored file missing", status_code=404) if isinstance(dataset.metadata_json, dict) and dataset.metadata_json.get("driver"): return dataset.metadata_json metadata = extract_raster_metadata(dataset.storage_path) dataset.metadata_json = dict(dataset.metadata_json or {}) dataset.metadata_json.update(metadata) dataset.status = "ready" db.add(dataset) db.commit() db.refresh(dataset) return metadata