From c0943fe7d4d017a2c6bf3339afd1785e291f261b Mon Sep 17 00:00:00 2001 From: Codex Date: Tue, 14 Jul 2026 15:19:42 +0200 Subject: [PATCH] feat: add temporal Mol explorer --- CHANGELOG.md | 9 + backend/README.md | 30 + ...02607140001_temporal_dataset_foundation.py | 72 +++ backend/app/api/routes/__init__.py | 2 +- backend/app/api/routes/datasets.py | 49 +- backend/app/api/routes/temporal.py | 29 + backend/app/main.py | 3 +- backend/app/models/entities.py | 35 +- backend/app/schemas/__init__.py | 2 + backend/app/schemas/dataset.py | 32 + backend/app/schemas/operations.py | 11 + backend/app/schemas/temporal.py | 71 +++ backend/app/services/dataset_service.py | 279 ++++++--- backend/app/services/demo_workflow_service.py | 17 +- .../app/services/raster_operations_service.py | 31 +- .../app/services/temporal_analysis_service.py | 310 ++++++++++ .../app/services/vector_feature_service.py | 103 +++- .../app/services/vector_operations_service.py | 26 +- .../tests/test_raster_operations_service.py | 20 +- .../test_sprint187_temporal_map_foundation.py | 250 ++++++++ deploy/unraid/Dockerfile.all-in-one | 2 + docs/API_CONTRACTS.md | 50 ++ docs/CODEX_EXECUTION_LOG.md | 26 + docs/DATABASE_IMPLEMENTATION_PLAN.md | 19 + docs/DATASET_STRATEGY.md | 17 + docs/DATA_SOURCES.md | 26 + frontend/README.md | 13 + frontend/src/components/GeoMap.tsx | 4 + frontend/src/components/map/MapWorkspace.tsx | 340 +++++++++-- frontend/src/hooks/useTemporalComparison.ts | 62 ++ frontend/src/services/api/datasets.ts | 24 + frontend/src/services/api/index.ts | 1 + frontend/src/services/api/temporal.ts | 13 + frontend/src/styles/app.css | 158 ++++- frontend/src/types.ts | 109 ++++ scripts/provision_mol_context_layers.py | 38 ++ scripts/provision_mol_historical_landuse.py | 552 ++++++++++++++++++ .../provision_mol_municipality_workspace.py | 56 ++ scripts/provision_mol_population_history.py | 337 +++++++++++ 39 files changed, 3065 insertions(+), 163 deletions(-) create mode 100644 backend/alembic/versions/202607140001_temporal_dataset_foundation.py create mode 100644 backend/app/api/routes/temporal.py create mode 100644 backend/app/schemas/temporal.py create mode 100644 backend/app/services/temporal_analysis_service.py create mode 100644 backend/tests/test_sprint187_temporal_map_foundation.py create mode 100644 frontend/src/hooks/useTemporalComparison.ts create mode 100644 frontend/src/services/api/temporal.ts create mode 100644 scripts/provision_mol_historical_landuse.py create mode 100644 scripts/provision_mol_population_history.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 48d06a2c..7b5fcf9c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,15 @@ # Changelog +## Sprint 187 Temporal Mol explorer (2026-07-14) + +- Added first-class temporal dataset metadata and immutable dataset-version provenance for uploaded and derived vector/raster datasets. +- Added project temporal-series discovery and bounded snapshot comparison APIs with explicit observation dates, metric deltas, warnings and optional stable-identity object changes. +- Extended bbox selection with source-governed PostGIS aggregations so population is reported as inhabitants and land cover as intersected hectares instead of misleading feature counts. +- Added a calm Latest state/Evolution flow to the map-first explorer, including period selection, metric comparison and added/removed/modified overlays where source identities support them. +- Added explicit operator provisioners for official Statbel Mol population snapshots (2021-2025) and Digitaal Vlaanderen historical land-use snapshots (1778, 1873 and 1969); no source is fetched during application startup. +- Preserved methodological honesty: partial statistical sectors are labelled area-weighted estimates, historical land-use identity changes are not fabricated and all source URLs, versions and processing limitations are persisted. + ## Sprint 186 Map-first Mol geographic explorer (2026-07-14) - Replaced the default dashboard entry with a calm map-first workflow: choose a real data theme, drag a rectangle, query PostGIS automatically and review results. diff --git a/backend/README.md b/backend/README.md index 61957897..95528d12 100644 --- a/backend/README.md +++ b/backend/README.md @@ -935,6 +935,36 @@ the same bbox-selected FeatureCollection as a normal export record with `export_type="vector_selection_geojson"`. This creates a handoff artifact only; it does not create a derived dataset. +## Temporal Mol data and evolution + +Dataset uploads accept `temporal_series_key`, `observed_at`, `valid_from`, +`valid_to`, `temporal_granularity` and `source_version`. Every new source or +derived dataset also writes dataset version 1 in the same transaction. + +After the Mol municipality workspace is available, import the official source +snapshots explicitly: + +```bash +docker exec geointel python /app/scripts/provision_mol_population_history.py +docker exec geointel python /app/scripts/provision_mol_historical_landuse.py +``` + +The first command imports Statbel sector population for 2021-2025. The second +imports Digitaal Vlaanderen historical land use for 1778, 1873 and 1969. Both +are idempotent, use the normal API/DatasetService flow and retain fetched +artifacts in persistent operator storage. They never run on app startup. + +Historical land-use work can be bounded explicitly: + +```bash +docker exec geointel python /app/scripts/provision_mol_historical_landuse.py --years 1778 1969 --themes forest water +``` + +`GET /api/v1/projects/{project_id}/temporal/series` discovers the series and +`POST /api/v1/projects/{project_id}/temporal/compare` compares two snapshots +inside one EPSG:4326 bbox. Partial statistical sectors are estimates; old map +editions without stable identities do not produce invented object changes. + ## Helpful repository scripts - `bash scripts/backend_install.sh` diff --git a/backend/alembic/versions/202607140001_temporal_dataset_foundation.py b/backend/alembic/versions/202607140001_temporal_dataset_foundation.py new file mode 100644 index 00000000..600c937e --- /dev/null +++ b/backend/alembic/versions/202607140001_temporal_dataset_foundation.py @@ -0,0 +1,72 @@ +"""Add temporal dataset metadata and durable dataset-version provenance.""" + +from alembic import op +import sqlalchemy as sa + + +revision = "202607140001" +down_revision = "202606120900" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.add_column("datasets", sa.Column("temporal_series_key", sa.String(length=255), nullable=True)) + op.add_column("datasets", sa.Column("observed_at", sa.DateTime(timezone=True), nullable=True)) + op.add_column("datasets", sa.Column("valid_from", sa.DateTime(timezone=True), nullable=True)) + op.add_column("datasets", sa.Column("valid_to", sa.DateTime(timezone=True), nullable=True)) + op.add_column("datasets", sa.Column("temporal_granularity", sa.String(length=32), nullable=True)) + op.add_column("datasets", sa.Column("source_version", sa.String(length=120), nullable=True)) + + op.add_column("dataset_versions", sa.Column("source_version", sa.String(length=120), nullable=True)) + op.add_column("dataset_versions", sa.Column("observed_at", sa.DateTime(timezone=True), nullable=True)) + op.add_column("dataset_versions", sa.Column("valid_from", sa.DateTime(timezone=True), nullable=True)) + op.add_column("dataset_versions", sa.Column("valid_to", sa.DateTime(timezone=True), nullable=True)) + op.add_column("dataset_versions", sa.Column("checksum_sha256", sa.String(length=64), nullable=True)) + op.add_column("dataset_versions", sa.Column("source_metadata", sa.JSON(), nullable=True)) + op.add_column("dataset_versions", sa.Column("provenance_metadata", sa.JSON(), nullable=True)) + + op.create_index( + "ix_datasets_project_temporal_series_observed", + "datasets", + ["project_id", "temporal_series_key", "observed_at"], + ) + op.create_index("ix_dataset_versions_dataset_version", "dataset_versions", ["dataset_id", "version"], unique=True) + op.create_index( + "ix_vector_features_dataset_source_feature", + "vector_features", + ["dataset_id", "source_feature_id"], + ) + op.create_check_constraint( + "ck_datasets_temporal_valid_range", + "datasets", + "valid_to IS NULL OR valid_from IS NULL OR valid_to >= valid_from", + ) + op.create_check_constraint( + "ck_dataset_versions_temporal_valid_range", + "dataset_versions", + "valid_to IS NULL OR valid_from IS NULL OR valid_to >= valid_from", + ) + + +def downgrade() -> None: + op.drop_constraint("ck_dataset_versions_temporal_valid_range", "dataset_versions", type_="check") + op.drop_constraint("ck_datasets_temporal_valid_range", "datasets", type_="check") + op.drop_index("ix_vector_features_dataset_source_feature", table_name="vector_features") + op.drop_index("ix_dataset_versions_dataset_version", table_name="dataset_versions") + op.drop_index("ix_datasets_project_temporal_series_observed", table_name="datasets") + + op.drop_column("dataset_versions", "provenance_metadata") + op.drop_column("dataset_versions", "source_metadata") + op.drop_column("dataset_versions", "checksum_sha256") + op.drop_column("dataset_versions", "valid_to") + op.drop_column("dataset_versions", "valid_from") + op.drop_column("dataset_versions", "observed_at") + op.drop_column("dataset_versions", "source_version") + + op.drop_column("datasets", "source_version") + op.drop_column("datasets", "temporal_granularity") + op.drop_column("datasets", "valid_to") + op.drop_column("datasets", "valid_from") + op.drop_column("datasets", "observed_at") + op.drop_column("datasets", "temporal_series_key") diff --git a/backend/app/api/routes/__init__.py b/backend/app/api/routes/__init__.py index 43122e71..0042096c 100644 --- a/backend/app/api/routes/__init__.py +++ b/backend/app/api/routes/__init__.py @@ -1 +1 @@ -__all__ = ["analysis", "areas", "datasets", "health", "projects", "exports", "jobs", "external", "qa"] +__all__ = ["analysis", "areas", "datasets", "health", "projects", "exports", "jobs", "external", "qa", "temporal"] diff --git a/backend/app/api/routes/datasets.py b/backend/app/api/routes/datasets.py index 80843f72..73478131 100644 --- a/backend/app/api/routes/datasets.py +++ b/backend/app/api/routes/datasets.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +from datetime import datetime from typing import Any from uuid import UUID from uuid import UUID as _UUID @@ -30,7 +31,7 @@ from app.schemas import ( VectorSelectionResponse, ) from app.schemas.job import JobCreate -from app.schemas.dataset import DatasetCreateResponse +from app.schemas.dataset import DatasetCreateResponse, DatasetTemporalUpdate from app.schemas.operations import VectorOperationResult from app.services.job_service import JobService from app.services.raster_operations_service import RasterOperationsService @@ -87,6 +88,12 @@ async def upload_dataset( reference_layer_name: str | None = Form(None), source_metadata_json: str | None = Form(None), provenance_metadata_json: str | None = Form(None), + temporal_series_key: str | None = Form(None), + observed_at: datetime | None = Form(None), + valid_from: datetime | None = Form(None), + valid_to: datetime | None = Form(None), + temporal_granularity: str | None = Form(None), + source_version: str | None = Form(None), db: Session = Depends(get_db), ): if area_id is not None: @@ -108,6 +115,12 @@ async def upload_dataset( source_metadata=_parse_metadata_json(source_metadata_json, "source_metadata_json"), provenance_metadata=_parse_metadata_json(provenance_metadata_json, "provenance_metadata_json"), area_id=area_id, + 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, ) return envelope(created.model_dump()) @@ -135,6 +148,33 @@ def get_dataset( return envelope(DatasetCreateResponse.model_validate(dataset).model_dump()) +@router.patch("/datasets/{dataset_id}/temporal", response_model=dict) +def update_dataset_temporal_metadata( + project_id: UUID, + dataset_id: UUID, + payload: DatasetTemporalUpdate, + db: Session = Depends(get_db), +): + dataset = DatasetService.get_dataset(db, dataset_id) + if dataset.project_id != project_id: + raise HTTPException(status_code=404, detail="Dataset not found") + updated = DatasetService.update_temporal_metadata(db, dataset_id, payload) + return envelope(updated.model_dump()) + + +@router.get("/datasets/{dataset_id}/versions", response_model=dict) +def list_dataset_versions( + project_id: UUID, + dataset_id: UUID, + db: Session = Depends(get_db), +): + dataset = DatasetService.get_dataset(db, dataset_id) + if dataset.project_id != project_id: + raise HTTPException(status_code=404, detail="Dataset not found") + versions = DatasetService.list_versions(db, dataset_id) + return envelope({"items": [item.model_dump() for item in versions], "total": len(versions)}) + + @router.post("/datasets/{dataset_id}/metadata/refresh", response_model=dict) def refresh_dataset_metadata( project_id: UUID, @@ -203,6 +243,13 @@ def select_vector_features( bbox=payload.bbox.model_dump(), limit=payload.limit, ) + if isinstance(dataset.source_metadata, dict) and dataset.source_metadata.get("selection_aggregation"): + result["summary"] = VectorFeatureService.summarize_features_by_bbox( + db, + dataset=dataset, + bbox=payload.bbox.model_dump(), + total_feature_count=result.get("total_feature_count"), + ) return envelope(VectorSelectionResponse(**result).model_dump(exclude_none=True)) diff --git a/backend/app/api/routes/temporal.py b/backend/app/api/routes/temporal.py new file mode 100644 index 00000000..ac576de6 --- /dev/null +++ b/backend/app/api/routes/temporal.py @@ -0,0 +1,29 @@ +from __future__ import annotations + +from uuid import UUID + +from fastapi import APIRouter, Depends +from sqlalchemy.orm import Session + +from app.db.session import get_db +from app.schemas.temporal import TemporalComparisonRequest +from app.services.temporal_analysis_service import TemporalAnalysisService +from app.utils.response import envelope + + +router = APIRouter(prefix="/projects/{project_id}/temporal", tags=["temporal"]) + + +@router.get("/series", response_model=dict) +def list_temporal_series(project_id: UUID, db: Session = Depends(get_db)): + series = TemporalAnalysisService.list_series(db, project_id) + return envelope({"items": [item.model_dump() for item in series], "total": len(series)}) + + +@router.post("/compare", response_model=dict) +def compare_temporal_snapshots( + project_id: UUID, + payload: TemporalComparisonRequest, + db: Session = Depends(get_db), +): + return envelope(TemporalAnalysisService.compare(db, project_id=project_id, payload=payload).model_dump()) diff --git a/backend/app/main.py b/backend/app/main.py index 666f23f3..32a51b67 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -5,7 +5,7 @@ from fastapi.exceptions import RequestValidationError from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import JSONResponse -from app.api.routes import analysis, areas, datasets, demo, detection, exports, external, health, jobs, projects, qa, quality_checks, segmentation +from app.api.routes import analysis, areas, datasets, demo, detection, exports, external, health, jobs, projects, qa, quality_checks, segmentation, temporal from app.core.config import get_settings from app.core.errors import AppError from app.core.logging import configure_logging @@ -57,6 +57,7 @@ def create_app() -> FastAPI: app.include_router(qa.router, prefix=settings.api_prefix) app.include_router(detection.router, prefix=settings.api_prefix) app.include_router(segmentation.router, prefix=settings.api_prefix) + app.include_router(temporal.router, prefix=settings.api_prefix) @app.exception_handler(AppError) async def app_error(request: Request, exc: AppError): # noqa: ARG001 diff --git a/backend/app/models/entities.py b/backend/app/models/entities.py index 5d441731..d283a795 100644 --- a/backend/app/models/entities.py +++ b/backend/app/models/entities.py @@ -4,7 +4,7 @@ import uuid from datetime import datetime from geoalchemy2 import Geometry -from sqlalchemy import DateTime, ForeignKey, Float, Index, JSON, String, Text, func +from sqlalchemy import CheckConstraint, DateTime, ForeignKey, Float, Index, JSON, String, Text, func from sqlalchemy.sql.sqltypes import Integer from sqlalchemy.dialects.postgresql import UUID from sqlalchemy.orm import Mapped, mapped_column, relationship @@ -44,6 +44,18 @@ class Area(Base): class Dataset(Base): __tablename__ = "datasets" + __table_args__ = ( + CheckConstraint( + "valid_to IS NULL OR valid_from IS NULL OR valid_to >= valid_from", + name="ck_datasets_temporal_valid_range", + ), + Index( + "ix_datasets_project_temporal_series_observed", + "project_id", + "temporal_series_key", + "observed_at", + ), + ) id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) project_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("projects.id", ondelete="CASCADE"), nullable=False) @@ -73,6 +85,12 @@ class Dataset(Base): source_metadata: Mapped[dict | None] = mapped_column(JSON, nullable=True) provenance_metadata: Mapped[dict | None] = mapped_column(JSON, nullable=True) imported_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) + temporal_series_key: Mapped[str | None] = mapped_column(String(255), nullable=True) + observed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + valid_from: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + valid_to: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + temporal_granularity: Mapped[str | None] = mapped_column(String(32), nullable=True) + source_version: Mapped[str | None] = mapped_column(String(120), nullable=True) status: Mapped[str] = mapped_column(String(32), default="uploaded") created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now(), onupdate=func.now()) @@ -92,11 +110,25 @@ class Dataset(Base): class DatasetVersion(Base): __tablename__ = "dataset_versions" + __table_args__ = ( + CheckConstraint( + "valid_to IS NULL OR valid_from IS NULL OR valid_to >= valid_from", + name="ck_dataset_versions_temporal_valid_range", + ), + Index("ix_dataset_versions_dataset_version", "dataset_id", "version", unique=True), + ) id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) dataset_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("datasets.id", ondelete="CASCADE"), nullable=False) version: Mapped[int] = mapped_column(Integer, default=1) storage_path: Mapped[str | None] = mapped_column(String(500), nullable=True) + source_version: Mapped[str | None] = mapped_column(String(120), nullable=True) + observed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + valid_from: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + valid_to: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + checksum_sha256: Mapped[str | None] = mapped_column(String(64), nullable=True) + source_metadata: Mapped[dict | None] = mapped_column(JSON, nullable=True) + provenance_metadata: Mapped[dict | None] = mapped_column(JSON, nullable=True) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) dataset: Mapped[Dataset] = relationship("Dataset", back_populates="versions") @@ -107,6 +139,7 @@ class VectorFeature(Base): __table_args__ = ( Index("ix_vector_features_dataset_id", "dataset_id"), Index("ix_vector_features_geometry", "geometry", postgresql_using="gist"), + Index("ix_vector_features_dataset_source_feature", "dataset_id", "source_feature_id"), ) id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) diff --git a/backend/app/schemas/__init__.py b/backend/app/schemas/__init__.py index 3e4f324f..1dfead60 100644 --- a/backend/app/schemas/__init__.py +++ b/backend/app/schemas/__init__.py @@ -77,6 +77,7 @@ from .operations import ( VectorSelectionDeriveRequest, VectorSelectionRequest, VectorSelectionResponse, + VectorSelectionSummary, VectorStatsRequest, VectorStatsResponse, ) @@ -134,6 +135,7 @@ __all__ = [ "VectorSelectionDeriveRequest", "VectorSelectionRequest", "VectorSelectionResponse", + "VectorSelectionSummary", "RasterClipRequest", "RasterStatsResponse", "RasterReprojectRequest", diff --git a/backend/app/schemas/dataset.py b/backend/app/schemas/dataset.py index c52d5ed1..5495380d 100644 --- a/backend/app/schemas/dataset.py +++ b/backend/app/schemas/dataset.py @@ -36,6 +36,12 @@ class DatasetCreateResponse(BaseModel): source_metadata: dict | None = None provenance_metadata: dict | None = None imported_at: datetime | 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 project_id: UUID area_id: UUID | None = None storage_path: str | None = None @@ -70,6 +76,32 @@ class DatasetMetadataRefresh(BaseModel): crs: str | None = None +class DatasetTemporalUpdate(BaseModel): + temporal_series_key: str + observed_at: datetime + valid_from: datetime | None = None + valid_to: datetime | None = None + temporal_granularity: str = "snapshot" + source_version: str | None = None + + +class DatasetVersionRead(BaseModel): + id: UUID + dataset_id: UUID + version: int + storage_path: str | None = None + source_version: str | None = None + observed_at: datetime | None = None + valid_from: datetime | None = None + valid_to: datetime | None = None + checksum_sha256: str | None = None + source_metadata: dict | None = None + provenance_metadata: dict | None = None + created_at: datetime | None = None + + model_config = {"from_attributes": True} + + class ExportRequest(BaseModel): dataset_id: UUID name: str | None = None diff --git a/backend/app/schemas/operations.py b/backend/app/schemas/operations.py index 199cf50f..7a8cae01 100644 --- a/backend/app/schemas/operations.py +++ b/backend/app/schemas/operations.py @@ -217,6 +217,16 @@ class VectorSelectionDeriveRequest(VectorSelectionRequest): output_name: str | None = None +class VectorSelectionSummary(BaseModel): + metric_label: str + metric_value: float + metric_unit: str + aggregation_method: str + feature_count: int + is_estimate: bool = False + warning: str | None = None + + class VectorSelectionResponse(BaseModel): selection_bbox: VectorSelectionBBox feature_count: int @@ -224,3 +234,4 @@ class VectorSelectionResponse(BaseModel): limit: int truncated: bool geojson: dict + summary: VectorSelectionSummary | None = None diff --git a/backend/app/schemas/temporal.py b/backend/app/schemas/temporal.py new file mode 100644 index 00000000..9544f656 --- /dev/null +++ b/backend/app/schemas/temporal.py @@ -0,0 +1,71 @@ +from __future__ import annotations + +from datetime import datetime +from uuid import UUID + +from pydantic import BaseModel, Field + +from app.schemas.operations import VectorSelectionBBox + + +class TemporalComparisonRequest(BaseModel): + earlier_dataset_id: UUID + later_dataset_id: UUID + bbox: VectorSelectionBBox + preview_limit: int = Field(default=500, ge=1, le=1000) + + +class TemporalDatasetRef(BaseModel): + id: UUID + name: str + observed_at: datetime + source_version: str | None = None + + +class TemporalMetricComparison(BaseModel): + label: str + unit: str + aggregation_method: str + earlier_value: float + later_value: float + absolute_change: float + percent_change: float | None = None + is_estimate: bool = False + + +class TemporalObjectChanges(BaseModel): + available: bool + added_count: int | None = None + removed_count: int | None = None + modified_count: int | None = None + unchanged_count: int | None = None + + +class TemporalComparisonResponse(BaseModel): + temporal_series_key: str + earlier: TemporalDatasetRef + later: TemporalDatasetRef + selection_bbox: VectorSelectionBBox + metric: TemporalMetricComparison + object_changes: TemporalObjectChanges + geojson: dict + warnings: list[str] + generated_at: datetime + + +class TemporalSeriesDataset(BaseModel): + id: UUID + name: str + observed_at: datetime + source_version: str | None = None + feature_count: int | None = None + + +class TemporalSeriesRead(BaseModel): + temporal_series_key: str + source_name: str | None = None + reference_layer_name: str | None = None + dataset_count: int + first_observed_at: datetime + last_observed_at: datetime + datasets: list[TemporalSeriesDataset] diff --git a/backend/app/services/dataset_service.py b/backend/app/services/dataset_service.py index 77896bf8..35f3db52 100644 --- a/backend/app/services/dataset_service.py +++ b/backend/app/services/dataset_service.py @@ -12,8 +12,14 @@ from fastapi import UploadFile from sqlalchemy.orm import Session from app.core.errors import AppError -from app.models import Dataset, Project -from app.schemas.dataset import DatasetCreateResponse, DatasetStorageResponse, DatasetVectorSummary +from app.models import 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 @@ -26,6 +32,105 @@ class DatasetService: 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: @@ -89,44 +194,7 @@ class DatasetService: .limit(limit) .all() ) - response_items = [] - for row in rows: - feature_count = None - metadata_json = row.metadata_json or {} - vector_summary = DatasetService._extract_vector_summary(row.dataset_type, metadata_json) - if isinstance(metadata_json, dict): - feature_count = metadata_json.get("feature_count") - response_items.append( - DatasetCreateResponse( - id=row.id, - name=row.name, - dataset_type=row.dataset_type, - source=row.source, - dataset_role=row.dataset_role, - source_name=row.source_name, - reference_layer_name=row.reference_layer_name, - source_metadata=row.source_metadata, - provenance_metadata=row.provenance_metadata, - imported_at=row.imported_at, - project_id=row.project_id, - area_id=row.area_id, - storage_path=row.storage_path, - original_filename=row.original_filename, - stored_filename=row.stored_filename, - content_type=row.content_type, - size_bytes=row.size_bytes, - checksum_sha256=row.checksum_sha256, - crs=row.crs, - bounds_json=row.bounds_json, - metadata_json=row.metadata_json, - vector_summary=vector_summary, - status=row.status, - derived_from_dataset_id=row.derived_from_dataset_id, - created_at=row.created_at, - feature_count=feature_count, - ) - ) - return response_items, total + return [DatasetService._to_response(row) for row in rows], total @staticmethod def _extract_vector_summary(dataset_type: str, metadata_json: dict) -> DatasetVectorSummary | None: @@ -195,6 +263,12 @@ class DatasetService: 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) @@ -202,6 +276,14 @@ class DatasetService: 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" @@ -280,6 +362,7 @@ class DatasetService: 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"], @@ -294,6 +377,20 @@ class DatasetService: status=status, ) 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) @@ -306,34 +403,7 @@ class DatasetService: feature_class=feature_class, ) - 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, - 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, - derived_from_dataset_id=dataset.derived_from_dataset_id, - bounds_json=dataset.bounds_json, - metadata_json=dataset.metadata_json, - vector_summary=DatasetService._extract_vector_summary(dataset.dataset_type, dataset.metadata_json or {}), - status=dataset.status, - created_at=dataset.created_at, - feature_count=metadata.get("feature_count") if isinstance(metadata, dict) else None, - ) + return DatasetService._to_response(dataset) @staticmethod def refresh_metadata(db: Session, dataset_id: UUID) -> DatasetCreateResponse: @@ -380,34 +450,53 @@ class DatasetService: db.commit() db.refresh(dataset) - 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, - 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, - derived_from_dataset_id=dataset.derived_from_dataset_id, - bounds_json=dataset.bounds_json, - metadata_json=dataset.metadata_json, - vector_summary=DatasetService._extract_vector_summary(dataset.dataset_type, dataset.metadata_json or {}), - status=dataset.status, - created_at=dataset.created_at, - feature_count=metadata.get("feature_count") if isinstance(metadata, dict) else None, + 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: diff --git a/backend/app/services/demo_workflow_service.py b/backend/app/services/demo_workflow_service.py index 0af85da1..5820881b 100644 --- a/backend/app/services/demo_workflow_service.py +++ b/backend/app/services/demo_workflow_service.py @@ -10,7 +10,7 @@ from uuid import UUID, uuid4 from geoalchemy2.shape import from_shape from sqlalchemy.orm import Session -from app.models import Area, Dataset, Metric, Project, QualityCheck +from app.models import Area, Dataset, DatasetVersion, Metric, Project, QualityCheck from app.schemas.demo import DemoWorkflowResponse from app.services.geojson_service import parse_geojson_payload from app.services.qa_service import QaService @@ -29,6 +29,19 @@ class DemoWorkflowService: RASTER_FILENAME = "demo_context_raster.tif" EXPECTED_METRICS_FILENAME = "expected_qa_metrics.json" + @staticmethod + def _add_initial_version(db: Session, dataset: Dataset) -> None: + db.add( + DatasetVersion( + dataset_id=dataset.id, + version=1, + storage_path=dataset.storage_path, + checksum_sha256=dataset.checksum_sha256, + source_metadata=dataset.source_metadata, + provenance_metadata=dataset.provenance_metadata, + ) + ) + @staticmethod def _repo_root() -> Path: return Path(__file__).resolve().parents[3] @@ -213,6 +226,7 @@ class DemoWorkflowService: status="ready", ) db.add(dataset) + DemoWorkflowService._add_initial_version(db, dataset) db.commit() db.refresh(dataset) VectorFeatureService.persist_geojson_features( @@ -297,6 +311,7 @@ class DemoWorkflowService: status="ready", ) db.add(dataset) + DemoWorkflowService._add_initial_version(db, dataset) db.commit() db.refresh(dataset) return dataset diff --git a/backend/app/services/raster_operations_service.py b/backend/app/services/raster_operations_service.py index c580187f..b72b8b75 100644 --- a/backend/app/services/raster_operations_service.py +++ b/backend/app/services/raster_operations_service.py @@ -12,7 +12,7 @@ from shapely.ops import transform as shapely_transform from shapely.validation import make_valid from app.core.errors import AppError -from app.models import Area, Dataset +from app.models import Area, Dataset, DatasetVersion from app.services.raster_service import extract_raster_metadata from app.services.storage_service import StorageService @@ -333,6 +333,21 @@ class RasterOperationsService: name=output_name, dataset_type="raster", source=f"operation:{operation_name}", + dataset_role="derived", + source_name=source_dataset.source_name, + source_metadata=source_dataset.source_metadata, + provenance_metadata=provenance, + imported_at=datetime.now(timezone.utc), + temporal_series_key=( + f"{source_dataset.temporal_series_key}:{operation_name}" + if source_dataset.temporal_series_key + else None + ), + observed_at=source_dataset.observed_at, + valid_from=source_dataset.valid_from, + valid_to=source_dataset.valid_to, + temporal_granularity=source_dataset.temporal_granularity, + source_version=source_dataset.source_version, storage_path=str(output_file), original_filename=storage_metadata["original_filename"], stored_filename=storage_metadata["stored_filename"], @@ -348,6 +363,20 @@ class RasterOperationsService: status="ready", ) db.add(derived_dataset) + db.add( + DatasetVersion( + dataset_id=derived_dataset.id, + version=1, + storage_path=derived_dataset.storage_path, + source_version=derived_dataset.source_version, + observed_at=derived_dataset.observed_at, + valid_from=derived_dataset.valid_from, + valid_to=derived_dataset.valid_to, + checksum_sha256=derived_dataset.checksum_sha256, + source_metadata=derived_dataset.source_metadata, + provenance_metadata=derived_dataset.provenance_metadata, + ) + ) db.commit() db.refresh(derived_dataset) return derived_id diff --git a/backend/app/services/temporal_analysis_service.py b/backend/app/services/temporal_analysis_service.py new file mode 100644 index 00000000..e5a8e5c5 --- /dev/null +++ b/backend/app/services/temporal_analysis_service.py @@ -0,0 +1,310 @@ +from __future__ import annotations + +from datetime import datetime, timezone +from typing import Any +from uuid import UUID + +from geoalchemy2.functions import ST_Intersects, ST_MakeEnvelope +from geoalchemy2.shape import to_shape +from shapely.geometry import mapping +from sqlalchemy.orm import Session + +from app.core.errors import AppError +from app.models import Dataset, VectorFeature +from app.schemas.temporal import ( + TemporalComparisonRequest, + TemporalComparisonResponse, + TemporalDatasetRef, + TemporalMetricComparison, + TemporalObjectChanges, + TemporalSeriesDataset, + TemporalSeriesRead, +) +from app.services.vector_feature_service import VectorFeatureService + + +class TemporalAnalysisService: + IDENTITY_COMPARISON_LIMIT = 5_000 + + @staticmethod + def list_series(db: Session, project_id: UUID) -> list[TemporalSeriesRead]: + rows = ( + db.query(Dataset) + .filter(Dataset.project_id == project_id) + .filter(Dataset.temporal_series_key.isnot(None)) + .filter(Dataset.observed_at.isnot(None)) + .order_by(Dataset.temporal_series_key.asc(), Dataset.observed_at.asc()) + .all() + ) + grouped: dict[str, list[Dataset]] = {} + for row in rows: + if row.temporal_series_key: + grouped.setdefault(row.temporal_series_key, []).append(row) + + result: list[TemporalSeriesRead] = [] + for key, datasets in grouped.items(): + observed = [item.observed_at for item in datasets if item.observed_at is not None] + if not observed: + continue + result.append( + TemporalSeriesRead( + temporal_series_key=key, + source_name=datasets[-1].source_name, + reference_layer_name=datasets[-1].reference_layer_name, + dataset_count=len(datasets), + first_observed_at=min(observed), + last_observed_at=max(observed), + datasets=[ + TemporalSeriesDataset( + id=item.id, + name=item.name, + observed_at=item.observed_at, + source_version=item.source_version, + feature_count=(item.metadata_json or {}).get("feature_count") + if isinstance(item.metadata_json, dict) + else None, + ) + for item in datasets + if item.observed_at is not None + ], + ) + ) + return result + + @staticmethod + def compare( + db: Session, + *, + project_id: UUID, + payload: TemporalComparisonRequest, + ) -> TemporalComparisonResponse: + if payload.earlier_dataset_id == payload.later_dataset_id: + raise AppError( + code="INVALID_TEMPORAL_COMPARISON", + message="Choose two different dataset snapshots", + status_code=400, + ) + earlier = TemporalAnalysisService._get_temporal_dataset(db, project_id, payload.earlier_dataset_id, "Earlier") + later = TemporalAnalysisService._get_temporal_dataset(db, project_id, payload.later_dataset_id, "Later") + if earlier.temporal_series_key != later.temporal_series_key: + raise AppError( + code="INCOMPATIBLE_TEMPORAL_SERIES", + message="Dataset snapshots must belong to the same temporal series", + details={ + "earlier_series": earlier.temporal_series_key, + "later_series": later.temporal_series_key, + }, + status_code=400, + ) + if earlier.observed_at >= later.observed_at: + raise AppError( + code="INVALID_TEMPORAL_ORDER", + message="Earlier snapshot must have an observation date before the later snapshot", + status_code=400, + ) + + bbox = payload.bbox.model_dump() + earlier_summary = VectorFeatureService.summarize_features_by_bbox(db, dataset=earlier, bbox=bbox) + later_summary = VectorFeatureService.summarize_features_by_bbox(db, dataset=later, bbox=bbox) + if ( + earlier_summary["aggregation_method"] != later_summary["aggregation_method"] + or earlier_summary["metric_unit"] != later_summary["metric_unit"] + ): + raise AppError( + code="INCOMPATIBLE_TEMPORAL_AGGREGATION", + message="Dataset snapshots use incompatible aggregation semantics", + status_code=400, + ) + + earlier_value = float(earlier_summary["metric_value"]) + later_value = float(later_summary["metric_value"]) + absolute_change = later_value - earlier_value + percent_change = (absolute_change / earlier_value * 100.0) if earlier_value else None + warnings = [ + warning + for warning in {earlier_summary.get("warning"), later_summary.get("warning")} + if warning + ] + + object_changes, geojson, identity_warnings = TemporalAnalysisService._compare_identity_features( + db, + earlier=earlier, + later=later, + bbox=bbox, + preview_limit=payload.preview_limit, + ) + warnings.extend(identity_warnings) + + return TemporalComparisonResponse( + temporal_series_key=earlier.temporal_series_key, + earlier=TemporalDatasetRef( + id=earlier.id, + name=earlier.name, + observed_at=earlier.observed_at, + source_version=earlier.source_version, + ), + later=TemporalDatasetRef( + id=later.id, + name=later.name, + observed_at=later.observed_at, + source_version=later.source_version, + ), + selection_bbox=payload.bbox, + metric=TemporalMetricComparison( + label=str(later_summary["metric_label"]), + unit=str(later_summary["metric_unit"]), + aggregation_method=str(later_summary["aggregation_method"]), + earlier_value=earlier_value, + later_value=later_value, + absolute_change=absolute_change, + percent_change=percent_change, + is_estimate=bool(earlier_summary["is_estimate"] or later_summary["is_estimate"]), + ), + object_changes=object_changes, + geojson=geojson, + warnings=warnings, + generated_at=datetime.now(timezone.utc), + ) + + @staticmethod + def _get_temporal_dataset(db: Session, project_id: UUID, dataset_id: UUID, label: str) -> Dataset: + dataset = db.get(Dataset, dataset_id) + if not dataset or dataset.project_id != project_id: + raise AppError(code="DATASET_NOT_FOUND", message=f"{label} dataset not found", status_code=404) + if dataset.dataset_type not in {"vector", "geojson"}: + raise AppError( + code="DATASET_NOT_VECTOR", + message="Temporal selection comparison currently requires vector datasets", + status_code=400, + ) + if not dataset.temporal_series_key or not dataset.observed_at: + raise AppError( + code="TEMPORAL_METADATA_MISSING", + message=f"{label} dataset has no explicit temporal series and observation date", + status_code=400, + ) + return dataset + + @staticmethod + def _compare_identity_features( + db: Session, + *, + earlier: Dataset, + later: Dataset, + bbox: dict[str, Any], + preview_limit: int, + ) -> tuple[TemporalObjectChanges, dict[str, Any], list[str]]: + earlier_config = earlier.source_metadata if isinstance(earlier.source_metadata, dict) else {} + later_config = later.source_metadata if isinstance(later.source_metadata, dict) else {} + if not earlier_config.get("identity_stable") or not later_config.get("identity_stable"): + return ( + TemporalObjectChanges(available=False), + {"type": "FeatureCollection", "features": []}, + ["Object-level changes are unavailable because the source does not guarantee stable feature identifiers."], + ) + + normalized_bbox = VectorFeatureService._normalize_selection_bbox(bbox) + envelope = ST_MakeEnvelope( + normalized_bbox["min_x"], + normalized_bbox["min_y"], + normalized_bbox["max_x"], + normalized_bbox["max_y"], + 4326, + ) + + def load(dataset_id: UUID) -> list[VectorFeature]: + return ( + db.query(VectorFeature) + .filter(VectorFeature.dataset_id == dataset_id) + .filter(ST_Intersects(VectorFeature.geometry, envelope)) + .filter(VectorFeature.source_feature_id.isnot(None)) + .order_by(VectorFeature.source_feature_id.asc()) + .limit(TemporalAnalysisService.IDENTITY_COMPARISON_LIMIT + 1) + .all() + ) + + earlier_rows = load(earlier.id) + later_rows = load(later.id) + if ( + len(earlier_rows) > TemporalAnalysisService.IDENTITY_COMPARISON_LIMIT + or len(later_rows) > TemporalAnalysisService.IDENTITY_COMPARISON_LIMIT + ): + return ( + TemporalObjectChanges(available=False), + {"type": "FeatureCollection", "features": []}, + ["Object-level preview was skipped because the selection exceeds the 5,000 feature safety limit."], + ) + + earlier_by_id = {str(row.source_feature_id): row for row in earlier_rows if row.source_feature_id} + later_by_id = {str(row.source_feature_id): row for row in later_rows if row.source_feature_id} + earlier_ids = set(earlier_by_id) + later_ids = set(later_by_id) + added_ids = sorted(later_ids - earlier_ids) + removed_ids = sorted(earlier_ids - later_ids) + common_ids = sorted(earlier_ids & later_ids) + comparison_property = str(later_config.get("comparison_property") or "").strip() or None + modified_ids: list[str] = [] + unchanged_ids: list[str] = [] + + for feature_id in common_ids: + earlier_row = earlier_by_id[feature_id] + later_row = later_by_id[feature_id] + geometry_changed = not to_shape(earlier_row.geometry).equals(to_shape(later_row.geometry)) + value_changed = False + if comparison_property: + value_changed = (earlier_row.properties_json or {}).get(comparison_property) != ( + later_row.properties_json or {} + ).get(comparison_property) + (modified_ids if geometry_changed or value_changed else unchanged_ids).append(feature_id) + + features: list[dict[str, Any]] = [] + for change_type, feature_ids, rows in ( + ("added", added_ids, later_by_id), + ("removed", removed_ids, earlier_by_id), + ("modified", modified_ids, later_by_id), + ): + for feature_id in feature_ids: + if len(features) >= preview_limit: + break + row = rows[feature_id] + properties = dict(row.properties_json or {}) + properties.update( + { + "change_type": change_type, + "source_feature_id": feature_id, + "earlier_dataset_id": str(earlier.id), + "later_dataset_id": str(later.id), + } + ) + if change_type == "modified" and comparison_property: + before = (earlier_by_id[feature_id].properties_json or {}).get(comparison_property) + after = (later_by_id[feature_id].properties_json or {}).get(comparison_property) + properties.update({"value_before": before, "value_after": after}) + if isinstance(before, (int, float)) and isinstance(after, (int, float)): + properties["value_delta"] = after - before + features.append( + { + "type": "Feature", + "id": str(row.id), + "geometry": mapping(to_shape(row.geometry)), + "properties": properties, + } + ) + + warnings: list[str] = [] + total_changes = len(added_ids) + len(removed_ids) + len(modified_ids) + if total_changes > preview_limit: + warnings.append( + f"The map shows the first {preview_limit} of {total_changes} changed features; counts remain complete." + ) + return ( + TemporalObjectChanges( + available=True, + added_count=len(added_ids), + removed_count=len(removed_ids), + modified_count=len(modified_ids), + unchanged_count=len(unchanged_ids), + ), + {"type": "FeatureCollection", "features": features}, + warnings, + ) diff --git a/backend/app/services/vector_feature_service.py b/backend/app/services/vector_feature_service.py index 02ce54c5..21c19f9a 100644 --- a/backend/app/services/vector_feature_service.py +++ b/backend/app/services/vector_feature_service.py @@ -9,9 +9,10 @@ from geoalchemy2.shape import to_shape from shapely.geometry import mapping from shapely.geometry import shape from shapely.validation import make_valid +from sqlalchemy import Float, cast, func from app.core.errors import AppError -from app.models import VectorFeature +from app.models import Dataset, VectorFeature class VectorFeatureService: @@ -88,6 +89,7 @@ class VectorFeatureService: dataset_id: UUID, bbox: dict[str, Any], limit: int = 100, + dataset: Dataset | None = None, ) -> dict[str, Any]: normalized_bbox = VectorFeatureService._normalize_selection_bbox(bbox) safe_limit = max(1, min(int(limit), 1000)) @@ -121,6 +123,14 @@ class VectorFeatureService: truncated = total_feature_count > safe_limit selected_rows = rows[:safe_limit] features = [VectorFeatureService._row_to_geojson_feature(row) for row in selected_rows] + summary = None + if dataset and isinstance(dataset.source_metadata, dict) and dataset.source_metadata.get("selection_aggregation"): + summary = VectorFeatureService.summarize_features_by_bbox( + db, + dataset=dataset, + bbox=normalized_bbox, + total_feature_count=total_feature_count, + ) return { "selection_bbox": normalized_bbox, @@ -132,6 +142,97 @@ class VectorFeatureService: "type": "FeatureCollection", "features": features, }, + "summary": summary, + } + + @staticmethod + def summarize_features_by_bbox( + db, + *, + dataset: Dataset, + bbox: dict[str, Any], + total_feature_count: int | None = None, + ) -> dict[str, Any]: + normalized_bbox = VectorFeatureService._normalize_selection_bbox(bbox) + envelope = ST_MakeEnvelope( + normalized_bbox["min_x"], + normalized_bbox["min_y"], + normalized_bbox["max_x"], + normalized_bbox["max_y"], + 4326, + ) + selection_filter = ( + VectorFeature.dataset_id == dataset.id, + ST_Intersects(VectorFeature.geometry, envelope), + ) + feature_count = total_feature_count + if feature_count is None: + feature_count = int(db.query(func.count(VectorFeature.id)).filter(*selection_filter).scalar() or 0) + + source_metadata = dataset.source_metadata if isinstance(dataset.source_metadata, dict) else {} + config = source_metadata.get("selection_aggregation") + if not isinstance(config, dict): + config = {} + method = str(config.get("method") or "feature_count") + label = str(config.get("label") or "Objecten") + unit = str(config.get("unit") or "objecten") + warning = str(config["warning"]) if config.get("warning") else None + is_estimate = bool(config.get("is_estimate", False)) + + metric_value = float(feature_count) + if method == "intersection_area": + intersection = func.ST_Intersection(VectorFeature.geometry, envelope) + area_expression = func.ST_Area(func.ST_Transform(intersection, 31370)) + area_m2 = db.query(func.coalesce(func.sum(area_expression), 0.0)).filter(*selection_filter).scalar() + divisor = 10_000.0 if unit == "ha" else 1.0 + metric_value = float(area_m2 or 0.0) / divisor + elif method == "intersection_length": + intersection = func.ST_Intersection(VectorFeature.geometry, envelope) + length_expression = func.ST_Length(func.ST_Transform(intersection, 31370)) + length_m = db.query(func.coalesce(func.sum(length_expression), 0.0)).filter(*selection_filter).scalar() + divisor = 1_000.0 if unit == "km" else 1.0 + metric_value = float(length_m or 0.0) / divisor + elif method in {"sum", "area_weighted_sum"}: + property_name = str(config.get("property") or "").strip() + if not property_name: + raise AppError( + code="INVALID_SELECTION_AGGREGATION", + message="Dataset selection aggregation requires a numeric property", + details={"dataset_id": str(dataset.id), "method": method}, + status_code=500, + ) + numeric_value = cast(VectorFeature.properties_json.op("->>")(property_name), Float) + value_expression = numeric_value + if method == "area_weighted_sum": + source_area = func.ST_Area(func.ST_Transform(VectorFeature.geometry, 31370)) + intersection_area = func.ST_Area( + func.ST_Transform(func.ST_Intersection(VectorFeature.geometry, envelope), 31370) + ) + value_expression = numeric_value * intersection_area / func.nullif(source_area, 0.0) + is_estimate = True + aggregate_value = ( + db.query(func.coalesce(func.sum(value_expression), 0.0)) + .filter(*selection_filter) + .filter(VectorFeature.properties_json.op("->>")(property_name).isnot(None)) + .scalar() + ) + metric_value = float(aggregate_value or 0.0) + elif method != "feature_count": + raise AppError( + code="INVALID_SELECTION_AGGREGATION", + message="Unsupported dataset selection aggregation", + details={"dataset_id": str(dataset.id), "method": method}, + status_code=500, + ) + + return { + "metric_label": label, + "metric_value": metric_value, + "metric_unit": unit, + "aggregation_method": method, + "feature_count": feature_count, + "is_estimate": is_estimate, + "warning": warning, } @staticmethod diff --git a/backend/app/services/vector_operations_service.py b/backend/app/services/vector_operations_service.py index 1aff6222..630e3591 100644 --- a/backend/app/services/vector_operations_service.py +++ b/backend/app/services/vector_operations_service.py @@ -15,7 +15,7 @@ from shapely.validation import make_valid from sqlalchemy.orm import Session from app.core.errors import AppError -from app.models import Area, Dataset +from app.models import Area, Dataset, DatasetVersion from app.schemas.dataset import DatasetCreateResponse from app.schemas.operations import VectorOperationResult from app.services.geojson_service import parse_geojson_payload @@ -444,6 +444,16 @@ class VectorOperationsService: source_metadata=source_metadata, provenance_metadata=provenance_metadata, imported_at=datetime.now(timezone.utc), + temporal_series_key=( + f"{source_dataset.temporal_series_key}:{operation}" + if source_dataset.temporal_series_key + else None + ), + observed_at=source_dataset.observed_at, + valid_from=source_dataset.valid_from, + valid_to=source_dataset.valid_to, + temporal_granularity=source_dataset.temporal_granularity, + source_version=source_dataset.source_version, storage_path=storage_info["storage_path"], original_filename=storage_info["original_filename"], stored_filename=storage_info["stored_filename"], @@ -459,6 +469,20 @@ class VectorOperationsService: status="ready", ) db.add(derived_dataset) + db.add( + DatasetVersion( + dataset_id=derived_dataset.id, + version=1, + storage_path=derived_dataset.storage_path, + source_version=derived_dataset.source_version, + observed_at=derived_dataset.observed_at, + valid_from=derived_dataset.valid_from, + valid_to=derived_dataset.valid_to, + checksum_sha256=derived_dataset.checksum_sha256, + source_metadata=derived_dataset.source_metadata, + provenance_metadata=derived_dataset.provenance_metadata, + ) + ) db.commit() db.refresh(derived_dataset) if persist_vector_features: diff --git a/backend/tests/test_raster_operations_service.py b/backend/tests/test_raster_operations_service.py index 73b8eabb..5740a316 100644 --- a/backend/tests/test_raster_operations_service.py +++ b/backend/tests/test_raster_operations_service.py @@ -7,7 +7,7 @@ import importlib from geoalchemy2.shape import from_shape from app.core.errors import AppError -from app.models import Area, Dataset +from app.models import Area, Dataset, DatasetVersion from app.services.raster_operations_service import RasterOperationsService from app.api.routes.datasets import _run_job_sync from shapely.geometry import box @@ -356,8 +356,12 @@ def test_raster_reproject_returns_persisted_derived_dataset(monkeypatch, tmp_pat ) assert result_id == output_id - assert len(db.added) == 1 + assert len(db.added) == 2 derived = db.added[0] + version = db.added[1] + assert isinstance(version, DatasetVersion) + assert version.dataset_id == output_id + assert version.version == 1 assert derived.id == output_id assert derived.metadata_json is not None assert derived.metadata_json["operation"] == "raster.reproject" @@ -783,9 +787,13 @@ def test_raster_clip_persists_derived_dataset(monkeypatch, tmp_path) -> None: result_id = RasterOperationsService.clip(db, dataset_id, area_id, "clip-result.tif") assert result_id == output_id - assert len(db.added) == 1 + assert len(db.added) == 2 derived = db.added[0] + version = db.added[1] assert isinstance(derived, Dataset) + assert isinstance(version, DatasetVersion) + assert version.dataset_id == output_id + assert version.version == 1 assert derived.id == output_id assert derived.source == "operation:raster.clip" assert derived.dataset_type == "raster" @@ -1295,8 +1303,12 @@ def test_raster_index_records_provenance_and_dtype(tmp_path, monkeypatch) -> Non result_dataset_id = RasterOperationsService.ndvi(db, dataset_id, nir_band=4, red_band=3, output_name="ndvi-test") assert result_dataset_id == output_dataset_id - assert len(db.added) == 1 + assert len(db.added) == 2 derived = db.added[0] + version = db.added[1] + assert isinstance(version, DatasetVersion) + assert version.dataset_id == output_dataset_id + assert version.version == 1 assert derived.id == output_dataset_id assert derived.metadata_json is not None assert derived.metadata_json["operation"] == "raster.ndvi" diff --git a/backend/tests/test_sprint187_temporal_map_foundation.py b/backend/tests/test_sprint187_temporal_map_foundation.py new file mode 100644 index 00000000..fe40af64 --- /dev/null +++ b/backend/tests/test_sprint187_temporal_map_foundation.py @@ -0,0 +1,250 @@ +from __future__ import annotations + +from datetime import datetime, timezone +from pathlib import Path +from types import SimpleNamespace +from uuid import uuid4 + +import pytest + +from app.core.errors import AppError +from app.models import Dataset, DatasetVersion +from app.schemas.dataset import DatasetTemporalUpdate +from app.schemas.temporal import TemporalComparisonRequest, TemporalObjectChanges +from app.services.dataset_service import DatasetService +from app.services.temporal_analysis_service import TemporalAnalysisService +from app.services.vector_feature_service import VectorFeatureService + + +ROOT = Path(__file__).parents[2] + + +class ScalarQuery: + def __init__(self, value: float): + self.value = value + + def filter(self, *args): # noqa: ANN002, ARG002 + return self + + def scalar(self): + return self.value + + +class ScalarSession: + def __init__(self, value: float): + self.value = value + + def query(self, *args): # noqa: ANN002, ARG002 + return ScalarQuery(self.value) + + +class VersionQuery: + def __init__(self, latest: DatasetVersion | None): + self.latest = latest + + def filter(self, *args): # noqa: ANN002, ARG002 + return self + + def order_by(self, *args): # noqa: ANN002, ARG002 + return self + + def first(self): + return self.latest + + +class TemporalUpdateSession: + def __init__(self, dataset: Dataset, latest: DatasetVersion | None): + self.dataset = dataset + self.latest = latest + self.added: list[object] = [] + + def get(self, model, item_id): # noqa: ANN001 + return self.dataset if model is Dataset and item_id == self.dataset.id else None + + def query(self, model): # noqa: ANN001 + assert model is DatasetVersion + return VersionQuery(self.latest) + + def add(self, item): # noqa: ANN001 + self.added.append(item) + + def commit(self): + return None + + def refresh(self, _item): + return None + + +def temporal_dataset(*, project_id, observed_year: int, metric_method: str = "feature_count") -> Dataset: + return Dataset( + id=uuid4(), + project_id=project_id, + name=f"snapshot-{observed_year}.geojson", + dataset_type="vector", + source="official", + dataset_role="reference", + temporal_series_key="official:test:mol", + observed_at=datetime(observed_year, 1, 1, tzinfo=timezone.utc), + source_version=str(observed_year), + source_metadata={ + "selection_aggregation": { + "method": metric_method, + "label": "Objecten", + "unit": "objecten", + } + }, + ) + + +def test_temporal_migration_and_models_align() -> None: + migration = (ROOT / "backend/alembic/versions/202607140001_temporal_dataset_foundation.py").read_text(encoding="utf-8") + for field in ( + "temporal_series_key", + "observed_at", + "valid_from", + "valid_to", + "temporal_granularity", + "source_version", + ): + assert field in migration + assert hasattr(Dataset, field) + assert "ix_vector_features_dataset_source_feature" in migration + assert 'down_revision = "202606120900"' in migration + + +def test_temporal_metadata_requires_an_explicit_series_and_observation_date() -> None: + with pytest.raises(AppError, match="observed_at is required"): + DatasetService._validate_temporal_metadata( + temporal_series_key="official:test:mol", + observed_at=None, + valid_from=None, + valid_to=None, + temporal_granularity="year", + source_version="2024", + ) + with pytest.raises(AppError, match="valid_to must be"): + DatasetService._validate_temporal_metadata( + temporal_series_key="official:test:mol", + observed_at=datetime(2024, 1, 1, tzinfo=timezone.utc), + valid_from=datetime(2024, 12, 31, tzinfo=timezone.utc), + valid_to=datetime(2024, 1, 1, tzinfo=timezone.utc), + temporal_granularity="year", + source_version="2024", + ) + + +def test_temporal_metadata_update_appends_provenance_version_and_is_idempotent() -> None: + project_id = uuid4() + dataset = temporal_dataset(project_id=project_id, observed_year=2024) + dataset.status = "ready" + dataset.metadata_json = {} + latest = DatasetVersion( + dataset_id=dataset.id, + version=3, + observed_at=dataset.observed_at, + source_version="2024", + ) + session = TemporalUpdateSession(dataset, latest) + payload = DatasetTemporalUpdate( + temporal_series_key="official:test:mol", + observed_at=datetime(2025, 1, 1, tzinfo=timezone.utc), + temporal_granularity="year", + source_version="2025", + ) + + updated = DatasetService.update_temporal_metadata(session, dataset.id, payload) + + assert updated.observed_at == payload.observed_at + assert latest.version == 3 + assert latest.observed_at == datetime(2024, 1, 1, tzinfo=timezone.utc) + assert len(session.added) == 2 + appended = session.added[1] + assert isinstance(appended, DatasetVersion) + assert appended.version == 4 + assert appended.observed_at == payload.observed_at + + session.added.clear() + DatasetService.update_temporal_metadata(session, dataset.id, payload) + assert session.added == [] + + +def test_selection_area_aggregation_returns_hectares_without_loading_all_features() -> None: + project_id = uuid4() + dataset = temporal_dataset(project_id=project_id, observed_year=1969, metric_method="intersection_area") + dataset.source_metadata["selection_aggregation"].update({"label": "Oppervlakte", "unit": "ha"}) + result = VectorFeatureService.summarize_features_by_bbox( + ScalarSession(125_000.0), + dataset=dataset, + bbox={"min_x": 5.0, "min_y": 51.1, "max_x": 5.2, "max_y": 51.3, "crs": "EPSG:4326"}, + total_feature_count=40, + ) + assert result["metric_value"] == 12.5 + assert result["metric_unit"] == "ha" + assert result["feature_count"] == 40 + + +def test_temporal_compare_returns_delta_and_canonical_change_payload(monkeypatch) -> None: + project_id = uuid4() + earlier = temporal_dataset(project_id=project_id, observed_year=2021) + later = temporal_dataset(project_id=project_id, observed_year=2024) + + def get_dataset(_db, _project_id, dataset_id, _label): + return earlier if dataset_id == earlier.id else later + + def summarize(_db, *, dataset, bbox): # noqa: ARG001 + value = 100.0 if dataset.id == earlier.id else 115.0 + return { + "metric_label": "Inwoners", + "metric_value": value, + "metric_unit": "inwoners", + "aggregation_method": "area_weighted_sum", + "feature_count": 10, + "is_estimate": True, + "warning": "Areal weighting", + } + + monkeypatch.setattr(TemporalAnalysisService, "_get_temporal_dataset", staticmethod(get_dataset)) + monkeypatch.setattr(VectorFeatureService, "summarize_features_by_bbox", staticmethod(summarize)) + monkeypatch.setattr( + TemporalAnalysisService, + "_compare_identity_features", + staticmethod( + lambda *args, **kwargs: ( + TemporalObjectChanges(available=True, added_count=1, removed_count=0, modified_count=2, unchanged_count=7), + {"type": "FeatureCollection", "features": []}, + [], + ) + ), + ) + result = TemporalAnalysisService.compare( + SimpleNamespace(), + project_id=project_id, + payload=TemporalComparisonRequest( + earlier_dataset_id=earlier.id, + later_dataset_id=later.id, + bbox={"min_x": 5.0, "min_y": 51.1, "max_x": 5.2, "max_y": 51.3}, + ), + ) + assert result.metric.absolute_change == 15.0 + assert result.metric.percent_change == 15.0 + assert result.metric.is_estimate is True + assert result.object_changes.modified_count == 2 + assert result.geojson["type"] == "FeatureCollection" + + +def test_temporal_frontend_and_official_operator_contracts_exist() -> None: + workspace = (ROOT / "frontend/src/components/map/MapWorkspace.tsx").read_text(encoding="utf-8") + temporal_api = (ROOT / "frontend/src/services/api/temporal.ts").read_text(encoding="utf-8") + population = (ROOT / "scripts/provision_mol_population_history.py").read_text(encoding="utf-8") + landuse = (ROOT / "scripts/provision_mol_historical_landuse.py").read_text(encoding="utf-8") + dockerfile = (ROOT / "deploy/unraid/Dockerfile.all-in-one").read_text(encoding="utf-8") + + assert "Laatste toestand" in workspace + assert "Evolutie" in workspace + assert "Vergelijk periode" in workspace + assert "/temporal/compare" in temporal_api + assert "Statbel" in population and "area_weighted_sum" in population + assert "HistLandgebruik" in landuse and "intersection_area" in landuse + assert "provision_mol_population_history.py" in dockerfile + assert "provision_mol_historical_landuse.py" in dockerfile + assert "fake" not in population.lower() diff --git a/deploy/unraid/Dockerfile.all-in-one b/deploy/unraid/Dockerfile.all-in-one index cc8cb111..623f9e82 100644 --- a/deploy/unraid/Dockerfile.all-in-one +++ b/deploy/unraid/Dockerfile.all-in-one @@ -74,6 +74,8 @@ RUN python scripts/gis_import_smoke.py \ COPY scripts/prepare_operator_real_data_samples.py /app/scripts/prepare_operator_real_data_samples.py COPY scripts/provision_mol_municipality_workspace.py /app/scripts/provision_mol_municipality_workspace.py COPY scripts/provision_mol_context_layers.py /app/scripts/provision_mol_context_layers.py +COPY scripts/provision_mol_population_history.py /app/scripts/provision_mol_population_history.py +COPY scripts/provision_mol_historical_landuse.py /app/scripts/provision_mol_historical_landuse.py COPY scripts/export_operator_yolo_tile_dataset.py /app/scripts/export_operator_yolo_tile_dataset.py COPY scripts/audit_operator_yolo_dataset_quality.py /app/scripts/audit_operator_yolo_dataset_quality.py COPY scripts/render_operator_yolo_label_qa_contact_sheets.py /app/scripts/render_operator_yolo_label_qa_contact_sheets.py diff --git a/docs/API_CONTRACTS.md b/docs/API_CONTRACTS.md index 3d72216d..b24b1959 100644 --- a/docs/API_CONTRACTS.md +++ b/docs/API_CONTRACTS.md @@ -1393,3 +1393,53 @@ GET /api/v1/exports/{export_id}/download ``` PDF/report-designer functionality can be added after core GeoAI workflows work. + +## Temporal datasets and area evolution + +Temporal metadata describes the source observation, not API run history. A +snapshot is a normal persisted dataset grouped by `temporal_series_key` and +ordered by `observed_at`. Optional validity uses `valid_from` and `valid_to`; +`temporal_granularity` is `snapshot`, `day`, `month`, `year` or `period`. + +Dataset upload accepts those temporal fields plus `source_version`. When a +dataset declares `source_metadata.selection_aggregation`, the vector bbox +selection response also contains a `summary` with metric label/value/unit, +aggregation method, feature count, estimate status and an optional warning. +Supported PostGIS aggregations are feature count, intersection area, +intersection length, numeric sum and area-weighted numeric sum. Area and length +are measured after transformation to EPSG:31370. + +### PATCH `/api/v1/projects/{project_id}/datasets/{dataset_id}/temporal` + +Updates the temporal provenance of an existing dataset. Series key and +observation date are required together. It does not alter features or +manufacture a historical observation. + +### GET `/api/v1/projects/{project_id}/datasets/{dataset_id}/versions` + +Lists immutable storage/provenance versions. Uploads and derived datasets +create version 1 in the same persistence transaction. + +### GET `/api/v1/projects/{project_id}/temporal/series` + +Returns dated project series in the canonical envelope. Each item contains its +source/layer identity, first and last observations and ordered datasets. + +### POST `/api/v1/projects/{project_id}/temporal/compare` + +```json +{ + "earlier_dataset_id": "uuid", + "later_dataset_id": "uuid", + "bbox": {"west": 5.0, "south": 51.0, "east": 5.2, "north": 51.2}, + "preview_limit": 500 +} +``` + +Both datasets must belong to the project and the same temporal series, with +the earlier observation preceding the later one. The response contains source +snapshot references, selection bbox, earlier/later metric values, +absolute/percentage change, estimate status, warnings and GeoJSON evidence. +Added/removed/modified object changes are calculated only when source +provenance declares stable feature identities; otherwise +`object_changes.available=false` and no object history is inferred. diff --git a/docs/CODEX_EXECUTION_LOG.md b/docs/CODEX_EXECUTION_LOG.md index 8821c811..5b4a0621 100644 --- a/docs/CODEX_EXECUTION_LOG.md +++ b/docs/CODEX_EXECUTION_LOG.md @@ -7768,3 +7768,29 @@ Limitations: Next: - Define authoritative population and land-cover source adapters, then reuse the proven municipality provisioner and bbox analysis flow for the complete Kempen. + +## Sprint 187 Temporal Mol explorer (2026-07-14) + +Implemented: +- Added observation/validity/source-version fields to datasets and immutable provenance fields to dataset versions, with one Alembic head and indexed temporal/source identity lookups. +- Persisted dataset version 1 atomically for uploads, demo fixtures and derived vector/raster operations. +- Added source-governed PostGIS selection summaries for object count, area, length, numeric sum and area-weighted sum. +- Added project temporal-series discovery and same-series bbox comparison with honest metric deltas and stable-identity-only object changes. +- Added the map-first `Laatste toestand` / `Evolutie` workflow with period selection, automatic rectangle analysis and change overlays. +- Added explicit, idempotent Statbel population (2021-2025) and Digitaal Vlaanderen historical land-use (1778/1873/1969) provisioners for Mol. +- Kept all source fetching operator-triggered; application startup and user queries never fabricate or silently download source data. + +Methodology: +- Population totals are source-published per statistical sector. Intersections with partial sectors are labelled area-weighted estimates. +- Historical land-use classes are clipped from official editions and measured in EPSG:31370. They do not claim stable cadastral object identity. +- Every snapshot carries its source URL, observation date, source version, checksum and processing limitations. + +Validation before live deployment: +- Script compilation passed. +- New temporal/API/static regression suite passed 6 tests, including append-only/idempotent temporal provenance updates. +- Raster and temporal focused suite passed 27 tests after extending existing assertions to require dataset-version persistence. +- Frontend TypeScript typecheck and production build passed. +- Offline Alembic SQL generation passed with head `202607140001`. + +Next: +- Run the complete readiness gate, deploy to Tower/PostGIS, provision the official snapshots and verify current/evolution selection end to end in the internal browser. diff --git a/docs/DATABASE_IMPLEMENTATION_PLAN.md b/docs/DATABASE_IMPLEMENTATION_PLAN.md index 2ffaeade..801f5d95 100644 --- a/docs/DATABASE_IMPLEMENTATION_PLAN.md +++ b/docs/DATABASE_IMPLEMENTATION_PLAN.md @@ -236,3 +236,22 @@ GRB and OSM live imports are intentionally `not_configured` in Sprint 7B. Manual - Multi-tenant row-level security. - User accounts. - Full model registry tables. + +## Temporal dataset foundation + +Historical observations remain normal `datasets` and `vector_features`; there +is no parallel temporal feature store. Snapshots are grouped by +`datasets.temporal_series_key` and carry `observed_at`, `valid_from`, +`valid_to`, `temporal_granularity` and `source_version`. Observation time is +kept separate from ingestion time (`imported_at`). + +`dataset_versions` records immutable storage provenance for every upload and +derived output: dataset-local `version`, storage path, source version, +observation/validity dates, checksum and source/provenance JSON. The +`(dataset_id, version)` pair is unique. Temporal series lookup is indexed by +`(project_id, temporal_series_key, observed_at)` and source-feature lookup by +`(dataset_id, source_feature_id)`. + +Time-series comparison is read-only and aggregates persisted geometry inside a +requested bbox. Object-level added/removed/modified evidence is only valid for +sources that explicitly declare stable source feature identifiers. diff --git a/docs/DATASET_STRATEGY.md b/docs/DATASET_STRATEGY.md index 9b7d2386..d5d59466 100644 --- a/docs/DATASET_STRATEGY.md +++ b/docs/DATASET_STRATEGY.md @@ -138,3 +138,20 @@ V1 dataset strategy is complete when: - an area can request/cache a reference building layer; - detection outputs can be compared with that reference layer; - exports include source metadata. + +## Temporal snapshots and evolution + +- A historical observation is one persisted dataset snapshot. Existing + datasets are not overwritten and features are not hidden in job JSON. +- Related observations share a stable `temporal_series_key`; `observed_at` + records when the source describes reality, while `imported_at` records + ingestion time. +- The latest-state map uses the latest available observation but must not call + an old source edition current reality. +- Evolution metrics use the same selection geometry, aggregation and units for + both snapshots. +- Partial-sector population is an area-weighted estimate. GeoIntel must not + imply address-level distribution when only sector totals are available. +- Historical cartographic classes can change meaning between editions. Source + classes and processing notes remain provenance, and object changes require + explicit stable source identity. diff --git a/docs/DATA_SOURCES.md b/docs/DATA_SOURCES.md index 9e719a82..d11344f8 100644 --- a/docs/DATA_SOURCES.md +++ b/docs/DATA_SOURCES.md @@ -29,6 +29,32 @@ start geen verborgen providerfetch. De resulterende GeoJSON-artefacten, checksums en bron-URL's worden onder persistent operator storage bewaard en via de bestaande DatasetService/vectorfeature-flow geïmporteerd. +### Mol population history + +`scripts/provision_mol_population_history.py` imports official Statbel +population-by-statistical-sector tables and matching sector geometries for +2021 through 2025. It keeps the published sector total as +`population_total`, clips the official geometry to Mol NIS `13025` and uploads +each year as a separate dataset in +`statbel:population-statistical-sector:mol`. + +Complete sectors use their published population total. A rectangle that cuts +through a sector uses an explicitly labelled area-weighted estimate; the +source does not justify a more precise intra-sector distribution. + +### Mol historical land use + +`scripts/provision_mol_historical_landuse.py` uses the official Digitaal +Vlaanderen Historical Land Use WFS for the 1778, 1873 and 1969 collections. +Buildings, forest, water and roads are filtered server-side, clipped to the +official Mol boundary and uploaded through the existing dataset service. +Source classes, request URLs, observation year, simplification tolerance and +methodological limitations remain in provenance. + +These historical map editions support exploratory area evolution, not +cadastral object lineage. Their feature identities are declared unstable and +GeoIntel does not fabricate added/removed object counts. + ## OSM - Naam: OpenStreetMap diff --git a/frontend/README.md b/frontend/README.md index 50133f78..ed87b39b 100644 --- a/frontend/README.md +++ b/frontend/README.md @@ -4,6 +4,19 @@ React + TypeScript + MapLibre foundation for project/area/dataset workflow. Mol is the primary operating context. On a fresh session the application opens the map-first geographic explorer, prefers the persisted `Mol Municipality Workbench`, selects the official NIS `13025` municipality boundary and activates the largest available authoritative building layer. Explicit project and dataset selections remain authoritative, and all broader Kempen workflows remain available. +The map-first explorer has two deliberate modes. `Latest state` selects the +latest explicitly dated source snapshot without claiming an old edition is +current, while `Evolution` lets the operator compare an earlier and later +snapshot from the same series over a drawn rectangle. Results show units, +absolute/percentage change, estimate status and source limitations. +Added/removed/modified overlays only appear for stable source identities. + +Selection results use dataset-specific PostGIS summaries. Object layers show +intersecting counts, population shows inhabitants with partial-sector +estimates clearly marked and land-cover sources show intersected hectares. The +advanced workbench remains available but is not required for the primary +choose-theme, draw-area, read-result flow. + The primary workflow is deliberately short: choose a data theme, drag a rectangle on the MapLibre map and read the resulting PostGIS evidence. Releasing the drag runs the active theme query and every other available theme query for the same EPSG:4326 bbox. The result panel shows selection area, exact intersection totals, active-theme density, source identity and bounded feature properties. Map rendering remains capped at 1,000 features while `total_feature_count` reports the exact database count. The theme catalog currently recognizes buildings, population, forest/green, water, roads and parcels from dataset names and canonical `reference_layer_name` metadata. A theme is enabled only when a ready persisted vector dataset exists; otherwise it states `Bron nog niet ingeladen`. This prevents missing population or land-cover sources from appearing as zero-valued observations. The previous technical Map workspace remains available through `Geavanceerde werkbank` for derived datasets, QA/QC evidence and export operations. diff --git a/frontend/src/components/GeoMap.tsx b/frontend/src/components/GeoMap.tsx index 3d21459f..621218a8 100644 --- a/frontend/src/components/GeoMap.tsx +++ b/frontend/src/components/GeoMap.tsx @@ -319,6 +319,8 @@ function GeoMap({ '#16a34a', 'removed', '#dc2626', + 'modified', + '#d97706', 'unchanged', '#2563eb', '#f97316', @@ -345,6 +347,8 @@ function GeoMap({ '#15803d', 'removed', '#b91c1c', + 'modified', + '#b45309', 'unchanged', '#1d4ed8', '#ea580c', diff --git a/frontend/src/components/map/MapWorkspace.tsx b/frontend/src/components/map/MapWorkspace.tsx index 25b46d11..03a34c93 100644 --- a/frontend/src/components/map/MapWorkspace.tsx +++ b/frontend/src/components/map/MapWorkspace.tsx @@ -3,6 +3,7 @@ import GeoMap from '../GeoMap' import type { AreaRead, DatasetCreateResponse, MapViewportState, QaComparisonResult, VectorSelectionBBox, VectorSelectionResponse } from '../../types' import { featureCollectionBounds } from '../../lib/geojsonBounds' import { useMapThemeSelectionInsights } from '../../hooks/useMapThemeSelectionInsights' +import { useTemporalComparison } from '../../hooks/useTemporalComparison' const DEFAULT_SELECTED_FEATURE_FILENAME = 'selected-feature.geojson' const DEFAULT_AREA_SELECTION_FILENAME = 'area-selection.geojson' @@ -90,12 +91,42 @@ function pickThemeDataset(datasets: DatasetCreateResponse[], theme: DataTheme): (dataset.reference_layer_name && theme.tokens.includes(dataset.reference_layer_name.toLowerCase()) ? 1_000_000 : 0) + (dataset.source_name === 'grb' ? 100_000 : 0) + (dataset.dataset_role === 'reference' ? 10_000 : 0) + + (dataset.observed_at ? new Date(dataset.observed_at).getTime() / 100_000_000 : 0) + (dataset.feature_count ?? dataset.vector_summary?.feature_count ?? 0) return score(right) - score(left) }) return candidates[0] ?? null } +function pickThemeTemporalSeries(datasets: DatasetCreateResponse[], theme: DataTheme): DatasetCreateResponse[] { + const groups = new Map() + for (const dataset of datasets) { + if (!datasetMatchesTheme(dataset, theme) || !dataset.temporal_series_key || !dataset.observed_at) { + continue + } + const items = groups.get(dataset.temporal_series_key) ?? [] + items.push(dataset) + groups.set(dataset.temporal_series_key, items) + } + return Array.from(groups.values()) + .filter((items) => items.length >= 2) + .sort((left, right) => { + if (right.length !== left.length) { + return right.length - left.length + } + const latest = (items: DatasetCreateResponse[]) => Math.max(...items.map((item) => new Date(item.observed_at ?? 0).getTime())) + return latest(right) - latest(left) + })[0] + ?.sort((left, right) => new Date(left.observed_at ?? 0).getTime() - new Date(right.observed_at ?? 0).getTime()) ?? [] +} + +function formatObservationDate(value: string | null | undefined): string { + if (!value) { + return 'Geen peildatum' + } + return new Intl.DateTimeFormat('nl-BE', { year: 'numeric', month: 'short', day: 'numeric' }).format(new Date(value)) +} + function selectionAreaSquareMetres(bbox: VectorSelectionBBox | null): number | null { if (!bbox) { return null @@ -134,6 +165,19 @@ function resultCountLabel(result: VectorSelectionResponse): string { return result.truncated && result.total_feature_count == null ? `${result.feature_count.toLocaleString('nl-BE')}+` : total.toLocaleString('nl-BE') } +function resultMetricLabel(result: VectorSelectionResponse): string { + if (!result.summary) { + return resultCountLabel(result) + } + const maximumFractionDigits = result.summary.metric_unit === 'inwoners' ? 0 : 2 + return `${result.summary.metric_value.toLocaleString('nl-BE', { maximumFractionDigits })} ${result.summary.metric_unit}` +} + +function formatTemporalMetric(value: number, unit: string): string { + const maximumFractionDigits = unit === 'inwoners' || unit === 'objecten' ? 0 : 2 + return `${value.toLocaleString('nl-BE', { maximumFractionDigits })} ${unit}` +} + function readablePropertyName(value: string): string { return value.replace(/_/g, ' ').replace(/\b\w/g, (character) => character.toUpperCase()) } @@ -444,6 +488,16 @@ export function MapWorkspace({ loadThemeInsights, clearThemeInsights, } = useMapThemeSelectionInsights(selectedProjectId) + const { + temporalComparison, + temporalComparisonLoading, + temporalComparisonError, + compareTemporalSnapshots, + clearTemporalComparison, + } = useTemporalComparison(selectedProjectId) + const [analysisMode, setAnalysisMode] = useState<'current' | 'evolution'>('current') + const [earlierDatasetId, setEarlierDatasetId] = useState('') + const [laterDatasetId, setLaterDatasetId] = useState('') const [bboxSelectionMode, setBboxSelectionMode] = useState(false) const [firstSelectionCorner, setFirstSelectionCorner] = useState<[number, number] | null>(null) const [bboxInput, setBboxInput] = useState(bboxToInputState(mapSelectionBbox)) @@ -482,6 +536,10 @@ export function MapWorkspace({ ) const activeTheme = DATA_THEMES.find((theme) => theme.id === activeThemeId) ?? DATA_THEMES[0] const activeThemeDataset = themeDatasetMap[activeTheme.id] + const activeTemporalSeries = useMemo( + () => pickThemeTemporalSeries(availableMapDatasets, activeTheme), + [activeTheme, availableMapDatasets], + ) const themeResults = useMemo( () => themeInsights.flatMap((insight) => { @@ -502,6 +560,14 @@ export function MapWorkspace({ const selectedDensity = selectedAreaSquareMetres && selectedAreaSquareMetres > 0 ? selectedResultTotal / (selectedAreaSquareMetres / 1_000_000) : null + const activeMetricValue = activeSelectionResult?.summary?.metric_value ?? selectedResultTotal + const activeMetricUnit = activeSelectionResult?.summary?.metric_unit ?? 'objecten' + const activeMetricLabel = activeSelectionResult?.summary?.metric_label ?? activeTheme.shortLabel + const activeSecondaryMetric = selectedAreaSquareMetres && selectedAreaSquareMetres > 0 + ? activeMetricUnit === 'ha' + ? `${((activeMetricValue * 10_000) / selectedAreaSquareMetres * 100).toLocaleString('nl-BE', { maximumFractionDigits: 1 })}% dekking` + : `${(activeMetricValue / (selectedAreaSquareMetres / 1_000_000)).toLocaleString('nl-BE', { maximumFractionDigits: 1 })} ${activeMetricUnit} / km2` + : null const selectedResultProperties = useMemo(() => { const keys = new Map>() for (const feature of activeSelectionResult?.geojson.features ?? []) { @@ -526,6 +592,14 @@ export function MapWorkspace({ setBboxInput(bboxToInputState(mapSelectionBbox)) }, [mapSelectionBbox]) + useEffect(() => { + const first = activeTemporalSeries[0] + const last = activeTemporalSeries[activeTemporalSeries.length - 1] + setEarlierDatasetId(first?.id ?? '') + setLaterDatasetId(last?.id ?? '') + clearTemporalComparison() + }, [activeTemporalSeries]) + useEffect(() => { if (advancedMode || !activeThemeDataset || (selectedMapDataset && datasetMatchesTheme(selectedMapDataset, activeTheme))) { return @@ -562,6 +636,7 @@ export function MapWorkspace({ const startBboxSelection = () => { setFirstSelectionCorner(null) clearThemeInsights() + clearTemporalComparison() setBboxSelectionMode(true) } @@ -590,6 +665,7 @@ export function MapWorkspace({ setFirstSelectionCorner(null) setBboxInput(bboxToInputState(null)) clearThemeInsights() + clearTemporalComparison() onClearMapSelectionExtract() } @@ -644,9 +720,15 @@ export function MapWorkspace({ return } setActiveThemeId(theme.id) + clearTemporalComparison() onOpenDatasetInMap(dataset) } + const setExplorerMode = (mode: 'current' | 'evolution') => { + setAnalysisMode(mode) + clearTemporalComparison() + } + const loadAllThemeResults = async (bbox: VectorSelectionBBox) => { const availableThemes = DATA_THEMES.flatMap((theme) => { const dataset = themeDatasetMap[theme.id] @@ -657,7 +739,18 @@ export function MapWorkspace({ const analyzeSelection = async (bbox: VectorSelectionBBox) => { setSelectionBbox(bbox) - await Promise.all([onRunMapSelectionExtract(bbox), loadAllThemeResults(bbox)]) + const tasks: Array> = [onRunMapSelectionExtract(bbox), loadAllThemeResults(bbox)] + if (analysisMode === 'evolution' && earlierDatasetId && laterDatasetId) { + tasks.push(compareTemporalSnapshots(earlierDatasetId, laterDatasetId, bbox)) + } + await Promise.all(tasks) + } + + const runTemporalComparison = () => { + if (!mapSelectionBbox || !earlierDatasetId || !laterDatasetId) { + return + } + void compareTemporalSnapshots(earlierDatasetId, laterDatasetId, mapSelectionBbox) } const handleMapBboxPreview = (bbox: VectorSelectionBBox) => { @@ -757,6 +850,26 @@ export function MapWorkspace({

Wat bevindt zich in dit gebied?

Kies een datathema, teken een rechthoek en lees de beschikbare gegevens meteen uit.

+
+ + +
@@ -796,15 +909,52 @@ export function MapWorkspace({
- Actieve bron - {activeThemeDataset?.name ?? 'Geen databron beschikbaar'} + {analysisMode === 'evolution' ? 'Tijdreeks' : 'Actieve bron'} + + {analysisMode === 'evolution' + ? activeTemporalSeries[0]?.temporal_series_key ?? 'Geen tijdreeks beschikbaar' + : activeThemeDataset?.name ?? 'Geen databron beschikbaar'} + - {activeThemeDataset - ? `${activeThemeDataset.source_name ?? activeThemeDataset.source} · ${activeThemeDataset.dataset_role ?? 'source'} · EPSG:4326` - : activeTheme.description} + {analysisMode === 'evolution' + ? activeTemporalSeries.length >= 2 + ? `${activeTemporalSeries.length} officiële meetmomenten · ${formatObservationDate(activeTemporalSeries[0].observed_at)} tot ${formatObservationDate(activeTemporalSeries[activeTemporalSeries.length - 1].observed_at)}` + : 'Minstens twee expliciet gedateerde snapshots zijn vereist.' + : activeThemeDataset + ? `${activeThemeDataset.source_name ?? activeThemeDataset.source} · ${activeThemeDataset.dataset_role ?? 'source'} · ${formatObservationDate(activeThemeDataset.observed_at)}` + : activeTheme.description}
+ {analysisMode === 'evolution' ? ( +
+ + + +
+ ) : null} +