diff --git a/geonode/upload/api/tests.py b/geonode/upload/api/tests.py index ddba0ebaf02..5ebc31bdb7d 100644 --- a/geonode/upload/api/tests.py +++ b/geonode/upload/api/tests.py @@ -31,6 +31,7 @@ from geonode.base.populate_test_data import create_single_dataset from django.http import HttpResponse, QueryDict +from geonode.resource.models import ExecutionRequest from geonode.upload.models import ResourceHandlerInfo from geonode.upload.tests.utils import ImporterBaseTestSupport from geonode.upload.utils import UploadLimitValidator @@ -400,6 +401,35 @@ def test_geojson_task_is_called(self, patch_upload): self.assertTrue(201, response.status_code) + @patch("geonode.upload.api.views.import_orchestrator") + def test_skip_existing_layers_is_propagated_to_execution_params(self, patch_upload): + patch_upload.apply_async.side_effect = MagicMock() + + self.client.force_login(get_user_model().objects.get(username="admin")) + for flag in (True, False, None): + with self.subTest(flag=flag): + payload = { + "base_file": SimpleUploadedFile( + name="issue14377.geojson", + content=( + b'{"type": "FeatureCollection", "features": ' + b'[{"type": "Feature", "properties": {}, "geometry": ' + b'{"type": "Point", "coordinates": [1, 2]}}]}' + ), + ), + "store_spatial_files": True, + "action": "upload", + } + + if flag is not None: + payload["skip_existing_layers"] = flag + response = self.client.post(self.url, data=payload) + + self.assertEqual(201, response.status_code) + execution = ExecutionRequest.objects.get(exec_id=response.json()["execution_id"]) + self.assertIs(execution.input_params.get("skip_existing_layer"), flag is True) + self.assertNotIn("skip_existing_layers", execution.input_params) + @patch("geonode.upload.api.views.import_orchestrator") def test_geojson_mixed_geometry_succed(self, patch_upload): patch_upload.apply_async.side_effect = MagicMock() diff --git a/geonode/upload/handlers/common/raster.py b/geonode/upload/handlers/common/raster.py index 535ddec9148..dde6c4ee9a5 100644 --- a/geonode/upload/handlers/common/raster.py +++ b/geonode/upload/handlers/common/raster.py @@ -127,7 +127,7 @@ def extract_params_from_data(_data, action=None): return {"title": data.pop("title"), "store_spatial_file": True}, _data return { - "skip_existing_layers": _data.pop("skip_existing_layers", "False"), + "skip_existing_layer": _data.pop("skip_existing_layers", False), "resource_pk": _data.pop("resource_pk", None), "store_spatial_file": _data.pop("store_spatial_files", "True"), "action": _data.pop("action", "upload"), @@ -341,7 +341,10 @@ def import_resource(self, files: dict, execution_id: str, **kwargs) -> str: except Exception as e: logger.error(e) raise e - return + raise ImportException( + "No new layers were detected in your upload. " + "Existing layers were left unchanged, so no updates were made." + ) def create_geonode_resource( self, layer_name: str, alternate: str, execution_id: str, resource_type: Dataset = Dataset, asset=None, **kwargs diff --git a/geonode/upload/handlers/common/tests_raster.py b/geonode/upload/handlers/common/tests_raster.py index 66699a18d0a..633c75e6655 100644 --- a/geonode/upload/handlers/common/tests_raster.py +++ b/geonode/upload/handlers/common/tests_raster.py @@ -25,6 +25,7 @@ from geonode.upload.orchestrator import orchestrator from geonode.base.populate_test_data import create_single_dataset from geonode.resource.models import ExecutionRequest +from geonode.upload.api.exceptions import ImportException class TestBaseRasterFileHandler(TestCase): @@ -77,6 +78,30 @@ def test_import_resource_should_not_be_imported(self, celery_chord, ogr2ogr_driv if exec_id: ExecutionRequest.objects.filter(exec_id=exec_id).delete() + @patch("geonode.upload.handlers.common.raster.import_orchestrator.apply_async") + def test_import_resource_should_skip_existing_layer(self, import_orchestrator): + exec_id = None + try: + create_single_dataset(name="test_raster", owner=self.owner) + exec_id = orchestrator.create_execution_request( + user=self.owner, + func_name="funct1", + step="step", + input_params={"files": self.valid_files, "skip_existing_layer": True}, + ) + + with self.assertRaisesMessage( + ImportException, + "No new layers were detected in your upload. " + "Existing layers were left unchanged, so no updates were made.", + ): + self.handler.import_resource(files=self.valid_files, execution_id=str(exec_id)) + + import_orchestrator.assert_not_called() + finally: + if exec_id: + ExecutionRequest.objects.filter(exec_id=exec_id).delete() + @patch("geonode.upload.handlers.common.raster.import_orchestrator.apply_async") def test_import_resource_should_work(self, import_orchestrator): try: diff --git a/geonode/upload/handlers/common/tests_vector.py b/geonode/upload/handlers/common/tests_vector.py index c5ecfaabad9..1aab256b32a 100644 --- a/geonode/upload/handlers/common/tests_vector.py +++ b/geonode/upload/handlers/common/tests_vector.py @@ -609,6 +609,82 @@ def test_select_valid_layers(self): self.assertEqual(1, len(valid_layer)) self.assertEqual("mattia_test", valid_layer[0].GetName()) + def test_existing_layers_are_filtered_only_during_import(self): + handler = GPKGFileHandler() + execution_id = orchestrator.create_execution_request( + user=self.owner, + func_name="import_resource", + step="geonode.upload.import_resource", + action="upload", + input_params={"skip_existing_layer": True}, + ) + source = handler.open_source_file(self.valid_files) + layers = handler._select_valid_layers(source, execution_id=str(execution_id)) + self.assertEqual([self.layer.name], [layer.GetName() for layer in layers]) + published = handler.extract_resource_to_publish(self.valid_files, "upload", self.layer.name, self.layer.name) + self.assertEqual(self.layer.name, published[0]["name"]) + with self.assertLogs("importer", level="INFO") as logs: + with self.assertRaisesMessage(ImportException, "No new layers were detected in your upload."): + handler._select_valid_layers(source, execution_id=str(execution_id), filter_existing=True) + self.assertIn(f"Skipping existing layer: {self.layer.name}", [record.message for record in logs.records]) + + def test_skip_filter_preserves_disabled_flag_and_other_owners(self): + handler = GPKGFileHandler() + other_user, _ = get_user_model().objects.get_or_create(username="skip-existing-other-owner") + for user, params in ( + (self.owner, {}), + (self.owner, {"skip_existing_layer": False}), + (other_user, {"skip_existing_layer": True}), + ): + with self.subTest(user=user, params=params): + execution_id = orchestrator.create_execution_request( + user=user, + func_name="import_resource", + step="geonode.upload.import_resource", + action="upload", + input_params=params, + ) + layers = handler._select_valid_layers( + handler.open_source_file(self.valid_files), execution_id=str(execution_id), filter_existing=True + ) + self.assertEqual([self.layer.name], [layer.GetName() for layer in layers]) + + @patch("geonode.upload.handlers.common.vector.BaseVectorFileHandler.open_source_file", return_value=[None]) + def test_import_resource_uses_extracted_skip_flag_for_filter_existing(self, open_source): + for flag in (True, False, None): + with self.subTest(flag=flag): + data = {} if flag is None else {"skip_existing_layers": flag} + input_params, _ = self.handler.extract_params_from_data(data) + execution_id = orchestrator.create_execution_request( + user=self.owner, + func_name="import_resource", + step="geonode.upload.import_resource", + action="upload", + input_params=input_params, + ) + + with patch.object(self.handler, "_select_valid_layers", return_value=[]) as select_layers: + with self.assertRaisesMessage(Exception, "No valid layers found"): + self.handler.import_resource(self.valid_files, str(execution_id)) + + select_layers.assert_called_once_with( + [None], + execution_id=str(execution_id), + filter_existing=flag is True, + ) + + @patch("geonode.upload.handlers.common.vector.BaseVectorFileHandler.open_source_file", return_value=[None]) + def test_empty_input_does_not_report_existing_layers(self, open_source): + execution_id = orchestrator.create_execution_request( + user=self.owner, + func_name="import_resource", + step="geonode.upload.import_resource", + action="upload", + input_params={"skip_existing_layer": True}, + ) + with self.assertRaisesMessage(Exception, "No valid layers found"): + self.handler.import_resource(self.valid_files, str(execution_id)) + @override_settings(MEDIA_ROOT="/tmp") def test_perform_last_step(self): """ diff --git a/geonode/upload/handlers/common/vector.py b/geonode/upload/handlers/common/vector.py index c8249fcdfdb..5ed3821bfa5 100644 --- a/geonode/upload/handlers/common/vector.py +++ b/geonode/upload/handlers/common/vector.py @@ -195,7 +195,7 @@ def extract_params_from_data(_data, action=None): return {"title": data.pop("title"), "store_spatial_file": True}, _data return { - "skip_existing_layers": _data.pop("skip_existing_layers", "False"), + "skip_existing_layer": _data.pop("skip_existing_layers", False), "resource_pk": _data.pop("resource_pk", None), "store_spatial_file": _data.pop("store_spatial_files", "True"), "action": _data.pop("action", "upload"), @@ -457,11 +457,15 @@ def import_resource(self, files: dict, execution_id: str, **kwargs) -> str: data inside the geonode_data database """ gdal_proxy = self.open_source_file(files) - layers = self._select_valid_layers(gdal_proxy, execution_id=execution_id) + _exec = self._get_execution_request_object(execution_id) + layers = self._select_valid_layers( + gdal_proxy, + execution_id=execution_id, + filter_existing=_exec.input_params.get("skip_existing_layer", False), + ) # for the moment we skip the dyanamic model creation layer_count = len(layers) logger.info(f"Total number of layers available: {layer_count}") - _exec = self._get_execution_request_object(execution_id) _input = {**_exec.input_params, **{"total_layers": layer_count}} orchestrator.update_execution_request_status(execution_id=str(execution_id), input_params=_input) dynamic_model = None @@ -482,71 +486,61 @@ def import_resource(self, files: dict, execution_id: str, **kwargs) -> str: layer_name = self.fixup_name(layer.GetName()) should_be_overwritten = _exec.action == ira.REPLACE.value - # should_be_imported check if the user+layername already exists or not - if ( - should_be_imported( - layer_name, - _exec.user, - skip_existing_layer=_exec.input_params.get("skip_existing_layer"), - overwrite_existing_layer=should_be_overwritten, - ) - # and layer.GetGeometryColumn() is not None - ): - _files = files.copy() - vrt_layer_name = None - if has_incompatible_field_names(layer): - vrt_filename, vrt_layer_name = create_vrt_file(layer, files.get("base_file")) - _files["temp_vrt_file"] = vrt_filename - - # update the execution request object - # setup dynamic model and retrieve the group task needed for tun the async workflow - # create the async task for create the resource into geonode_data with ogr2ogr - if settings.IMPORTER_ENABLE_DYN_MODELS: - ( - dynamic_model, - alternate, - celery_group, - ) = self.setup_dynamic_model( - layer, - execution_id, - should_be_overwritten, - username=_exec.user, - ) - else: - alternate = self.find_alternate_by_dataset(_exec, layer_name, should_be_overwritten) - - layer_names.append(layer_name) - alternates.append(alternate) - - ogr_res = self.get_ogr2ogr_task_group( + _files = files.copy() + vrt_layer_name = None + if has_incompatible_field_names(layer): + vrt_filename, vrt_layer_name = create_vrt_file(layer, files.get("base_file")) + _files["temp_vrt_file"] = vrt_filename + + # update the execution request object + # setup dynamic model and retrieve the group task needed for tun the async workflow + # create the async task for create the resource into geonode_data with ogr2ogr + if settings.IMPORTER_ENABLE_DYN_MODELS: + ( + dynamic_model, + alternate, + celery_group, + ) = self.setup_dynamic_model( + layer, execution_id, - _files, - vrt_layer_name or layer.GetName().lower(), should_be_overwritten, - alternate, + username=_exec.user, + ) + else: + alternate = self.find_alternate_by_dataset(_exec, layer_name, should_be_overwritten) + + layer_names.append(layer_name) + alternates.append(alternate) + + ogr_res = self.get_ogr2ogr_task_group( + execution_id, + _files, + vrt_layer_name or layer.GetName().lower(), + should_be_overwritten, + alternate, + ) + + if settings.IMPORTER_ENABLE_DYN_MODELS: + group_to_call = group( + celery_group.set(link_error=["dynamic_model_error_callback"]), + ogr_res.set(link_error=["dynamic_model_error_callback"]), + ) + else: + group_to_call = group( + ogr_res.set(link_error=["dynamic_model_error_callback"]), ) - if settings.IMPORTER_ENABLE_DYN_MODELS: - group_to_call = group( - celery_group.set(link_error=["dynamic_model_error_callback"]), - ogr_res.set(link_error=["dynamic_model_error_callback"]), - ) - else: - group_to_call = group( - ogr_res.set(link_error=["dynamic_model_error_callback"]), - ) - - # prepare the async chord workflow with the on_success and on_fail methods - workflow = chord(group_to_call)( # noqa - import_next_step.s( - execution_id, - str(self), # passing the handler module path - task_name, - layer_name, - alternate, - **kwargs, - ) + # prepare the async chord workflow with the on_success and on_fail methods + workflow = chord(group_to_call)( # noqa + import_next_step.s( + execution_id, + str(self), # passing the handler module path + task_name, + layer_name, + alternate, + **kwargs, ) + ) except Exception as e: logger.error(e) @@ -582,6 +576,7 @@ def _select_valid_layers(self, all_layers, **kwargs): to extract all the layers. Is possible to pass a filter_layer argument with the name of the layer to retrieve only the needed one + The import step sets filter_existing to exclude existing layers before counting them. """ filter_layer = kwargs.get("filter_layer", None) layers = [] @@ -609,6 +604,26 @@ def _select_valid_layers(self, all_layers, **kwargs): if filter_layer and not layers: logger.warning(f"No layer matching filter '{filter_layer}' was found.") + if layers and kwargs.get("filter_existing"): + _exec = self._get_execution_request_object(kwargs["execution_id"]) + selected_layers = [] + for layer in layers: + layer_name = self.fixup_name(layer.GetName()) + if should_be_imported( + layer_name, + _exec.user, + skip_existing_layer=_exec.input_params.get("skip_existing_layer"), + overwrite_existing_layer=_exec.action == ira.REPLACE.value, + ): + selected_layers.append(layer) + else: + logger.info(f"Skipping existing layer: {layer_name}") + if not selected_layers: + raise ImportException( + "No new layers were detected in your upload. " + "Existing layers were left unchanged, so no updates were made." + ) + layers = selected_layers return layers def can_overwrite(self, _exec_obj, dataset): diff --git a/geonode/upload/handlers/gpkg/handler.py b/geonode/upload/handlers/gpkg/handler.py index 5bd494ce5bf..21b5712038e 100644 --- a/geonode/upload/handlers/gpkg/handler.py +++ b/geonode/upload/handlers/gpkg/handler.py @@ -152,10 +152,9 @@ def handle_xml_file(self, saved_dataset, _exec): pass def _select_valid_layers(self, all_layers, **kwargs): - layers = super()._select_valid_layers(all_layers=all_layers, **kwargs) if execution_id := kwargs.get("execution_id", None): exec_obj = orchestrator.get_execution_object(execution_id) if exec_obj.action in (ira.REPLACE.value, ira.UPSERT.value): - if len(layers) > 1: + if len(super()._select_valid_layers(all_layers=all_layers)) > 1: raise InvalidGeopackageException("For Upsert and Replace, only one layer is allowed in the GPKG") - return layers + return super()._select_valid_layers(all_layers=all_layers, **kwargs) diff --git a/geonode/upload/handlers/gpkg/tests.py b/geonode/upload/handlers/gpkg/tests.py index 90ffc85f12c..f0379472cd1 100644 --- a/geonode/upload/handlers/gpkg/tests.py +++ b/geonode/upload/handlers/gpkg/tests.py @@ -31,6 +31,7 @@ from osgeo import ogr from geonode.upload.celery_tasks import UpdateTaskClass +from unittest.mock import patch class TestGPKGHandler(TestCase): @@ -167,7 +168,8 @@ def test_single_message_error_handler(self): errors = exec_request_obj.output_params.get("errors", []) assert any("exception raised" in str(e) for e in errors) - def test_select_valid_layers(self): + @patch("geonode.upload.handlers.common.vector.should_be_imported", return_value=False) + def test_select_valid_layers(self, should_import): """ The function should return only the datasets with a geometry The other one are discarded @@ -190,12 +192,11 @@ def test_select_valid_layers(self): all_layers = GPKGFileHandler().open_source_file({"base_file": ("/tmp/multiple_layers.gpkg")}) - with self.assertRaises(Exception) as exp: - GPKGFileHandler()._select_valid_layers(all_layers, execution_id=str(exec_id)) + for action in ("replace", "upsert"): + ExecutionRequest.objects.filter(exec_id=exec_id).update(action=action) + with self.assertRaisesMessage(Exception, "For Upsert and Replace, only one layer is allowed in the GPKG"): + GPKGFileHandler()._select_valid_layers(all_layers, execution_id=str(exec_id), filter_existing=True) - self.assertIn( - "For Upsert and Replace, only one layer is allowed in the GPKG", - exp.exception.args[0], - ) + should_import.assert_not_called() os.remove("/tmp/multiple_layers.gpkg") diff --git a/geonode/upload/handlers/shapefile/handler.py b/geonode/upload/handlers/shapefile/handler.py index 009ba91d24f..74f66956b65 100644 --- a/geonode/upload/handlers/shapefile/handler.py +++ b/geonode/upload/handlers/shapefile/handler.py @@ -113,7 +113,7 @@ def extract_params_from_data(_data, action=None): return {"title": data.pop("title"), "store_spatial_file": True}, _data additional_params = { - "skip_existing_layers": _data.pop("skip_existing_layers", "False"), + "skip_existing_layer": _data.pop("skip_existing_layers", False), "resource_pk": _data.pop("resource_pk", None), "store_spatial_file": _data.pop("store_spatial_files", "True"), "action": _data.pop("action", "upload"), diff --git a/geonode/upload/handlers/tiles3d/handler.py b/geonode/upload/handlers/tiles3d/handler.py index 90720933891..4b6f24fa37d 100755 --- a/geonode/upload/handlers/tiles3d/handler.py +++ b/geonode/upload/handlers/tiles3d/handler.py @@ -30,11 +30,12 @@ from geonode.upload.orchestrator import orchestrator from geonode.upload.celery_tasks import import_orchestrator from geonode.upload.handlers.common.vector import BaseVectorFileHandler -from geonode.upload.handlers.utils import create_alternate, should_be_imported +from geonode.upload.handlers.utils import create_alternate from geonode.upload.utils import ImporterRequestAction as ira from geonode.base.models import ResourceBase from geonode.upload.handlers.tiles3d.exceptions import Invalid3DTilesException from geonode.resource.registry import resource_manager_registry +from geonode.upload.api.exceptions import ImportException logger = logging.getLogger("importer") @@ -191,7 +192,7 @@ def extract_params_from_data(_data, action=None): return {"title": data.pop("title"), "store_spatial_file": True}, _data return { - "skip_existing_layers": _data.pop("skip_existing_layers", "False"), + "skip_existing_layer": _data.pop("skip_existing_layers", False), "store_spatial_file": _data.pop("store_spatial_files", "True"), "action": _data.pop("action", "upload"), "original_zip_name": _data.pop("original_zip_name", None), @@ -209,18 +210,16 @@ def import_resource(self, files: dict, execution_id: str, **kwargs) -> str: # start looping on the layers available layer_name = self.fixup_name(filename) should_be_overwritten = _exec.action == ira.REPLACE.value - # should_be_imported check if the user+layername already exists or not - if should_be_imported( - layer_name, - _exec.user, - skip_existing_layer=_exec.input_params.get("skip_existing_layer"), - overwrite_existing_layer=should_be_overwritten, + # Check for name collisions with any resource owned by the user + user_datasets = ResourceBase.objects.filter(owner=_exec.user, alternate=layer_name) + + dataset_exists = user_datasets.exists() + if not ( + dataset_exists + and _exec.input_params.get("skip_existing_layer") + and user_datasets.filter(resource_type="dataset", subtype="3dtiles").exists() ): - user_datasets = ResourceBase.objects.filter(owner=_exec.user, alternate=layer_name) - - dataset_exists = user_datasets.exists() - if dataset_exists and should_be_overwritten: layer_name, alternate = ( layer_name, @@ -230,6 +229,11 @@ def import_resource(self, files: dict, execution_id: str, **kwargs) -> str: alternate = layer_name else: alternate = create_alternate(layer_name, execution_id) + else: + raise ImportException( + "No new layers were detected in your upload. " + "Existing layers were left unchanged, so no updates were made." + ) import_orchestrator.apply_async( ( diff --git a/geonode/upload/handlers/tiles3d/tests.py b/geonode/upload/handlers/tiles3d/tests.py index ab6a5a0088d..5b553d1d1bb 100755 --- a/geonode/upload/handlers/tiles3d/tests.py +++ b/geonode/upload/handlers/tiles3d/tests.py @@ -19,6 +19,7 @@ import json import os import shutil +from unittest.mock import patch from django.test import TestCase from geonode.upload.handlers.tiles3d.exceptions import Invalid3DTilesException from geonode.upload.handlers.tiles3d.handler import Tiles3DFileHandler @@ -26,10 +27,11 @@ from geonode.upload import project_dir from geonode.upload.orchestrator import orchestrator from geonode.upload.models import UploadParallelismLimit -from geonode.upload.api.exceptions import UploadParallelismLimitException +from geonode.upload.api.exceptions import UploadParallelismLimitException, ImportException from geonode.base.populate_test_data import create_single_dataset from osgeo import ogr from geonode.assets.handlers import asset_handler_registry +from geonode.base.models import ResourceBase class TestTiles3DFileHandler(TestCase): @@ -69,6 +71,101 @@ def test_task_list_is_the_expected_one_copy(self): self.assertEqual(len(self.handler.TASKS["copy"]), 2) self.assertTupleEqual(expected, self.handler.TASKS["copy"]) + @patch("geonode.upload.handlers.tiles3d.handler.import_orchestrator.apply_async") + def test_import_resource_should_skip_existing_layer(self, import_orchestrator): + ResourceBase.objects.create( + title="valid_3dtiles", + alternate="valid_3dtiles", + owner=self.owner, + resource_type="dataset", + subtype="3dtiles", + ) + exec_id = orchestrator.create_execution_request( + user=self.owner, + func_name="funct1", + step="step", + input_params={"files": self.valid_files, "skip_existing_layer": True}, + ) + + with self.assertRaisesMessage( + ImportException, + "No new layers were detected in your upload. " + "Existing layers were left unchanged, so no updates were made.", + ): + self.handler.import_resource(files=self.valid_files, execution_id=str(exec_id)) + + import_orchestrator.assert_not_called() + + @patch("geonode.upload.handlers.tiles3d.handler.import_orchestrator.apply_async") + def test_import_resource_should_import_when_layer_is_not_skipped(self, import_orchestrator): + other_owner, _ = get_user_model().objects.get_or_create(username="other-3dtiles-owner") + cases = ( + (False, self.owner, False), + (None, self.owner, False), + (True, other_owner, False), + (True, self.owner, True), + ) + for index, (flag, owner, vector) in enumerate(cases): + with self.subTest(flag=flag, owner=owner, vector=vector): + layer_name = f"tiles_{index}" + if vector: + create_single_dataset(name=layer_name, owner=owner) + else: + ResourceBase.objects.create( + title=layer_name, + alternate=layer_name, + owner=owner, + resource_type="dataset", + subtype="3dtiles", + ) + params = {"original_zip_name": layer_name} + if flag is not None: + params["skip_existing_layer"] = flag + exec_id = orchestrator.create_execution_request( + user=self.owner, + func_name="funct1", + step="step", + input_params=params, + ) + result = self.handler.import_resource(files=self.valid_files, execution_id=str(exec_id)) + self.assertEqual(layer_name, result[0]) + if flag: + self.assertEqual(layer_name, result[1]) + else: + self.assertNotEqual(layer_name, result[1]) + import_orchestrator.assert_called_once() + import_orchestrator.reset_mock() + + @patch("geonode.upload.handlers.tiles3d.handler.import_orchestrator.apply_async") + def test_import_resource_should_not_skip_non_3d_resource(self, import_orchestrator): + for resource_type, subtype in (("document", "pdf"), ("dataset", "vector")): + with self.subTest(resource_type=resource_type, subtype=subtype): + layer_name = f"non_3d_{resource_type}" + resource = ResourceBase.objects.create( + title=layer_name, + alternate=layer_name, + owner=self.owner, + resource_type=resource_type, + subtype=subtype, + ) + original = ResourceBase.objects.filter(pk=resource.pk).values().get() + exec_id = orchestrator.create_execution_request( + user=self.owner, + func_name="funct1", + step="step", + action="upload", + input_params={"original_zip_name": layer_name, "skip_existing_layer": True}, + ) + + result = self.handler.import_resource(files=self.valid_files, execution_id=str(exec_id)) + + self.assertEqual(layer_name, result[0]) + self.assertNotEqual(layer_name, result[1]) + import_orchestrator.assert_called_once() + self.assertEqual(result[1], import_orchestrator.call_args.args[0][5]) + self.assertEqual(original, ResourceBase.objects.filter(pk=resource.pk).values().get()) + import_orchestrator.reset_mock() + def test_is_valid_should_raise_exception_if_the_parallelism_is_met(self): parallelism, created = UploadParallelismLimit.objects.get_or_create(slug="default_max_parallel_uploads") old_value = parallelism.max_number diff --git a/geonode/upload/tests/unit/test_dastore.py b/geonode/upload/tests/unit/test_dastore.py index fd9b504c3cf..5216abe2677 100644 --- a/geonode/upload/tests/unit/test_dastore.py +++ b/geonode/upload/tests/unit/test_dastore.py @@ -16,12 +16,27 @@ # along with this program. If not, see . # ######################################################################### -from django.test import TestCase +from django.test import TestCase, override_settings from geonode.upload import project_dir from geonode.upload.orchestrator import orchestrator from geonode.upload.datastore import DataStoreManager from django.contrib.auth import get_user_model +from unittest.mock import patch + +from geonode.base.populate_test_data import create_single_dataset +from geonode.resource.models import ExecutionRequest +from geonode.upload.models import ResourceHandlerInfo +from geonode.upload.utils import UploadLimitValidator +from geonode.upload.api.exceptions import ImportException +from geonode.upload.celery_tasks import UpdateTaskClass +from geonode.resource.api.serializer import ExecutionRequestSerializer +from geonode.upload.handlers.common.vector import BaseVectorFileHandler +from geonode.upload.handlers.common.raster import BaseRasterFileHandler +from geonode.upload.handlers.shapefile.handler import ShapeFileHandler +from geonode.upload.handlers.tiles3d.handler import Tiles3DFileHandler +from geonode.base.models import ResourceBase + class TestDataStoreManager(TestCase): """ """ @@ -30,6 +45,7 @@ class TestDataStoreManager(TestCase): def setUpClass(cls): super().setUpClass() cls.files = {"base_file": f"{project_dir}/tests/fixture/valid.gpkg"} + cls.raster_files = {"base_file": f"{project_dir}/tests/fixture/test_raster.tif"} def setUp(self): self.user = get_user_model().objects.first() @@ -58,8 +74,151 @@ def setUp(self): ) self.gpkg_path = f"{project_dir}/tests/fixture/valid.gpkg" + def _create_datastore(self, files, handler_module_path, **input_params): + execution_id = orchestrator.create_execution_request( + user=self.user, + func_name="create", + step="geonode.upload.import_resource", + action="upload", + input_params={ + "files": files, + "handler_module_path": handler_module_path, + **input_params, + }, + ) + datastore = DataStoreManager(files, handler_module_path, self.user, execution_id) + return execution_id, datastore + + def _assert_all_skipped(self, execution_id, datastore, resource): + message = ( + "No new layers were detected in your upload. " + "Existing layers were left unchanged, so no updates were made." + ) + original = type(resource).objects.filter(pk=resource.pk).values().get() + self.assertEqual(1, UploadLimitValidator(self.user)._get_parallel_uploads_count()) + task = UpdateTaskClass() + task.name = "geonode.upload.import_resource" + task.bulk = True + args = (str(execution_id), str(datastore.handler()), "upload") + task.before_start("test-skip", args, {}) + with self.assertRaisesMessage(ImportException, message) as error: + datastore._import_and_register(str(execution_id), task.name) + task.on_failure(error.exception, "test-skip", args, {}, None) + execution = ExecutionRequest.objects.get(exec_id=execution_id) + self.assertEqual(ExecutionRequest.STATUS_FAILED, execution.status) + self.assertIsNotNone(execution.finished) + self.assertEqual(message, execution.log) + self.assertEqual(message, ExecutionRequestSerializer(execution).data["log"]) + self.assertEqual(0, UploadLimitValidator(self.user)._get_parallel_uploads_count()) + self.assertFalse(ResourceHandlerInfo.objects.filter(execution_request=execution).exists()) + self.assertEqual(original, type(resource).objects.filter(pk=resource.pk).values().get()) + + def test_handlers_normalize_skip_existing_layers(self): + for handler in (BaseVectorFileHandler, BaseRasterFileHandler, ShapeFileHandler, Tiles3DFileHandler): + for flag in (True, False, None): + with self.subTest(handler=handler.__name__, flag=flag): + data = {} if flag is None else {"skip_existing_layers": flag} + params, files = handler.extract_params_from_data(data) + self.assertIs(params["skip_existing_layer"], flag is True) + self.assertNotIn("skip_existing_layers", params) + self.assertNotIn("skip_existing_layers", files) + def test_input_is_valid_with_files(self): self.assertTrue(self.datastore.input_is_valid()) def test_input_is_valid_with_urls(self): self.assertTrue(self.datastore_url.input_is_valid()) + + @patch("geonode.upload.handlers.common.raster.import_orchestrator.apply_async") + def test_skip_existing_raster_fails_execution_and_releases_parallel_slot(self, schedule_import): + resource = create_single_dataset(name="test_raster", owner=self.user) + handler_module_path = "geonode.upload.handlers.common.raster.BaseRasterFileHandler" + execution_id, datastore = self._create_datastore( + self.raster_files, + handler_module_path, + skip_existing_layer=True, + total_layers=1, + ) + + self._assert_all_skipped(execution_id, datastore, resource) + + schedule_import.assert_not_called() + + @patch("geonode.upload.handlers.common.vector.chord") + def test_skip_existing_geojson_fails_execution(self, celery_chord): + files = {"base_file": f"{project_dir}/tests/fixture/valid.geojson"} + resource = create_single_dataset(name="valid", owner=self.user) + handler_module_path = "geonode.upload.handlers.geojson.handler.GeoJsonFileHandler" + execution_id, datastore = self._create_datastore( + files, + handler_module_path, + skip_existing_layer=True, + total_layers=1, + ) + + self._assert_all_skipped(execution_id, datastore, resource) + + celery_chord.assert_not_called() + + @patch("geonode.upload.handlers.tiles3d.handler.import_orchestrator.apply_async") + def test_skip_existing_3dtiles_fails_execution(self, schedule_import): + resource = ResourceBase.objects.create( + title="valid_3dtiles", + alternate="valid_3dtiles", + owner=self.user, + resource_type="dataset", + subtype="3dtiles", + ) + execution_id, datastore = self._create_datastore( + {"base_file": f"{project_dir}/tests/fixture/3dtilesample/tileset.json"}, + "geonode.upload.handlers.tiles3d.handler.Tiles3DFileHandler", + skip_existing_layer=True, + original_zip_name="valid_3dtiles", + total_layers=1, + ) + self._assert_all_skipped(execution_id, datastore, resource) + schedule_import.assert_not_called() + + @override_settings(IMPORTER_ENABLE_DYN_MODELS=False) + @patch("geonode.upload.handlers.common.vector.chord") + def test_partial_skip_uses_imported_layer_count_for_progress(self, celery_chord): + handler_module_path = "geonode.upload.handlers.gpkg.handler.GPKGFileHandler" + execution_id, datastore = self._create_datastore( + {"base_file": f"{project_dir}/tests/fixture/multiple_layers.gpkg"}, + handler_module_path, + skip_existing_layer=True, + ) + + handler = datastore.handler() + layers = handler._select_valid_layers(handler.open_source_file(datastore.files)) + names = [handler.fixup_name(layer.GetName()) for layer in layers] + self.assertGreater(len(names), 1) + create_single_dataset(name=names[0], owner=self.user) + datastore._import_and_register(str(execution_id), "geonode.upload.import_resource") + + execution = ExecutionRequest.objects.get(exec_id=execution_id) + self.assertEqual(len(names) - 1, execution.input_params["total_layers"]) + self.assertEqual(len(names) - 1, celery_chord.call_count) + scheduled_names = [call.args[0].args[3] for call in celery_chord.return_value.call_args_list] + self.assertEqual(names[1:], scheduled_names) + + for name in names[1:]: + resource = create_single_dataset(name=name, owner=self.user) + ResourceHandlerInfo.objects.create( + execution_request=execution, + handler_module_path=handler_module_path, + resource=resource, + ) + + tasks = execution.tasks + for status_by_task in tasks.values(): + for task_name in status_by_task: + status_by_task[task_name] = "SUCCESS" + execution.tasks = tasks + execution.save(update_fields=["tasks"]) + + orchestrator.evaluate_execution_progress(str(execution_id), handler_module_path=handler_module_path) + + execution.refresh_from_db() + self.assertEqual(ExecutionRequest.STATUS_FINISHED, execution.status) + self.assertIsNotNone(execution.finished)