From d9fe79b83bb98c910f5cb6e7c56ca251e451e77a Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Mon, 13 Jul 2026 15:03:41 +0200 Subject: [PATCH 1/6] New API version v2_1_0 with POST /workflows/discover --- CHANGELOG.md | 5 + pyproject.toml | 2 +- src/ewoksserver/app/lifespan.py | 15 +- .../app/routes/common/discovery.py | 165 ++++++++++++++++++ .../app/routes/execution/__init__.py | 2 + .../app/routes/execution/router.py | 6 +- src/ewoksserver/app/routes/icons/__init__.py | 1 + src/ewoksserver/app/routes/tasks/__init__.py | 1 + src/ewoksserver/app/routes/tasks/discovery.py | 90 ---------- src/ewoksserver/app/routes/tasks/router.py | 2 +- .../app/routes/workflows/__init__.py | 10 +- .../app/routes/workflows/models.py | 10 ++ .../app/routes/workflows/router.py | 58 +++++- src/ewoksserver/tests/api_versions.py | 5 +- 14 files changed, 261 insertions(+), 111 deletions(-) create mode 100644 src/ewoksserver/app/routes/common/discovery.py delete mode 100644 src/ewoksserver/app/routes/tasks/discovery.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 15217f7..b3ef961 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - In test, use httpx2 instead of httpx (Deprecation Warning from starlette dependency) +### Added + +- New API version `v2_1_0`: +- New endpoint `POST /api/workflows/discover` to discover ewoks workflows from python packages. + ## [2.1.2] - 2026-03-06 ### Changed diff --git a/pyproject.toml b/pyproject.toml index 12a4d9d..77fa1d9 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -19,7 +19,7 @@ dependencies = [ "fastapi", "uvicorn[standard]", "python-socketio", - "ewoksjob[worker] >=1.1,<2", + "ewoksjob[worker] >=1.6.0rc1,<2", "ewokscore >=1.0.0", "pydantic-settings", "packaging", diff --git a/src/ewoksserver/app/lifespan.py b/src/ewoksserver/app/lifespan.py index 1326aa5..5982f22 100644 --- a/src/ewoksserver/app/lifespan.py +++ b/src/ewoksserver/app/lifespan.py @@ -15,8 +15,8 @@ from .. import resources from . import config from .backends import json_backend +from .routes.common import discovery from .routes.execution import socketio -from .routes.tasks.discovery import discover_tasks logger = logging.getLogger(__name__) @@ -31,7 +31,7 @@ async def fastapi_lifespan(app: FastAPI) -> Generator[None, None, None]: _copy_default_resources(ewoks_settings) _enable_execution_events(ewoks_settings) with _enable_execution(ewoks_settings): - _rediscover_tasks(ewoks_settings) + _rediscover_resources(ewoks_settings) _print_ewoks_settings(ewoks_settings) yield @@ -64,11 +64,12 @@ def _copy_default_resources(ewoks_settings: config.EwoksSettings) -> None: shutil.copy(src, dest) -def _rediscover_tasks(ewoks_settings: config.EwoksSettings) -> None: +def _rediscover_resources(ewoks_settings: config.EwoksSettings) -> None: if not ewoks_settings.ewoks_discovery.on_start_up: return + try: - tasks = discover_tasks(ewoks_settings) + tasks = discovery.discover_tasks(ewoks_settings) except Exception as ex: tasks = [] logger.exception("Task discovery failed: %s", ex) @@ -76,6 +77,12 @@ def _rediscover_tasks(ewoks_settings: config.EwoksSettings) -> None: for resource in tasks: json_backend.save_resource(root_url, resource["task_identifier"], resource) + try: + discovery.discover_workflows(ewoks_settings) + except Exception as ex: + logger.exception("Workflow discovery failed: %s", ex) + logger.warning("Discovered workflows not used yet") + def _enable_execution_events(ewoks_settings: config.EwoksSettings) -> None: """Set default ewoks event handler when nothing has been configured""" diff --git a/src/ewoksserver/app/routes/common/discovery.py b/src/ewoksserver/app/routes/common/discovery.py new file mode 100644 index 0000000..5a851d5 --- /dev/null +++ b/src/ewoksserver/app/routes/common/discovery.py @@ -0,0 +1,165 @@ +import logging + +from ewoksjob.client import discover_all_tasks +from ewoksjob.client import discover_all_workflows +from ewoksjob.client import discover_tasks_from_modules +from ewoksjob.client import discover_workflows_from_modules +from ewoksjob.client import get_queues +from ewoksjob.client.local import discover_all_tasks as discover_all_tasks_local +from ewoksjob.client.local import discover_all_workflows as discover_all_workflows_local +from ewoksjob.client.local import ( + discover_tasks_from_modules as discover_tasks_from_modules_local, +) +from ewoksjob.client.local import ( + discover_workflows_from_modules as discover_workflows_from_modules_local, +) + +from ...config import EwoksSettings +from ...models import EwoksSchedulingType + +logger = logging.getLogger(__name__) + + +def discover_tasks( + settings: EwoksSettings, + modules: list[str] | None = None, + reload: bool | None = None, + task_type: str | None = None, + worker_options: dict | None = None, +) -> list[dict[str, str]]: + """ + :raises ModuleNotFoundError: failed importing tasks. + :raises TimeoutError: timeout when asking a remote worker for tasks. + :raises Exception: any other import or remote error. + """ + if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: + if modules: + discover = discover_tasks_from_modules_local + else: + discover = discover_all_tasks_local + else: + if modules: + discover = discover_tasks_from_modules + else: + discover = discover_all_tasks + + discover_kwargs = dict() + if reload is not None: + discover_kwargs["reload"] = reload + if task_type is not None: + discover_kwargs["task_type"] = task_type + + tasks = _discover( + discover, + settings, + modules=modules, + discover_kwargs=discover_kwargs, + worker_options=worker_options, + key=lambda task: task["task_identifier"], + ) + + for task in tasks: + _set_default_task_properties(task) + return tasks + + +def discover_workflows( + settings: EwoksSettings, + modules: list[str] | None = None, + workflow_extension: str | None = None, + worker_options: dict | None = None, +) -> list[str]: + """ + :raises ModuleNotFoundError: failed importing workflows. + :raises TimeoutError: timeout when asking a remote worker for workflows. + :raises Exception: any other import or remote error. + """ + if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: + if modules: + discover = discover_workflows_from_modules_local + else: + discover = discover_all_workflows_local + else: + if modules: + discover = discover_workflows_from_modules + else: + discover = discover_all_workflows + + discover_kwargs = dict() + if workflow_extension is not None: + discover_kwargs["workflow_extension"] = workflow_extension + + return _discover( + discover, + settings, + modules=modules, + discover_kwargs=discover_kwargs, + worker_options=worker_options, + key=lambda workflow: workflow, + ) + + +def _discover( + discover, + settings: EwoksSettings, + modules: list[str] | None, + discover_kwargs: dict, + worker_options: dict | None, + key, +) -> list: + """ + :raises ModuleNotFoundError: failed importing tasks or workflows. + :raises TimeoutError: timeout when asking a remote worker. + :raises Exception: any other import or remote error. + """ + if worker_options is None: + kwargs = dict() + else: + kwargs = dict(worker_options) + + # Discovery: position arguments + if modules: + kwargs["args"] = modules + + # Discovery: named arguments + kwargs["kwargs"] = discover_kwargs + + timeout = settings.ewoks_discovery.timeout + if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: + return _discover_locally(discover, kwargs, timeout=timeout) + else: + return _discover_in_all_queues(discover, kwargs, key, timeout=timeout) + + +def _discover_locally(discover, kwargs: dict, timeout: float | None = None) -> list: + return discover(**kwargs).result(timeout=timeout) + + +def _discover_in_all_queues( + discover, kwargs: dict, key, timeout: float | None = None +) -> list: + futures = [discover(**kwargs, queue=queue) for queue in get_queues()] + + # Store items in a dict to avoid duplicates + item_dict = {} + for future in futures: + # Ignore failures of a single queue to not prevent discovery on other queues + new_items = future.result(timeout=timeout) + exc = future.exception() + if exc: + logger.warning(f"Discovery failed for {future.queue}: {exc}") + continue + if new_items is None: + continue + for item in new_items: + item_dict[key(item)] = item + return list(item_dict.values()) + + +def _set_default_task_properties(task: dict) -> None: + if not task.get("icon"): + task["icon"] = "default.png" + if not task.get("label"): + task_identifier = task.get("task_identifier") + if task_identifier: + task["label"] = task_identifier.split(".")[-1] diff --git a/src/ewoksserver/app/routes/execution/__init__.py b/src/ewoksserver/app/routes/execution/__init__.py index 3acba71..0c9d861 100644 --- a/src/ewoksserver/app/routes/execution/__init__.py +++ b/src/ewoksserver/app/routes/execution/__init__.py @@ -7,10 +7,12 @@ (1, 0, 0): _v1_0_0_router, (1, 1, 0): _v1_1_0_router, (2, 0, 0): _v2_0_0_router, + (2, 1, 0): _v2_0_0_router, } app_creators = { (1, 0, 0): _create_socketio_app, (1, 1, 0): _create_socketio_app, (2, 0, 0): _create_socketio_app, + (2, 1, 0): _create_socketio_app, } diff --git a/src/ewoksserver/app/routes/execution/router.py b/src/ewoksserver/app/routes/execution/router.py index ca03546..69b6ed8 100644 --- a/src/ewoksserver/app/routes/execution/router.py +++ b/src/ewoksserver/app/routes/execution/router.py @@ -23,8 +23,6 @@ from .utils import submit_workflow v1_0_0_router = APIRouter() -v1_1_0_router = APIRouter() -v2_0_0_router = APIRouter() @v1_0_0_router.post( @@ -123,6 +121,7 @@ def execute_events_v1( return {"jobs": list(jobs.values())} +v1_1_0_router = APIRouter() v1_1_0_router.include_router(v1_0_0_router) @@ -140,6 +139,9 @@ def workers(settings: EwoksSettingsType) -> dict[str, list[str] | None]: return {"workers": get_queues()} +v2_0_0_router = APIRouter() + + @v2_0_0_router.post( "/execute/{identifier}", summary="Execute workflow", diff --git a/src/ewoksserver/app/routes/icons/__init__.py b/src/ewoksserver/app/routes/icons/__init__.py index a25287a..31edf27 100644 --- a/src/ewoksserver/app/routes/icons/__init__.py +++ b/src/ewoksserver/app/routes/icons/__init__.py @@ -4,4 +4,5 @@ (1, 0, 0): _router, (1, 1, 0): _router, (2, 0, 0): _router, + (2, 1, 0): _router, } diff --git a/src/ewoksserver/app/routes/tasks/__init__.py b/src/ewoksserver/app/routes/tasks/__init__.py index a25287a..31edf27 100644 --- a/src/ewoksserver/app/routes/tasks/__init__.py +++ b/src/ewoksserver/app/routes/tasks/__init__.py @@ -4,4 +4,5 @@ (1, 0, 0): _router, (1, 1, 0): _router, (2, 0, 0): _router, + (2, 1, 0): _router, } diff --git a/src/ewoksserver/app/routes/tasks/discovery.py b/src/ewoksserver/app/routes/tasks/discovery.py deleted file mode 100644 index 4b89b85..0000000 --- a/src/ewoksserver/app/routes/tasks/discovery.py +++ /dev/null @@ -1,90 +0,0 @@ -import logging - -from ewoksjob.client import discover_all_tasks -from ewoksjob.client import discover_tasks_from_modules -from ewoksjob.client import get_queues -from ewoksjob.client.local import discover_all_tasks as discover_all_tasks_local -from ewoksjob.client.local import ( - discover_tasks_from_modules as discover_tasks_from_modules_local, -) - -from ...config import EwoksSettings -from ...models import EwoksSchedulingType - -logger = logging.getLogger(__name__) - - -def discover_tasks( - settings: EwoksSettings, - modules: list[str] | None = None, - reload: bool | None = None, - task_type: str | None = None, - worker_options: dict | None = None, -) -> list[dict[str, str]]: - """ - :raises ModuleNotFoundError: failed importing tasks. - :raises TimeoutError: timeout when asking a remote worker for tasks. - :raises Exception: any other import or remote error. - """ - if worker_options is None: - kwargs = dict() - else: - kwargs = dict(worker_options) - - # Task discovery: position arguments - if modules: - kwargs["args"] = modules - # Task discovery: named arguments - kwargs["kwargs"] = dict() - if reload is not None: - kwargs["kwargs"]["reload"] = reload - if task_type is not None: - kwargs["kwargs"]["task_type"] = task_type - - timeout = settings.ewoks_discovery.timeout - if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: - if modules: - future = discover_tasks_from_modules_local(**kwargs) - else: - future = discover_all_tasks_local(**kwargs) - tasks = future.result(timeout=timeout) - else: - tasks = _discover_tasks_in_all_queues(kwargs, timeout=timeout) - - for task in tasks: - _set_default_task_properties(task) - return tasks - - -def _discover_tasks_in_all_queues( - kwargs: dict, timeout: float | None = None -) -> list[dict[str, str]]: - discover_from_modules = "args" in kwargs and bool(kwargs["args"]) - discover = ( - discover_tasks_from_modules if discover_from_modules else discover_all_tasks - ) - futures = [discover(**kwargs, queue=queue) for queue in get_queues()] - - # Store tasks in a dict to avoid duplicates - task_dict = {} - for future in futures: - # Ignore failures of a single queue to not prevent discovery on other queues - new_tasks = future.result(timeout=timeout) - exc = future.exception() - if exc: - logger.warning(f"Task discovery failed for {future.queue}: {exc}") - continue - if new_tasks is None: - continue - for task in new_tasks: - task_dict[task["task_identifier"]] = task - return list(task_dict.values()) - - -def _set_default_task_properties(task: dict) -> None: - if not task.get("icon"): - task["icon"] = "default.png" - if not task.get("label"): - task_identifier = task.get("task_identifier") - if task_identifier: - task["label"] = task_identifier.split(".")[-1] diff --git a/src/ewoksserver/app/routes/tasks/router.py b/src/ewoksserver/app/routes/tasks/router.py index c206618..1b9a75a 100644 --- a/src/ewoksserver/app/routes/tasks/router.py +++ b/src/ewoksserver/app/routes/tasks/router.py @@ -10,8 +10,8 @@ from ...backends import json_backend from ...config import EwoksSettingsType from .. import status +from ..common import discovery from ..common import models as common_models -from . import discovery from . import models logger = logging.getLogger(__name__) diff --git a/src/ewoksserver/app/routes/workflows/__init__.py b/src/ewoksserver/app/routes/workflows/__init__.py index a25287a..1969f7d 100644 --- a/src/ewoksserver/app/routes/workflows/__init__.py +++ b/src/ewoksserver/app/routes/workflows/__init__.py @@ -1,7 +1,9 @@ -from .router import router as _router +from .router import v1_0_0_router as _v1_0_0_router +from .router import v2_1_0_router as _v2_1_0_router routers = { - (1, 0, 0): _router, - (1, 1, 0): _router, - (2, 0, 0): _router, + (1, 0, 0): _v1_0_0_router, + (1, 1, 0): _v1_0_0_router, + (2, 0, 0): _v1_0_0_router, + (2, 1, 0): _v2_1_0_router, } diff --git a/src/ewoksserver/app/routes/workflows/models.py b/src/ewoksserver/app/routes/workflows/models.py index a6ff986..1d4dee1 100644 --- a/src/ewoksserver/app/routes/workflows/models.py +++ b/src/ewoksserver/app/routes/workflows/models.py @@ -31,3 +31,13 @@ class EwoksWorkflowIdentifiers(BaseModel): class EwoksWorkflowDescriptions(BaseModel): items: list[EwoksWorkflowDescription] = Field(title="Workflow descriptions") + + +class EwoksWorkflowDiscovery(BaseModel): + modules: list[str] | None = Field( + title="Ewoks workflow modules to discover", default=None + ) + workflow_extension: str | None = Field( + title="Workflow file extension to discover", default=None + ) + worker_options: dict | None = Field(title="Worker options", default=None) diff --git a/src/ewoksserver/app/routes/workflows/router.py b/src/ewoksserver/app/routes/workflows/router.py index 116f50c..e752e87 100644 --- a/src/ewoksserver/app/routes/workflows/router.py +++ b/src/ewoksserver/app/routes/workflows/router.py @@ -10,14 +10,15 @@ from ...backends import json_backend from ...config import EwoksSettingsType from .. import status +from ..common import discovery from ..common import models as common_models from . import descriptions from . import models -router = APIRouter() +v1_0_0_router = APIRouter() -@router.get( +@v1_0_0_router.get( "/workflow/{identifier}", summary="Get ewoks workflow", response_model=models.EwoksWorkflow, @@ -69,7 +70,7 @@ def get_workflow( ) -@router.get( +@v1_0_0_router.get( "/workflows", summary="Get all ewoks workflow identifiers", response_model=models.EwoksWorkflowIdentifiers, @@ -105,7 +106,7 @@ def _compile_keywords(kw: list[str] | None) -> dict | None: return keywords -@router.get( +@v1_0_0_router.get( "/workflows/descriptions", summary="Get all ewoks workflow descriptions", response_model=models.EwoksWorkflowDescriptions, @@ -126,7 +127,7 @@ def get_workflows( } -@router.put( +@v1_0_0_router.put( "/workflow/{identifier}", summary="Update ewoks workflow", response_model=models.EwoksWorkflow, @@ -205,7 +206,7 @@ def update_workflow( return workflow -@router.post( +@v1_0_0_router.post( "/workflows", summary="Create ewoks workflow", response_model=models.EwoksWorkflow, @@ -289,7 +290,7 @@ def create_workflow( return workflow -@router.delete( +@v1_0_0_router.delete( "/workflow/{identifier}", summary="Delete ewoks workflow", response_model=common_models.ResourceInfo, @@ -339,3 +340,46 @@ def delete_workflow( status_code=status.HTTP_404_NOT_FOUND, ) return {"identifier": identifier, "type": "workflow"} + + +v2_1_0_router = APIRouter() +v2_1_0_router.include_router(v1_0_0_router) + + +@v2_1_0_router.post( + "/workflows/discover", + summary="Discover ewoks workflow identifiers from a worker environment", + response_model=models.EwoksWorkflowIdentifiers, + response_description="Discovered ewoks workflow identifiers", + status_code=200, + responses={ + status.HTTP_404_NOT_FOUND: { + "description": "Module not found", + "model": common_models.ResourceError, + }, + }, +) +def discover_workflows( + settings: EwoksSettingsType, + options: Annotated[ + models.EwoksWorkflowDiscovery, Body(title="Ewoks workflow discovery options") + ] = None, +) -> dict[str, list[str]]: + if options: + discover_options = options.model_dump() + else: + discover_options = dict() + try: + identifiers = discovery.discover_workflows(settings, **discover_options) + except ModuleNotFoundError as e: + return JSONResponse( + { + "message": str(e), + "type": "workflow", + }, + status_code=status.HTTP_404_NOT_FOUND, + ) + + print("Discovered workflows not used yet:", identifiers) + + return {"identifiers": identifiers} diff --git a/src/ewoksserver/tests/api_versions.py b/src/ewoksserver/tests/api_versions.py index 3833ca4..2839fe3 100644 --- a/src/ewoksserver/tests/api_versions.py +++ b/src/ewoksserver/tests/api_versions.py @@ -5,8 +5,9 @@ from ..app import routes _API_VERSIONS = { - (): (2, 0, 0), - (2,): (2, 0, 0), + (): (2, 1, 0), + (2,): (2, 1, 0), + (2, 1, 0): (2, 1, 0), (2, 0, 0): (2, 0, 0), (1,): (1, 1, 0), (1, 1, 0): (1, 1, 0), From bc638d35ced6f8d3efff023a7a11b9865ee294cd Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Mon, 13 Jul 2026 15:53:00 +0200 Subject: [PATCH 2/6] tests: POST /api/workflows/discover --- src/ewoksserver/tests/_loadtest/__init__.py | 0 src/ewoksserver/tests/_loadtest/graph.json | 41 ++++++++++++ src/ewoksserver/tests/_loadtest/subgraph.json | 29 +++++++++ src/ewoksserver/tests/conftest.py | 21 +++++++ ...test_discover.py => test_task_discover.py} | 0 .../tests/test_workflow_discover.py | 62 +++++++++++++++++++ 6 files changed, 153 insertions(+) create mode 100644 src/ewoksserver/tests/_loadtest/__init__.py create mode 100644 src/ewoksserver/tests/_loadtest/graph.json create mode 100644 src/ewoksserver/tests/_loadtest/subgraph.json rename src/ewoksserver/tests/{test_discover.py => test_task_discover.py} (100%) create mode 100644 src/ewoksserver/tests/test_workflow_discover.py diff --git a/src/ewoksserver/tests/_loadtest/__init__.py b/src/ewoksserver/tests/_loadtest/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/ewoksserver/tests/_loadtest/graph.json b/src/ewoksserver/tests/_loadtest/graph.json new file mode 100644 index 0000000..5b7555d --- /dev/null +++ b/src/ewoksserver/tests/_loadtest/graph.json @@ -0,0 +1,41 @@ +{ + "graph": { + "id": "graph", + "schema_version": "1.1" + }, + "nodes": [ + { + "id": "node1", + "task_type": "method", + "task_identifier": "dummy", + "default_inputs": [ + { + "name": "name", + "value": "node1" + }, + { + "name": "value", + "value": 0 + } + ] + }, + { + "id": "node2", + "task_type": "graph", + "task_identifier": "subgraph" + } + ], + "links": [ + { + "source": "node1", + "target": "node2", + "sub_target": "in", + "data_mapping": [ + { + "target_input": "value", + "source_output": "return_value" + } + ] + } + ] +} diff --git a/src/ewoksserver/tests/_loadtest/subgraph.json b/src/ewoksserver/tests/_loadtest/subgraph.json new file mode 100644 index 0000000..848fad3 --- /dev/null +++ b/src/ewoksserver/tests/_loadtest/subgraph.json @@ -0,0 +1,29 @@ +{ + "graph": { + "id": "subgraph", + "schema_version": "1.1", + "input_nodes": [ + { + "id": "in", + "node": "subnode1" + } + ] + }, + "nodes": [ + { + "id": "subnode1", + "task_type": "method", + "task_identifier": "dummy", + "default_inputs": [ + { + "name": "name", + "value": "subnode1" + }, + { + "name": "value", + "value": 0 + } + ] + } + ] +} diff --git a/src/ewoksserver/tests/conftest.py b/src/ewoksserver/tests/conftest.py index 8f97141..6aa8f4b 100644 --- a/src/ewoksserver/tests/conftest.py +++ b/src/ewoksserver/tests/conftest.py @@ -5,6 +5,8 @@ import pytest from ewokscore import events +from ewokscore import workflow_discovery +from ewoksjob.client import local as local_client from ewoksjob.tests.conftest import celery_config # noqa F401 from ewoksjob.tests.conftest import celery_includes # noqa F401 from fastapi.testclient import TestClient @@ -46,6 +48,25 @@ def get_ewoks_settings_for_tests(): yield client +@pytest.fixture +def local_patched_ewoks_worker(monkeypatch): + """Only works for local (in-process) ewoksjob worker.""" + monkeypatch.setattr(workflow_discovery, "entry_points", _mock_workflow_entry_points) + + with local_client.pool_context(pool_type="thread"): + yield + + +class _MockEntryPoint: + def __init__(self, name): + self.name = name + + +def _mock_workflow_entry_points(group): + assert group == "ewoks.workflows" + return [_MockEntryPoint("ewoksserver.tests._loadtest.*")] + + @pytest.fixture() def ewoks_handlers(tmpdir): uri = f"file:{tmpdir / 'ewoks_events.db'}" diff --git a/src/ewoksserver/tests/test_discover.py b/src/ewoksserver/tests/test_task_discover.py similarity index 100% rename from src/ewoksserver/tests/test_discover.py rename to src/ewoksserver/tests/test_task_discover.py diff --git a/src/ewoksserver/tests/test_workflow_discover.py b/src/ewoksserver/tests/test_workflow_discover.py new file mode 100644 index 0000000..9ca2167 --- /dev/null +++ b/src/ewoksserver/tests/test_workflow_discover.py @@ -0,0 +1,62 @@ +import pytest +from ewoksjob.client.futures import TimeoutError + +from .api_versions import api_version_bounds + + +@api_version_bounds(min_version="2.1.0") +def test_discover_workflows_from_a_module(rest_client, api_root): + module_pattern = "ewoksserver.tests._loadtest.*" + + response = rest_client.post( + f"{api_root}/workflows/discover", json={"modules": [module_pattern]} + ) + data = response.json() + assert response.status_code == 200, data + expected = [ + "ewoksserver.tests._loadtest.graph", + "ewoksserver.tests._loadtest.subgraph", + ] + assert sorted(data["identifiers"]) == sorted(expected) + + +@api_version_bounds(min_version="2.1.0") +def test_discover_workflow_extension(rest_client, api_root): + module_pattern = "ewoksserver.tests._loadtest.*" + + response = rest_client.post( + f"{api_root}/workflows/discover", + json={"modules": [module_pattern], "workflow_extension": "yaml"}, + ) + data = response.json() + assert response.status_code == 200, data + assert data["identifiers"] == [] + + +@api_version_bounds(min_version="2.1.0") +def test_discover_all_workflows(local_patched_ewoks_worker, rest_client, api_root): + response = rest_client.post(f"{api_root}/workflows/discover") + data = response.json() + assert response.status_code == 200, data + expected = [ + "ewoksserver.tests._loadtest.graph", + "ewoksserver.tests._loadtest.subgraph", + ] + assert set(expected) <= set(data["identifiers"]) + + +@api_version_bounds(min_version="2.1.0") +def test_discover_workflows_in_a_non_existing_module(rest_client, api_root): + response = rest_client.post( + f"{api_root}/workflows/discover", json={"modules": ["not_a_module.foo"]} + ) + data = response.json() + assert response.status_code == 404, data + assert "No module named" in data["message"] + + +@api_version_bounds(min_version="2.1.0") +def test_discover_timeout(celery_discover_timeout_client, api_root): + rest_client, _ = celery_discover_timeout_client + with pytest.raises(TimeoutError): + rest_client.post(f"{api_root}/workflows/discover") From b029971f99d07359a85591c93a6a90c0d65d439e Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Mon, 13 Jul 2026 18:58:04 +0200 Subject: [PATCH 3/6] fix test_new_client_new_events: capture events when the manager starts, not when the reader thread starts --- src/ewoksserver/app/routes/execution/socketio.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/ewoksserver/app/routes/execution/socketio.py b/src/ewoksserver/app/routes/execution/socketio.py index 1f48954..0d02547 100644 --- a/src/ewoksserver/app/routes/execution/socketio.py +++ b/src/ewoksserver/app/routes/execution/socketio.py @@ -59,9 +59,10 @@ async def _start(self) -> None: return self._stop_event.clear() + starttime = datetime.now().astimezone() loop = asyncio.get_running_loop() self._fetch_events_future = loop.run_in_executor( - self._executor, self._fetch_events_main, loop + self._executor, self._fetch_events_main, loop, starttime ) async def _stop(self, timeout: float | None = None) -> None: @@ -71,12 +72,11 @@ async def _stop(self, timeout: float | None = None) -> None: self._stop_event.set() await asyncio.wait_for(future, timeout=timeout) - def _fetch_events_main(self, loop) -> None: + def _fetch_events_main(self, loop, starttime: datetime) -> None: try: with events.reader_context(self._ewoks_settings) as reader: if reader is None: raise RuntimeError("Ewoks event handlers not configured") - starttime = datetime.now().astimezone() for event in reader.wait_events( starttime=starttime, stop_event=self._stop_event ): From 98c236819e68cea9601c571ce1329f0e570d7029 Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Tue, 14 Jul 2026 06:51:52 +0200 Subject: [PATCH 4/6] fix test_discover_timeout: ensure the discovery takes longer than the timeout --- src/ewoksserver/tests/conftest.py | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/src/ewoksserver/tests/conftest.py b/src/ewoksserver/tests/conftest.py index 6aa8f4b..ca44707 100644 --- a/src/ewoksserver/tests/conftest.py +++ b/src/ewoksserver/tests/conftest.py @@ -5,6 +5,7 @@ import pytest from ewokscore import events +from ewokscore import task_discovery from ewokscore import workflow_discovery from ewoksjob.client import local as local_client from ewoksjob.tests.conftest import celery_config # noqa F401 @@ -136,9 +137,19 @@ def get_settings_override(): @pytest.fixture def celery_discover_timeout_client( - tmpdir, celery_session_registered_worker, ewoks_handlers + tmpdir, celery_session_registered_worker, ewoks_handlers, monkeypatch ): """Client to the REST server and Socket.IO (with a very small timeout for discovery)""" + + timeout = 0.1 + + def slow_entry_points(group): + time.sleep(20 * timeout) + return [] + + monkeypatch.setattr(task_discovery, "entry_points", slow_entry_points) + monkeypatch.setattr(workflow_discovery, "entry_points", slow_entry_points) + app = newserver.create_app() def get_settings_override(): @@ -150,7 +161,7 @@ def get_settings_override(): ), ewoks_execution=EwoksExecutionSettings(handlers=ewoks_handlers), # Disable discovery since this client is used to test manual discovery timeout - ewoks_discovery=EwoksDiscoverySettings(on_start_up=False, timeout=0.1), + ewoks_discovery=EwoksDiscoverySettings(on_start_up=False, timeout=timeout), ) app.dependency_overrides[serverconfig.get_ewoks_settings] = get_settings_override From ff286fed972ba9433b9d152c23eaec0436bc9e4d Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Tue, 14 Jul 2026 07:01:59 +0200 Subject: [PATCH 5/6] show events received so far upon timeout --- src/ewoksserver/tests/test_execute.py | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/src/ewoksserver/tests/test_execute.py b/src/ewoksserver/tests/test_execute.py index ef8e1a2..e09fc6b 100644 --- a/src/ewoksserver/tests/test_execute.py +++ b/src/ewoksserver/tests/test_execute.py @@ -200,10 +200,19 @@ def get_events(api_root, sclient, nevents, timeout=10): break time.sleep(0.1) if time.time() - t0 > timeout: - raise TimeoutError(f"Received {len(events)} instead of {nevents}") + received = "\n".join(_event_summary(event) for event in events) + raise TimeoutError( + f"Received {len(events)} instead of {nevents} events:\n{received}" + ) return events +def _event_summary(event: dict) -> str: + if event.get("context") == "node": + return f"{event.get('context')} {event.get('type')} ({event.get('node_id')})" + return f"{event.get('context')} {event.get('type')}" + + def _assert_events(response, events, expected): n = 2 * (len(expected) + 2) assert len(events) == n From 3d1121bb2a2f2386691383c94794c87fc800536c Mon Sep 17 00:00:00 2001 From: woutdenolf Date: Mon, 20 Jul 2026 09:01:37 +0200 Subject: [PATCH 6/6] improve variable names --- .../app/routes/common/discovery.py | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/src/ewoksserver/app/routes/common/discovery.py b/src/ewoksserver/app/routes/common/discovery.py index 5a851d5..d2558bf 100644 --- a/src/ewoksserver/app/routes/common/discovery.py +++ b/src/ewoksserver/app/routes/common/discovery.py @@ -1,4 +1,6 @@ import logging +from typing import Any +from typing import Callable from ewoksjob.client import discover_all_tasks from ewoksjob.client import discover_all_workflows @@ -55,7 +57,7 @@ def discover_tasks( modules=modules, discover_kwargs=discover_kwargs, worker_options=worker_options, - key=lambda task: task["task_identifier"], + id_extractor=lambda task: task["task_identifier"], ) for task in tasks: @@ -95,7 +97,7 @@ def discover_workflows( modules=modules, discover_kwargs=discover_kwargs, worker_options=worker_options, - key=lambda workflow: workflow, + id_extractor=lambda graph_id: graph_id, ) @@ -105,7 +107,7 @@ def _discover( modules: list[str] | None, discover_kwargs: dict, worker_options: dict | None, - key, + id_extractor: Callable[[Any], str], ) -> list: """ :raises ModuleNotFoundError: failed importing tasks or workflows. @@ -128,7 +130,7 @@ def _discover( if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: return _discover_locally(discover, kwargs, timeout=timeout) else: - return _discover_in_all_queues(discover, kwargs, key, timeout=timeout) + return _discover_in_all_queues(discover, kwargs, id_extractor, timeout=timeout) def _discover_locally(discover, kwargs: dict, timeout: float | None = None) -> list: @@ -136,7 +138,10 @@ def _discover_locally(discover, kwargs: dict, timeout: float | None = None) -> l def _discover_in_all_queues( - discover, kwargs: dict, key, timeout: float | None = None + discover, + kwargs: dict, + id_extractor: Callable[[Any], str], + timeout: float | None = None, ) -> list: futures = [discover(**kwargs, queue=queue) for queue in get_queues()] @@ -147,12 +152,12 @@ def _discover_in_all_queues( new_items = future.result(timeout=timeout) exc = future.exception() if exc: - logger.warning(f"Discovery failed for {future.queue}: {exc}") + logger.warning(f"Discovery failed on queue {future.queue!r}: {exc}") continue if new_items is None: continue for item in new_items: - item_dict[key(item)] = item + item_dict[id_extractor(item)] = item return list(item_dict.values())