diff --git a/notebooks/test_delete_items_from_collection.ipynb b/notebooks/test_delete_items_from_collection.ipynb new file mode 100644 index 0000000..13e9c9b --- /dev/null +++ b/notebooks/test_delete_items_from_collection.ipynb @@ -0,0 +1,331 @@ +{ + "cells": [ + { + "metadata": {}, + "cell_type": "markdown", + "source": "# tests to delete unwanted modelIDs from the modelSTAC", + "id": "7187881d346266b1" + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:25:06.246567Z", + "start_time": "2026-09-15T08:24:59.582135700Z" + } + }, + "cell_type": "code", + "source": [ + "import pystac_client\n", + "import pandas as pd\n", + "from eo_processing.utils.storage import WEED_storage" + ], + "id": "187fe67715113a74", + "outputs": [], + "execution_count": 1 + }, + { + "cell_type": "code", + "id": "initial_id", + "metadata": { + "collapsed": true, + "ExecuteTime": { + "end_time": "2026-09-15T08:25:18.518515200Z", + "start_time": "2026-09-15T08:25:18.506034800Z" + } + }, + "source": [ + "username = 'VAULT_TOKEN' # or your terrascope user name to hand over credentials manually\n", + "# details to STAC\n", + "stac_url:str = 'https://catalogue.weed.apex.esa.int'\n", + "collection_id:str = 'model-STAC'" + ], + "outputs": [], + "execution_count": 2 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:25:22.668653500Z", + "start_time": "2026-09-15T08:25:20.543609200Z" + } + }, + "cell_type": "code", + "source": [ + "# get all itemIDs from modelSTAC\n", + "client = pystac_client.Client.open(stac_url)\n", + "\n", + "search = client.search(\n", + " collections=[collection_id],\n", + " fields=[\"id\", \"properties.topology\", \"properties.training_year\", \"properties.model_version\", \"properties.name_spatial_region\", \"properties.name_spatial_zone\", \"properties.modelID\"],\n", + ")\n", + "\n", + "# run the search\n", + "results = []\n", + "for item in search.items_as_dicts():\n", + " results.append([item['id'], item['properties']['modelID'], item['properties']['topology'], item['properties']['training_year'], item['properties']['model_version'], item['properties']['name_spatial_region'], item['properties']['name_spatial_zone']])\n", + "\n", + "df_result = pd.DataFrame(results, columns=['item_id', 'modelID', 'typology', 'training_year', 'model_version', 'spatial_region', 'spatial_zone'])" + ], + "id": "15d488e9fe1470d0", + "outputs": [], + "execution_count": 3 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:25:25.181893800Z", + "start_time": "2026-09-15T08:25:25.131787900Z" + } + }, + "cell_type": "code", + "source": [ + "# get the overview\n", + "df_result.head()" + ], + "id": "7e0052d61288cea", + "outputs": [ + { + "data": { + "text/plain": [ + " item_id modelID typology \\\n", + "0 IUCNGET_global_v321_2024_IND IUCNGET_global_v321_2024_IND IUCNGET \n", + "1 IUCNGET_global_v315_2024_AFR IUCNGET_global_v315_2024_AFR IUCNGET \n", + "2 IUCNGET_global_v314_2024_NEO IUCNGET_global_v314_2024_NEO IUCNGET \n", + "3 IUCNGET_global_v314_2024_IND IUCNGET_global_v314_2024_IND IUCNGET \n", + "4 IUCNGET_global_v314_2024_AFR IUCNGET_global_v314_2024_AFR IUCNGET \n", + "\n", + " training_year model_version spatial_region spatial_zone \n", + "0 2024 3.21 global IndoMalesian \n", + "1 2024 3.15 global African \n", + "2 2024 3.14 global Neotropical \n", + "3 2024 3.14 global IndoMalesian \n", + "4 2024 3.14 global African " + ], + "text/html": [ + "
\n", + "\n", + "\n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + " \n", + "
item_idmodelIDtypologytraining_yearmodel_versionspatial_regionspatial_zone
0IUCNGET_global_v321_2024_INDIUCNGET_global_v321_2024_INDIUCNGET20243.21globalIndoMalesian
1IUCNGET_global_v315_2024_AFRIUCNGET_global_v315_2024_AFRIUCNGET20243.15globalAfrican
2IUCNGET_global_v314_2024_NEOIUCNGET_global_v314_2024_NEOIUCNGET20243.14globalNeotropical
3IUCNGET_global_v314_2024_INDIUCNGET_global_v314_2024_INDIUCNGET20243.14globalIndoMalesian
4IUCNGET_global_v314_2024_AFRIUCNGET_global_v314_2024_AFRIUCNGET20243.14globalAfrican
\n", + "
" + ] + }, + "execution_count": 4, + "metadata": {}, + "output_type": "execute_result" + } + ], + "execution_count": 4 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:25:26.828837600Z", + "start_time": "2026-09-15T08:25:26.806822200Z" + } + }, + "cell_type": "code", + "source": [ + "# selecte items to delete\n", + "df_del = df_result[~(df_result['typology'] == 'IUCNGET')]\n", + "\n", + "ldelete = df_del[df_del['model_version'] < 3].item_id.tolist()" + ], + "id": "73903f37bab568f", + "outputs": [], + "execution_count": 5 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:25:28.072842500Z", + "start_time": "2026-09-15T08:25:28.011634200Z" + } + }, + "cell_type": "code", + "source": "ldelete", + "id": "fea9dbbc06a3d4b8", + "outputs": [ + { + "data": { + "text/plain": [ + "['EUNIS2021plus_panEU_v201_2024_OneZone',\n", + " 'EUNIS2021plus_EU_v1_2024_STE',\n", + " 'EUNIS2021plus_EU_v1_2024_PAN',\n", + " 'EUNIS2021plus_EU_v1_2024_MED',\n", + " 'EUNIS2021plus_EU_v1_2024_CON',\n", + " 'EUNIS2021plus_EU_v1_2024_BOR',\n", + " 'EUNIS2021plus_EU_v1_2024_ATL',\n", + " 'EUNIS2021plus_EU_v1_2024_ALP']" + ] + }, + "execution_count": 6, + "metadata": {}, + "output_type": "execute_result" + } + ], + "execution_count": 6 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:25:29.386166Z", + "start_time": "2026-09-15T08:25:29.308068100Z" + } + }, + "cell_type": "code", + "source": [ + "# init the storage object\n", + "st = WEED_storage(username=username, s3_bucket='model', project='WEED', stac_env='prod')" + ], + "id": "4d5fe40f5ffc5142", + "outputs": [], + "execution_count": 7 + }, + { + "metadata": { + "ExecuteTime": { + "end_time": "2026-09-15T08:26:50.696816200Z", + "start_time": "2026-09-15T08:26:48.614625600Z" + } + }, + "cell_type": "code", + "source": [ + "# now we delete the items we do not want\n", + "for item in ldelete:\n", + " print(f'delete item: {item}')\n", + " st.delete_collection_item(collection_name=collection_id, item_id=item)" + ], + "id": "66216cd955542da7", + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "delete item: EUNIS2021plus_panEU_v201_2024_OneZone\n", + "Item 'EUNIS2021plus_panEU_v201_2024_OneZone' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_STE\n", + "Item 'EUNIS2021plus_EU_v1_2024_STE' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_PAN\n", + "Item 'EUNIS2021plus_EU_v1_2024_PAN' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_MED\n", + "Item 'EUNIS2021plus_EU_v1_2024_MED' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_CON\n", + "Item 'EUNIS2021plus_EU_v1_2024_CON' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_BOR\n", + "Item 'EUNIS2021plus_EU_v1_2024_BOR' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_ATL\n", + "Item 'EUNIS2021plus_EU_v1_2024_ATL' deleted successfully from collection 'model-STAC'\n", + "delete item: EUNIS2021plus_EU_v1_2024_ALP\n", + "Item 'EUNIS2021plus_EU_v1_2024_ALP' deleted successfully from collection 'model-STAC'\n" + ] + } + ], + "execution_count": 8 + }, + { + "metadata": {}, + "cell_type": "code", + "outputs": [], + "execution_count": null, + "source": "", + "id": "857e74f1803a5361" + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 2 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython2", + "version": "2.7.6" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/src/eo_processing/resources/udf_force_into_proba.py b/src/eo_processing/resources/udf_force_into_proba.py new file mode 100644 index 0000000..6e5d6bc --- /dev/null +++ b/src/eo_processing/resources/udf_force_into_proba.py @@ -0,0 +1,228 @@ +import xarray as xr +from typing import Dict +from openeo.udf import inspect +from openeo.metadata import CubeMetadata + +def apply_metadata(metadata: CubeMetadata, context:Dict) -> CubeMetadata: + """ Rename the bands by using openeo apply_metadata function + :param metadata: Metadata of the input data cube + :param context: Context of the UDF + :return: renamed labels + """ + band_names = context.get("band_names") + typology_schema = context.get("typology_schema") + force_ocean = context.get("force_ocean") + force_snow = context.get("force_snow") + mask_non_mangrove = context.get("mask_non_mangrove") + + if force_ocean: + if typology_schema == "EUNIS2021plus": + needed = "Level1_class-0_habitat-MH-20000" + elif typology_schema == "IUCNGET": + needed = "Level1_class-0_habitat-M-80000" + else: + raise Exception("unknown typology schema - please adjust the code base of UDF") + if needed not in band_names: + band_names.append(needed) + + if force_snow: + if typology_schema == "EUNIS2021plus": + needed = "Level1_class-0_habitat-U-90000" + elif typology_schema == "IUCNGET": + needed = "Level1_class-0_habitat-T-10000" + else: + raise Exception("unknown typology schema - please adjust the code base of UDF") + if needed not in band_names: + band_names.append(needed) + + if mask_non_mangrove: + if typology_schema == "IUCNGET": + needed = "Level1_class-0_habitat-MFT-100000" + else: + raise Exception("currently mangrove masking is only possible in IUCNGET. set parameter to False for " + "EUNIS or adjust the code base of UDF for new typology.") + if needed not in band_names: + band_names.append(needed) + + return metadata.rename_labels(dimension="bands", target=band_names) + +def apply_datacube(cube: xr.DataArray, context: Dict) -> xr.DataArray: + """ imprint or mask external data into the probability cube + + :param cube: data cube + :param context: dictionary to provide external data - key class_mapping is used to provide external re-mapping dict + :return: data cube with remapped values (new cube) + """ + + # get all needed data together + output_band_names = context.get("band_names") + typology_schema = context.get("typology_schema") + force_ocean = context.get("force_ocean") + force_snow = context.get("force_snow") + mask_non_mangrove = context.get("mask_non_mangrove") + mask_non_ocean = context.get("mask_non_ocean") + inspect(message=f"settings: typology_schema: {typology_schema}, force_ocean: {force_ocean}, force_snow: " + f"{force_snow}, mask_non_mangrove: {mask_non_mangrove}, mask_non_ocean: {mask_non_ocean}," + f"output_band_names ({len(output_band_names)}): {output_band_names}") + + # get the band names of cube handed over to UDF + input_band_names = cube.indexes["bands"].values + inspect(message=f"input cube band names ({len(input_band_names)}): {input_band_names}") + + # check that we have to imprint "ocean" + if force_ocean: + inspect(message="imprinting ocean") + # check that we have the needed band + if "ocean" not in input_band_names: + raise ValueError("no ocean band in cube - please add a ocean band to the input cube") + + cube_ocean = cube.sel(bands="ocean") + ocean_value = 255 + + if typology_schema == "EUNIS2021plus": + imprint_band = "Level1_class-0_habitat-MH-20000" + elif typology_schema == "IUCNGET": + imprint_band = "Level1_class-0_habitat-M-80000" + else: + raise Exception("unknown typology schema - please adjust the code base of UDF") + + # check that we have the needed output band, if not create + if imprint_band not in input_band_names: + inspect(message=f"+ create extra proba layer for ocean ({imprint_band})") + output_band_names.append(imprint_band) + # Create zero-filled band with same spatial dimensions as cube + zero_band = xr.zeros_like(cube.isel(bands=0)) + zero_band = zero_band.assign_coords(bands=imprint_band) + # Add the new band to the cube + cube = xr.concat([cube, zero_band], dim="bands") + + # now we imprint the ocean in band "Level1_class-0_habitat-M-80000" + ocean_mask = cube_ocean == ocean_value + inspect(message=f"+ imprint ocean mask into band {imprint_band}") + cube.loc[dict(bands=imprint_band)] = xr.where(ocean_mask, 100, cube.loc[dict(bands=imprint_band)]) + + # reset the values of all other bands on level 1 to ZERO + other_bands_level1 = [x for x in input_band_names if x.startswith("Level1_class-0")] + if imprint_band in other_bands_level1: + other_bands_level1.remove(imprint_band) + inspect(message=f"+ imprint ocean mask into non-marine bands ({other_bands_level1})") + cube.loc[dict(bands=other_bands_level1)] = xr.where(ocean_mask, 0, cube.loc[dict(bands=other_bands_level1)]) + + # now mask areas wrongly classified as ocean (e.g. in the case of the ocean mask) + # NOTE: currently only needed for IUCN-GET typology + if mask_non_ocean and (typology_schema == "IUCNGET"): + inspect(message=f"+ reset areas in the {imprint_band} band which are not really ocean. " + f"(level1 second winner will be used)") + # get the mask of the non-ocean areas + non_ocean_mask = ~ocean_mask + # now mask the non-ocean areas + cube.loc[dict(bands=imprint_band)] = xr.where(non_ocean_mask, 0, cube.loc[dict(bands=imprint_band)]) + + # check that we have to imprint "snow" + if force_snow: + inspect(message="imprinting snow") + # check that we have the needed band + if "snow" not in input_band_names: + raise ValueError("no snow band in cube - please add a snow band to the input cube") + + snow_cube = cube.sel(bands="snow") + snow_value = 110 # value in LCFM 10m maps for permanent snow & ice + + if typology_schema == "EUNIS2021plus": + imprint_band = "Level1_class-0_habitat-U-90000" + elif typology_schema == "IUCNGET": + imprint_band = "Level1_class-0_habitat-T-10000" + else: + raise Exception("unknown typology schema - please adjust the code base of UDF") + + # check that we have the needed output band, if not create + if imprint_band not in input_band_names: + inspect(message=f"+ create extra proba layer for snow ({imprint_band})") + output_band_names.append(imprint_band) + # Create zero-filled band with same spatial dimensions as cube + zero_band = xr.zeros_like(cube.isel(bands=0)) + zero_band = zero_band.assign_coords(bands=imprint_band) + # Add the new band to the cube + cube = xr.concat([cube, zero_band], dim="bands") + + # create mask and imprint snow into band + snow_mask = snow_cube == snow_value + inspect(message=f"+ imprint snow mask into band {imprint_band}") + cube.loc[dict(bands=imprint_band)] = xr.where(snow_mask, 100, cube.loc[dict(bands=imprint_band)]) + + # reset the values of all other bands on level 1 to ZERO + other_bands_level1 = [x for x in input_band_names if x.startswith("Level1_class-0")] + if imprint_band in other_bands_level1: + other_bands_level1.remove(imprint_band) + inspect(message=f"+ imprint snow mask into non-snow bands ({other_bands_level1})") + cube.loc[dict(bands=other_bands_level1)] = xr.where(snow_mask, 0, cube.loc[dict(bands=other_bands_level1)]) + + # check that we have to mask mangrove + if mask_non_mangrove: + inspect(message="imprinting mangroves") + # check that we have the needed band + if "mangrove" not in input_band_names: + raise ValueError("no mangrove band in cube - please add a mangrove band to the input cube") + + mango_cube = cube.sel(bands="mangrove") + mango_value = 1 # raster value for potential areas of mangrove + + if typology_schema == "IUCNGET": + imprint_band = "Level1_class-0_habitat-MFT-100000" + else: + raise Exception("currently mangrove masking is only possible in IUCNGET. set parameter to False for " + "EUNIS or adjust the code base of UDF for new typology.") + + # check that we have the needed output band, if not create + if imprint_band not in input_band_names: + inspect(message=f"+ create extra proba layer for mangrove ({imprint_band})") + output_band_names.append(imprint_band) + # Create zero-filled band with same spatial dimensions as cube + zero_band = xr.zeros_like(cube.isel(bands=0)) + zero_band = zero_band.assign_coords(bands=imprint_band) + # Add the new band to the cube + cube = xr.concat([cube, zero_band], dim="bands") + + # We have to first determine the level3 winner in the level1 MFT results. When the winner is mangrove (MFT1.2) and + # outside the mangrove potential area mask THEN we can set the level1 pixel to ZERO. (Again under the assumption that + # these areas can not be MFT1.1 or MFT1.3) + # since level2 is a single class model we do not have to check it + + level3_bands = ['Level3_class-MFT1_habitat-MFT1.1-100101', + 'Level3_class-MFT1_habitat-MFT1.2-100102', + 'Level3_class-MFT1_habitat-MFT1.3-100103'] + + if all(band in input_band_names for band in level3_bands): + # get the level3 winner + level3_data = cube.sel(bands=level3_bands).copy() + level3_data = level3_data.fillna(0) + level3_winner = level3_data.sel(bands=level3_bands).argmax('bands') + # get mask of MFT1.2 winner + level3_mangrove_winner = level3_winner == 1 + # check if value is bigger than 0 + level3_mangrove_winner_value = level3_data.sel(bands='Level3_class-MFT1_habitat-MFT1.2-100102') > 0 + + # create mask where mangrove can exist + mango_mask = mango_cube == mango_value + + # create the removal mask -> MFT1.2 outside the mangrove potential area mask + removal_mask = ~mango_mask & level3_mangrove_winner & level3_mangrove_winner_value + + n_pixels_to_remove = int(removal_mask.sum().values) + + # reset the level1 MFT values to ZERO in these areas + inspect( + message=f"+ reset areas in the {imprint_band} band were wrongly mangroves were detected ({n_pixels_to_remove} pixel)" + f"(level1 second winner will be chosen)") + cube.loc[dict(bands=imprint_band)] = xr.where(removal_mask, 0, cube.loc[dict(bands=imprint_band)]) + else: + inspect(message=f"+ no level3 bands found in cube - cannot mask mangroves") + + # filter the cube to only the bands we need + cube = cube.sel(bands=output_band_names) + inspect(message=f"output cube band names {len(cube.indexes['bands'].values)}: {cube.indexes['bands'].values}") + # final check that we have the right number of bands + if len(cube.indexes['bands'].values) != len(output_band_names): + raise ValueError("wrong number of bands in output cube - please adjust the code base of UDF") + + return cube \ No newline at end of file diff --git a/src/eo_processing/utils/stac_helper.py b/src/eo_processing/utils/stac_helper.py index e88e2f3..046a87c 100644 --- a/src/eo_processing/utils/stac_helper.py +++ b/src/eo_processing/utils/stac_helper.py @@ -1,5 +1,10 @@ import pystac_client - +import logging +from urllib.request import urlopen +from io import BytesIO +from typing import Optional, List +import geopandas as gpd +import pandas as pd def get_stac_collection_url(collection_id: str, catalog_url: str = "https://catalogue.weed.apex.esa.int/") -> str: """ @@ -25,7 +30,7 @@ def get_stac_collection_url(collection_id: str, catalog_url: str = "https://cata def query_modelID_asset_url(model_id: str, catalog_url: str ="https://catalogue.weed.apex.esa.int/", - collection_id: str ="model-STAC-v2") -> str: + collection_id: str ="model-STAC") -> str: """ Queries the URL of a specific asset for a given modelID from a STAC catalog. @@ -50,4 +55,163 @@ def query_modelID_asset_url(model_id: str, filter_lang="cql2-json", ) item_collection = search.item_collection() - return item_collection.items[0].assets["model_valid_geometry"].href \ No newline at end of file + if not item_collection.items: + logging.error(f"No items found for model_id: {model_id}") + raise ValueError(f"No items found for model_id: {model_id}") + + return item_collection.items[0].assets["model_valid_geometry"].href + +def query_modelID_output_bands(model_id: str, + catalog_url: str ="https://catalogue.weed.apex.esa.int/", + collection_id: str ="model-STAC") -> List[str]: + """ + Queries the output bands associated with the specified model ID from a STAC catalog. + + This function communicates with a STAC catalog to retrieve the model's output + bands by searching for items using their model ID. It uses CQL2 filtering + to perform the search and returns the output band names from the first matching + item's properties. + + :param model_id: (str) The unique identifier for the model to query. + :param catalog_url: (str) The URL of the STAC catalog to connect to. Defaults to + "https://catalogue.weed.apex.esa.int/". + :param collection_id: (str) The ID of the STAC collection to search within. Defaults + to "model-STAC". + :return: A list of output band names (List[str]) associated with the given + model ID. + :raises ValueError: If no items are found in the catalog for the provided model ID. + """ + client = pystac_client.Client.open(catalog_url) + + search = client.search( + limit=20, + collections=[collection_id], + filter={"op": "=", "args": [{"property": "properties.modelID"}, model_id]}, + filter_lang="cql2-json", + ) + item_collection = search.item_collection() + if not item_collection.items: + logging.error(f"No items found for model_id: {model_id}") + raise ValueError(f"No items found for model_id: {model_id}") + + return item_collection.items[0].properties["output_band_names"] + +def get_modelID_asset_geometry_from_STAC(df_AOI: gpd.GeoDataFrame, typology_schema: str = 'IUCNGET', + model_version: Optional[str]=None, + stac_url:str = 'https://catalogue.weed.apex.esa.int', + collection_id:str = 'model-STAC', info_debug: bool = True) -> gpd.GeoDataFrame: + """ + Retrieves geometries associated with model IDs from a STAC (SpatioTemporal Asset Catalog) collection. + + This function queries a STAC catalog to find model records based on an input area of interest (AOI), + typology schema, and optionally a specific model version. The geometries associated with the + models are retrieved and returned as a GeoDataFrame. + + Parameters + ---------- + df_AOI : gpd.GeoDataFrame + The area of interest represented as a GeoDataFrame. Must have a valid coordinate reference + system (CRS). If not in EPSG:4326, it will be converted to this CRS. + typology_schema : str + A typology schema filter (e.g., "IUCNGET") to apply to the STAC search. This helps filter + results by their topology type. Default is "IUCNGET". + model_version : Optional[str] + A specific model version (e.g., "1.0") to filter the search results. If None, all versions + are considered. Default is None. + stac_url : str + The URL of the STAC catalog to query. Default is 'https://catalogue.weed.apex.esa.int'. + collection_id : str + The ID of the STAC collection to search within. Default is 'model-STAC'. + info_debug : bool + if additional debug messages should be printed. Default is True. + + Returns + ------- + gpd.GeoDataFrame + A GeoDataFrame containing the retrieved model IDs, their properties, and associated geometries. + The GeoDataFrame is in EPSG:4326 CRS. + + Raises + ------ + ValueError + If no intersecting model IDs are found in the specified STAC collection. + """ + if info_debug: print(f"get_modelID_asset_geometry_from_STAC") + # convert AOI into BBOX in 4326 + if df_AOI.crs != "EPSG:4326": + df_AOI_4326 = df_AOI.to_crs("EPSG:4326") + else: + df_AOI_4326 = df_AOI.copy() + + bbox_4326 = df_AOI_4326.total_bounds + + # init the pySTAC search + if info_debug: print(f"- searching for modelIDs in {bbox_4326} with typology schema {typology_schema}") + if info_debug: print(f"- using STAC url: {stac_url}") + if info_debug: print(f"- using collection id: {collection_id}") + client = pystac_client.Client.open(stac_url) + + search = client.search( + collections=[collection_id], + bbox=bbox_4326, + filter={"op": "=", "args": [{"property": "properties.topology"}, typology_schema]}, + filter_lang="cql2-json", + fields=["properties.modelID", "properties.topology", "properties.training_year", "properties.model_version", "properties.name_spatial_region", "properties.name_spatial_zone", "assets.model_valid_geometry"], + ) + + # run the search + results = [] + for item in search.items_as_dicts(): + if model_version is None: + results.append([item['properties']['modelID'], item['properties']['topology'], + item['properties']['training_year'], item['properties']['model_version'], + item['properties']['name_spatial_region'], item['properties']['name_spatial_zone'], item['assets']['model_valid_geometry']['href']]) + else: + if str(item['properties']['model_version']) == model_version: + results.append([item['properties']['modelID'], item['properties']['topology'], + item['properties']['training_year'], item['properties']['model_version'], + item['properties']['name_spatial_region'], item['properties']['name_spatial_zone'], item['assets']['model_valid_geometry']['href']]) + # build dataframe + df_result = pd.DataFrame(results, columns=['modelID', 'typology', 'training_year', 'model_version', 'spatial_region', 'spatial_zone', 'asset']) + if info_debug: print(f"- found {len(df_result)} intersecting modelIDs") + # check if there are any results + if df_result.empty: + ValueError("No intersecting modelIDs found in the modelSTAC.") + + ## load the modelIS asset geometry and add + if info_debug: print(f"- loading modelIDS geometries from STAC assets") + # Initialize an empty list to store geometries + geometries = [] + + # Loop over each row in df_result + for idx, row in df_result.iterrows(): + parquet_url = row['asset'] + + try: + # Download and read parquet file in memory + with urlopen(parquet_url) as response: + parquet_data = response.read() + + # Read from bytes + temp_gdf = gpd.read_parquet(BytesIO(parquet_data)) + + # Extract the geometry (assuming there's only one geometry per file) + # If there are multiple geometries, you can use .union_all() or take the first one + if len(temp_gdf) > 0: + geometry = temp_gdf.geometry.iloc[0] + else: + geometry = None + + geometries.append(geometry) + + if info_debug: print(f"- Successfully loaded geometry for {row['modelID']}") + + except Exception as e: + print(f"Error loading {row['modelID']}: {e}") + geometries.append(None) + + # Add geometries as a new column to df_result + df_result['geometry'] = geometries + + # Convert df_result to a GeoDataFrame + return gpd.GeoDataFrame(df_result, geometry='geometry', crs='EPSG:4326') diff --git a/src/eo_processing/utils/storage.py b/src/eo_processing/utils/storage.py index c319101..1609fe9 100644 --- a/src/eo_processing/utils/storage.py +++ b/src/eo_processing/utils/storage.py @@ -980,7 +980,7 @@ def QueryItems(self, table: str, lcolumns: List[str]) -> List[tuple]: return lresults - def AddColumns(self, table: str, column_names: lst(str)) -> None: + def AddColumns(self, table: str, column_names: List[str]) -> None: """ Adds a new column to a specified table in the database. The method dynamically constructs an SQL ALTER TABLE statement to add the column with the specified name, ensuring @@ -1451,6 +1451,26 @@ def delete_collection(self, collection_name: str) -> None: else: print(f"Failed to delete collection: HTTP {resp.status_code}\n{resp.text}") + def delete_collection_item(self, collection_name: str, item_id: str) -> None: + """ + Deletes a specified item from a collection in the catalog. This method constructs + the URL for the item, uses authentication to validate the request, and sends a + delete request to remove the item. If the deletion is successful, a confirmation + message is printed. Otherwise, an error message with the appropriate HTTP status + code and error details is displayed. + + :param collection_name: The name of the collection containing the item. + :param item_id: The ID of the item to delete. + """ + catalog_url = self.get_catalog_url().rstrip('/') + auth_token = self.get_bearer_auth() + item_url = f"{catalog_url}/collections/{collection_name}/items/{item_id}" + resp = delete(item_url, auth=auth_token) + if resp.status_code == 204: + print(f"Item '{item_id}' deleted successfully from collection '{collection_name}'") + else: + print(f"Failed to delete item: HTTP {resp.status_code}\n{resp.text}") + class ReadFaker: """ A class that mimics file reading behavior for data in tabular formats.