Skip to content
Merged
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
59 changes: 59 additions & 0 deletions alembic/versions/c30191fce04f_added_new_status.py
Original file line number Diff line number Diff line change
@@ -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,
)
6 changes: 5 additions & 1 deletion app/platforms/implementations/openeo.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
Expand Down Expand Up @@ -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}"
)
Expand Down
1 change: 1 addition & 0 deletions app/schemas/enum.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ class ProcessingStatusEnum(str, Enum):
FINISHED = "finished"
CANCELED = "canceled"
FAILED = "failed"
DELETED = "deleted"
UNKNOWN = "unknown"


Expand Down
10 changes: 5 additions & 5 deletions app/services/processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -13,20 +16,17 @@
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,
ProcessingJobSummary,
ServiceDetails,
)

from stac_pydantic import Collection

INACTIVE_JOB_STATUSES = {
ProcessingStatusEnum.DELETED,
ProcessingStatusEnum.CANCELED,
ProcessingStatusEnum.FAILED,
ProcessingStatusEnum.FINISHED,
Expand Down
40 changes: 39 additions & 1 deletion tests/platforms/test_openeo_platform.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down
60 changes: 60 additions & 0 deletions tests/services/test_processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Loading