Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions geonode/upload/api/tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down
7 changes: 5 additions & 2 deletions geonode/upload/handlers/common/raster.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand Down Expand Up @@ -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
Expand Down
25 changes: 25 additions & 0 deletions geonode/upload/handlers/common/tests_raster.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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:
Expand Down
76 changes: 76 additions & 0 deletions geonode/upload/handlers/common/tests_vector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
"""
Expand Down
141 changes: 78 additions & 63 deletions geonode/upload/handlers/common/vector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"),
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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 = []
Expand Down Expand Up @@ -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):
Expand Down
5 changes: 2 additions & 3 deletions geonode/upload/handlers/gpkg/handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Loading