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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
15 changes: 11 additions & 4 deletions src/ewoksserver/app/lifespan.py
Original file line number Diff line number Diff line change
Expand Up @@ -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__)

Expand All @@ -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

Expand Down Expand Up @@ -64,18 +64,25 @@ 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)
root_url = json_backend.root_url(ewoks_settings.resource_directory, "tasks")
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"""
Expand Down
165 changes: 165 additions & 0 deletions src/ewoksserver/app/routes/common/discovery.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
import logging

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Refactor the task discovery to do both task and workflow discovery.


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,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should it not be workflow id here?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

In contrast to task discovery which returns List[TaskDict] the workflow discovery returns List[str] which are graph ids.

For the variable name is graph_id ok?

)


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]
2 changes: 2 additions & 0 deletions src/ewoksserver/app/routes/execution/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
6 changes: 4 additions & 2 deletions src/ewoksserver/app/routes/execution/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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)


Expand All @@ -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",
Expand Down
6 changes: 3 additions & 3 deletions src/ewoksserver/app/routes/execution/socketio.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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
):
Expand Down
1 change: 1 addition & 0 deletions src/ewoksserver/app/routes/icons/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,4 +4,5 @@
(1, 0, 0): _router,
(1, 1, 0): _router,
(2, 0, 0): _router,
(2, 1, 0): _router,
}
1 change: 1 addition & 0 deletions src/ewoksserver/app/routes/tasks/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,4 +4,5 @@
(1, 0, 0): _router,
(1, 1, 0): _router,
(2, 0, 0): _router,
(2, 1, 0): _router,
}
Loading
Loading