From b61142b480417668379d36305f8fd93de56ed405 Mon Sep 17 00:00:00 2001 From: bramjanssen Date: Fri, 14 Aug 2026 14:42:32 +0200 Subject: [PATCH] fix(#45): added new status to represent deleted jobs --- .../versions/c30191fce04f_added_new_status.py | 59 ++++++++++++++++++ app/platforms/implementations/openeo.py | 6 +- app/schemas/enum.py | 1 + app/services/processing.py | 10 ++-- tests/platforms/test_openeo_platform.py | 40 ++++++++++++- tests/services/test_processing.py | 60 +++++++++++++++++++ 6 files changed, 169 insertions(+), 7 deletions(-) create mode 100644 alembic/versions/c30191fce04f_added_new_status.py diff --git a/alembic/versions/c30191fce04f_added_new_status.py b/alembic/versions/c30191fce04f_added_new_status.py new file mode 100644 index 0000000..03ad977 --- /dev/null +++ b/alembic/versions/c30191fce04f_added_new_status.py @@ -0,0 +1,59 @@ +"""added new status + +Revision ID: c30191fce04f +Revises: 833e4a41c2ad +Create Date: 2026-08-14 14:10:04.216822 + +""" +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = 'c30191fce04f' +down_revision: Union[str, Sequence[str], None] = '833e4a41c2ad' +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + +OLD_STATUS_ENUM = sa.Enum( + 'CREATED', 'QUEUED', 'RUNNING', 'FINISHED', 'CANCELED', 'FAILED', 'UNKNOWN', + name='processingstatusenum', +) +NEW_STATUS_ENUM = sa.Enum( + 'CREATED', 'QUEUED', 'RUNNING', 'FINISHED', 'CANCELED', 'FAILED', 'DELETED', 'UNKNOWN', + name='processingstatusenum', +) + + +def upgrade() -> None: + """Upgrade schema.""" + op.alter_column( + 'processing_jobs', 'status', + existing_type=OLD_STATUS_ENUM, + type_=NEW_STATUS_ENUM, + existing_nullable=False, + ) + op.alter_column( + 'upscaling_tasks', 'status', + existing_type=OLD_STATUS_ENUM, + type_=NEW_STATUS_ENUM, + existing_nullable=False, + ) + + +def downgrade() -> None: + """Downgrade schema.""" + op.alter_column( + 'upscaling_tasks', 'status', + existing_type=NEW_STATUS_ENUM, + type_=OLD_STATUS_ENUM, + existing_nullable=False, + ) + op.alter_column( + 'processing_jobs', 'status', + existing_type=NEW_STATUS_ENUM, + type_=OLD_STATUS_ENUM, + existing_nullable=False, + ) diff --git a/app/platforms/implementations/openeo.py b/app/platforms/implementations/openeo.py index 9822df6..fa8d25b 100644 --- a/app/platforms/implementations/openeo.py +++ b/app/platforms/implementations/openeo.py @@ -340,7 +340,7 @@ def _map_openeo_status(self, status: str) -> ProcessingStatusEnum: "created": ProcessingStatusEnum.CREATED, "queued": ProcessingStatusEnum.QUEUED, "running": ProcessingStatusEnum.RUNNING, - "cancelled": ProcessingStatusEnum.CANCELED, + "canceled": ProcessingStatusEnum.CANCELED, "finished": ProcessingStatusEnum.FINISHED, "error": ProcessingStatusEnum.FAILED, } @@ -368,6 +368,10 @@ async def get_job_status( f"job {job_id} after refresh: {retry_error}" ) return ProcessingStatusEnum.UNKNOWN + if e.http_status_code == 404 and e.code == "JobNotFound": + logger.warning(f"Job {job_id} not found on OpenEO backend.") + return ProcessingStatusEnum.DELETED + logger.error( f"Error occurred while fetching job status for job {job_id}: {e}" ) diff --git a/app/schemas/enum.py b/app/schemas/enum.py index e1d329c..fceec4e 100644 --- a/app/schemas/enum.py +++ b/app/schemas/enum.py @@ -13,6 +13,7 @@ class ProcessingStatusEnum(str, Enum): FINISHED = "finished" CANCELED = "canceled" FAILED = "failed" + DELETED = "deleted" UNKNOWN = "unknown" diff --git a/app/services/processing.py b/app/services/processing.py index 11c3a53..39b4f87 100644 --- a/app/services/processing.py +++ b/app/services/processing.py @@ -3,6 +3,9 @@ from fastapi import Response from loguru import logger +from sqlalchemy.orm import Session +from stac_pydantic import Collection + from app.auth import get_current_user_id from app.database.models.processing_job import ( ProcessingJobRecord, @@ -13,10 +16,8 @@ update_job_status_by_id, ) from app.platforms.dispatcher import get_processing_platform -from sqlalchemy.orm import Session - from app.schemas.enum import ProcessingStatusEnum -from app.schemas.parameters import ParamRequest, Parameter +from app.schemas.parameters import Parameter, ParamRequest from app.schemas.unit_job import ( BaseJobRequest, ProcessingJob, @@ -24,9 +25,8 @@ ServiceDetails, ) -from stac_pydantic import Collection - INACTIVE_JOB_STATUSES = { + ProcessingStatusEnum.DELETED, ProcessingStatusEnum.CANCELED, ProcessingStatusEnum.FAILED, ProcessingStatusEnum.FINISHED, diff --git a/tests/platforms/test_openeo_platform.py b/tests/platforms/test_openeo_platform.py index c9f66cc..9fc9f31 100644 --- a/tests/platforms/test_openeo_platform.py +++ b/tests/platforms/test_openeo_platform.py @@ -188,7 +188,7 @@ async def test_execute_job_process_openeo_error( ("created", ProcessingStatusEnum.CREATED), ("queued", ProcessingStatusEnum.QUEUED), ("running", ProcessingStatusEnum.RUNNING), - ("cancelled", ProcessingStatusEnum.CANCELED), + ("canceled", ProcessingStatusEnum.CANCELED), ("finished", ProcessingStatusEnum.FINISHED), ("error", ProcessingStatusEnum.FAILED), ("CrEaTeD", ProcessingStatusEnum.CREATED), # Case insensitivity @@ -297,6 +297,44 @@ async def test_get_job_status_non_auth_openeo_error_returns_unknown( assert result == ProcessingStatusEnum.UNKNOWN +@pytest.mark.asyncio +@patch.object(OpenEOPlatform, "_setup_connection", new_callable=AsyncMock) +async def test_get_job_status_returns_deleted_when_job_not_found( + mock_setup_connection, platform +): + job = MagicMock() + job.status.side_effect = OpenEoApiError( + message="not found", code="JobNotFound", http_status_code=404 + ) + connection = MagicMock() + connection.job.return_value = job + mock_setup_connection.return_value = connection + + details = ServiceDetails(endpoint="foo", application="bar") + result = await platform.get_job_status("foobar", "job123", details) + + assert result == ProcessingStatusEnum.DELETED + + +@pytest.mark.asyncio +@patch.object(OpenEOPlatform, "_setup_connection", new_callable=AsyncMock) +async def test_get_job_status_404_with_other_code_returns_unknown( + mock_setup_connection, platform +): + job = MagicMock() + job.status.side_effect = OpenEoApiError( + message="not found", code="SomethingElse", http_status_code=404 + ) + connection = MagicMock() + connection.job.return_value = job + mock_setup_connection.return_value = connection + + details = ServiceDetails(endpoint="foo", application="bar") + result = await platform.get_job_status("foobar", "job123", details) + + assert result == ProcessingStatusEnum.UNKNOWN + + @pytest.mark.asyncio @patch.object(OpenEOPlatform, "_setup_connection") async def test_get_job_results_success(mock_connection, platform, fake_result): diff --git a/tests/services/test_processing.py b/tests/services/test_processing.py index 81815db..45f1ad7 100644 --- a/tests/services/test_processing.py +++ b/tests/services/test_processing.py @@ -469,6 +469,66 @@ async def test_get_processing_job_by_user_id_inactive_status( assert result.parameters == {"param1": "value1"} +@pytest.mark.asyncio +@patch("app.services.processing._refresh_job_status") +@patch("app.services.processing.get_job_by_user_id") +@patch("app.services.processing.get_current_user_id") +async def test_get_processing_job_by_user_id_deleted_status_skips_refresh( + mock_current_user, mock_get_job, mock_refresh_status, fake_db_session +): + + fake_service_details = { + "endpoint": "https://openeofed.dataspace.copernicus.eu", + "application": "https://raw.githubusercontent.com/ESA-APEx/apex_algorithms/" + "32ea3c9a6fa24fe063cb59164cd318cceb7209b0/openeo_udp/variabilitymap/" + "variabilitymap.json", + } + fake_result = make_job_record(ProcessingStatusEnum.DELETED, fake_service_details) + mock_get_job.return_value = fake_result + + mock_current_user.return_value = "foobar" + + result = await get_processing_job_by_user_id("foobar-token", fake_db_session, 1) + + mock_get_job.assert_called_once_with(fake_db_session, 1, "foobar") + mock_refresh_status.assert_not_called() + assert isinstance(result, ProcessingJob) + assert result.status == ProcessingStatusEnum.DELETED + + +@pytest.mark.asyncio +@patch("app.services.processing.update_job_status_by_id") +@patch("app.services.processing.get_job_status") +@patch("app.services.processing.get_jobs_by_user_id") +@patch("app.services.processing.get_current_user_id") +async def test_get_processing_jobs_skips_refresh_for_deleted_status( + mock_current_user, + mock_get_jobs, + mock_get_job_status, + mock_update_job_status, + fake_db_session, +): + deleted_job = ProcessingJobRecord( + id=4, + platform_job_id="platform789", + label=ProcessTypeEnum.OPENEO, + title="Deleted Job", + status=ProcessingStatusEnum.DELETED, + parameters="{}", + service=json.dumps({"application": "foo", "endpoint": "bar"}), + ) + mock_get_jobs.return_value = [deleted_job] + + mock_current_user.return_value = "foobar" + + results = await get_processing_jobs_by_user_id("foobar-token", fake_db_session) + + assert len(results) == 1 + assert results[0].status == ProcessingStatusEnum.DELETED + mock_get_job_status.assert_not_called() + mock_update_job_status.assert_not_called() + + @pytest.mark.asyncio @patch("app.services.processing.get_job_by_user_id") @patch("app.services.processing.get_current_user_id")