Skip to content

Commit 080ad75

Browse files
committed
add worker tasks discover_workflows_from_modules and discover_all_workflows
1 parent cb8e2b1 commit 080ad75

11 files changed

Lines changed: 122 additions & 1 deletion

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
1212
- Add ewoksjob options under `EWOKSJOB_OPTIONS`: ``log_memory_usage`` and ``detect_memory_leaks``.
1313
- Added ability to handle remote ewoksutils exception types on the client side.
1414
- Add `client.task_utils.TaskSubmitter` helper to execute a single task.
15+
- Added worker tasks `discover_workflows_from_modules` and `discover_all_workflows`.
1516

1617
### Fixed
1718

pyproject.toml

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ Changelog = "https://github.com/ewoks-kit/ewoksjob/-/blob/main/CHANGELOG.md"
3030

3131
[project.optional-dependencies]
3232
worker = [
33-
"ewoks >=4.0.0",
33+
"ewoks >=6.0.0",
3434
"objgraph",
3535
]
3636
redis = [
@@ -95,6 +95,9 @@ package-dir = { "" = "src" }
9595
[tool.setuptools.packages.find]
9696
where = ["src"]
9797

98+
[tool.setuptools.package-data]
99+
"*" = ["*.json"]
100+
98101
[tool.coverage.run]
99102
omit = ['*/tests/*']
100103

src/ewoksjob/apps/_decorators/ewoks_tasks.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
import ewoks
77
from ewokscore import task_discovery
8+
from ewokscore import workflow_discovery
89

910
from ...worker.executor import get_execute_method
1011
from ..errors import replace_exception_for_client
@@ -49,6 +50,8 @@ def new_celery_task(*args, **kwargs) -> Any:
4950
"convert_graph": ewoks.convert_graph,
5051
"discover_tasks_from_modules": task_discovery.discover_tasks_from_modules,
5152
"discover_all_tasks": task_discovery.discover_all_tasks,
53+
"discover_workflows_from_modules": workflow_discovery.discover_workflows_from_modules,
54+
"discover_all_workflows": workflow_discovery.discover_all_workflows,
5255
}
5356

5457
_BOUND_TASKS = {"execute_graph"}

src/ewoksjob/apps/ewoks.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,3 +44,17 @@ def discover_tasks_from_modules(*args, **kwargs) -> List[dict]:
4444
@ewoks_tasks.ewoks_task
4545
def discover_all_tasks(*args, **kwargs) -> List[dict]:
4646
return ewoks_tasks.get_native_ewoks_task("discover_all_tasks")(*args, **kwargs)
47+
48+
49+
@app.task(bind=False)
50+
@ewoks_tasks.ewoks_task
51+
def discover_workflows_from_modules(*args, **kwargs) -> List[dict]:
52+
return ewoks_tasks.get_native_ewoks_task("discover_workflows_from_modules")(
53+
*args, **kwargs
54+
)
55+
56+
57+
@app.task(bind=False)
58+
@ewoks_tasks.ewoks_task
59+
def discover_all_workflows(*args, **kwargs) -> List[dict]:
60+
return ewoks_tasks.get_native_ewoks_task("discover_all_workflows")(*args, **kwargs)

src/ewoksjob/client/celery/tasks.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,8 @@
1212
"convert_workflow",
1313
"discover_tasks_from_modules",
1414
"discover_all_tasks",
15+
"discover_workflows_from_modules",
16+
"discover_all_workflows",
1517
]
1618

1719

@@ -55,3 +57,15 @@ def discover_tasks_from_modules(**kw) -> CeleryFuture:
5557
def discover_all_tasks(**kw) -> CeleryFuture:
5658
async_result = send_task("ewoksjob.apps.ewoks.discover_all_tasks", **kw)
5759
return CeleryFuture(async_result.id, async_result)
60+
61+
62+
def discover_workflows_from_modules(**kw) -> CeleryFuture:
63+
async_result = send_task(
64+
"ewoksjob.apps.ewoks.discover_workflows_from_modules", **kw
65+
)
66+
return CeleryFuture(async_result.id, async_result)
67+
68+
69+
def discover_all_workflows(**kw) -> CeleryFuture:
70+
async_result = send_task("ewoksjob.apps.ewoks.discover_all_workflows", **kw)
71+
return CeleryFuture(async_result.id, async_result)

src/ewoksjob/client/local/tasks.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66

77
import ewoks
88
from ewokscore import task_discovery
9+
from ewokscore import workflow_discovery
910

1011
from ..dummy_workflow import dummy_workflow
1112
from .futures import LocalFuture
@@ -18,6 +19,8 @@
1819
"convert_workflow",
1920
"discover_tasks_from_modules",
2021
"discover_all_tasks",
22+
"discover_workflows_from_modules",
23+
"discover_all_workflows",
2124
]
2225

2326

@@ -74,3 +77,21 @@ def discover_all_tasks(
7477
) -> LocalFuture:
7578
pool = get_active_pool()
7679
return pool.submit(task_discovery.discover_all_tasks, args=args, kwargs=kwargs)
80+
81+
82+
def discover_workflows_from_modules(
83+
args: Optional[Tuple] = tuple(), kwargs: Optional[Mapping] = None
84+
) -> LocalFuture:
85+
pool = get_active_pool()
86+
return pool.submit(
87+
workflow_discovery.discover_workflows_from_modules, args=args, kwargs=kwargs
88+
)
89+
90+
91+
def discover_all_workflows(
92+
args: Optional[Tuple] = tuple(), kwargs: Optional[Mapping] = None
93+
) -> LocalFuture:
94+
pool = get_active_pool()
95+
return pool.submit(
96+
workflow_discovery.discover_all_workflows, args=args, kwargs=kwargs
97+
)

src/ewoksjob/tests/_loadtest/__init__.py

Whitespace-only changes.
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
{"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"}]}]}
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
{"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}]}]}

src/ewoksjob/tests/conftest.py

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33

44
import pytest
55
from ewokscore import events
6+
from ewokscore import workflow_discovery
67

78
from ewoksjob.events.readers import read_ewoks_events
89
from ewoksjob.worker import options as worker_options
@@ -130,6 +131,24 @@ def local_ewoks_worker(slurm_client_kwargs):
130131
assert len(pool._tasks) == 0, str(list(pool._tasks.values()))
131132

132133

134+
@pytest.fixture
135+
def local_patched_ewoks_worker(monkeypatch):
136+
monkeypatch.setattr(workflow_discovery, "entry_points", _mock_entry_points)
137+
138+
with local.pool_context(pool_type="thread"):
139+
yield
140+
141+
142+
class _MockEntryPoint:
143+
def __init__(self, name):
144+
self.name = name
145+
146+
147+
def _mock_entry_points(group):
148+
assert group == "ewoks.workflows"
149+
return [_MockEntryPoint("ewoksjob.tests._loadtest.*")]
150+
151+
133152
@pytest.fixture()
134153
def sqlite3_ewoks_events(tmp_path):
135154
uri = f"file:{tmp_path / 'ewoks_events.db'}"

0 commit comments

Comments
 (0)