diff --git a/CHANGELOG.md b/CHANGELOG.md index 744875d..2ff11ad 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,7 +20,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added - New API version `v2_1_0`: -- New endpoint `POST /api/workflows/discover` to discover ewoks workflows from python packages. +- New endpoint `POST /api/workflows/discover` to discover external workflows + with a local copy that shadows the external content. + - Local shadowing on discovery can be disabled. + - Editing an external workflow creates a local copy to shadow it. + - Deleting an external workflow without a local shadow fails. + - Re-discovering does not override shadows. ## [2.1.2] - 2026-03-06 diff --git a/src/ewoksserver/app/lifespan.py b/src/ewoksserver/app/lifespan.py index 5982f22..287eab5 100644 --- a/src/ewoksserver/app/lifespan.py +++ b/src/ewoksserver/app/lifespan.py @@ -17,6 +17,7 @@ from .backends import json_backend from .routes.common import discovery from .routes.execution import socketio +from .routes.workflows import backend as workflow_backend logger = logging.getLogger(__name__) @@ -78,10 +79,14 @@ def _rediscover_resources(ewoks_settings: config.EwoksSettings) -> None: json_backend.save_resource(root_url, resource["task_identifier"], resource) try: - discovery.discover_workflows(ewoks_settings) + _, identifier_to_queue = discovery.discover_workflows(ewoks_settings) except Exception as ex: + identifier_to_queue = {} logger.exception("Workflow discovery failed: %s", ex) - logger.warning("Discovered workflows not used yet") + root_url = json_backend.root_url(ewoks_settings.resource_directory, "workflows") + workflow_backend.register_external_workflows( + ewoks_settings, root_url, identifier_to_queue + ) def _enable_execution_events(ewoks_settings: config.EwoksSettings) -> None: diff --git a/src/ewoksserver/app/models.py b/src/ewoksserver/app/models.py index 29a2d4a..b72912d 100644 --- a/src/ewoksserver/app/models.py +++ b/src/ewoksserver/app/models.py @@ -15,9 +15,15 @@ class EwoksSchedulingType(str, Enum): class EwoksDiscoverySettings(BaseModel): - on_start_up: bool = Field(default=True, title="Discover ewoks tasks on startup") + on_start_up: bool = Field( + default=True, title="Discover ewoks tasks/workflows on startup" + ) timeout: float | None = Field( - default=None, title="Timeout for task discovery (in seconds)" + default=None, title="Timeout for task/workflow discovery (in seconds)" + ) + cache_workflows: bool = Field( + default=True, + title="Create a local copy of a workflow when it is discovered", ) diff --git a/src/ewoksserver/app/routes/common/discovery.py b/src/ewoksserver/app/routes/common/discovery.py index d2558bf..1729cfb 100644 --- a/src/ewoksserver/app/routes/common/discovery.py +++ b/src/ewoksserver/app/routes/common/discovery.py @@ -51,7 +51,7 @@ def discover_tasks( if task_type is not None: discover_kwargs["task_type"] = task_type - tasks = _discover( + tasks, _identifier_to_queue = _discover( discover, settings, modules=modules, @@ -70,11 +70,14 @@ def discover_workflows( modules: list[str] | None = None, workflow_extension: str | None = None, worker_options: dict | None = None, -) -> list[str]: +) -> tuple[list[str], dict[str, str | None]]: """ :raises ModuleNotFoundError: failed importing workflows. :raises TimeoutError: timeout when asking a remote worker for workflows. :raises Exception: any other import or remote error. + :returns: the discovered workflow identifiers, and a mapping of each + identifier to the celery queue it was discovered on (`None` for + local scheduling). """ if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: if modules: @@ -108,11 +111,14 @@ def _discover( discover_kwargs: dict, worker_options: dict | None, id_extractor: Callable[[Any], str], -) -> list: +) -> tuple[list, dict[str, str | None]]: """ :raises ModuleNotFoundError: failed importing tasks or workflows. :raises TimeoutError: timeout when asking a remote worker. :raises Exception: any other import or remote error. + :returns: the discovered items, and a mapping of each item's identifier + (return value of `id_extractor`) to the celery queue it was discovered on + (`None` for local scheduling). """ if worker_options is None: kwargs = dict() @@ -128,7 +134,8 @@ def _discover( timeout = settings.ewoks_discovery.timeout if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: - return _discover_locally(discover, kwargs, timeout=timeout) + items = _discover_locally(discover, kwargs, timeout=timeout) + return items, {id_extractor(item): None for item in items} else: return _discover_in_all_queues(discover, kwargs, id_extractor, timeout=timeout) @@ -142,12 +149,13 @@ def _discover_in_all_queues( kwargs: dict, id_extractor: Callable[[Any], str], timeout: float | None = None, -) -> list: - futures = [discover(**kwargs, queue=queue) for queue in get_queues()] +) -> tuple[list, dict[str, str | None]]: + futures = [(queue, discover(**kwargs, queue=queue)) for queue in get_queues()] # Store items in a dict to avoid duplicates item_dict = {} - for future in futures: + identifier_to_queue: dict[str, str | None] = {} + for queue, 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() @@ -157,8 +165,10 @@ def _discover_in_all_queues( if new_items is None: continue for item in new_items: - item_dict[id_extractor(item)] = item - return list(item_dict.values()) + identifier = id_extractor(item) + item_dict[identifier] = item + identifier_to_queue[identifier] = queue + return list(item_dict.values()), identifier_to_queue def _set_default_task_properties(task: dict) -> None: diff --git a/src/ewoksserver/app/routes/workflows/backend.py b/src/ewoksserver/app/routes/workflows/backend.py new file mode 100644 index 0000000..3659faf --- /dev/null +++ b/src/ewoksserver/app/routes/workflows/backend.py @@ -0,0 +1,221 @@ +import json +import logging +from pathlib import Path +from typing import Any +from typing import Iterator + +from ewoksjob.client import convert_graph +from ewoksjob.client.local import convert_graph as convert_graph_local + +from ...backends import json_backend +from ...config import EwoksSettings +from ...models import EwoksSchedulingType + +logger = logging.getLogger(__name__) + + +def load_workflow( + settings: EwoksSettings, + root: json_backend.ResourceUrlType, + identifier: str, + worker_options: dict | None = None, +) -> json_backend.ResourceContentType: + """Load a local or external workflow. + + :raises FileNotFoundError: no local and external workflow + for this identifier. + """ + if json_backend.resource_exists(root, identifier): + return json_backend.load_resource(root, identifier) + + index = _load_external_workflow_index(settings) + if identifier not in index: + raise FileNotFoundError(identifier) + + graph = _load_external_workflow( + settings, identifier, queue=index[identifier], worker_options=worker_options + ) + if graph is None: + raise FileNotFoundError(identifier) + graph.setdefault("graph", {})["id"] = identifier + return graph + + +def save_workflow( + settings: EwoksSettings, + root: json_backend.ResourceUrlType, + identifier: str, + content: json_backend.ResourceContentType, +) -> None: + """Save a workflow, copying it locally when it is an + external workflow. In that case it shadows the external + workflow. + + :raises PermissionError: no permission to save the workflow. + """ + index = _load_external_workflow_index(settings) + if identifier in index: + del index[identifier] + _save_external_workflow_index(settings, index) + + json_backend.save_resource(root, identifier, content) + + +def delete_workflow(root: json_backend.ResourceUrlType, identifier: str) -> None: + """Delete a local workflow. When it shadows an external workflow + the external workflow needs to be re-discovered. + + :raises PermissionError: no permission to delete the workflow. + :raises FileNotFoundError: no local workflow for this identifier. + """ + json_backend.delete_resource(root, identifier) + + +def workflow_exists( + settings: EwoksSettings, root: json_backend.ResourceUrlType, identifier: str +) -> bool: + """Whether a local or external workflow exists for this identifier. + + :raises ValueError: invalid identifier. + """ + return json_backend.resource_exists(root, identifier) or is_external_workflow( + settings, identifier + ) + + +def workflow_identifiers( + settings: EwoksSettings, root: json_backend.ResourceUrlType +) -> list[str]: + """Identifiers of local and external workflows.""" + identifiers = set(json_backend.resource_identifiers(root)) + identifiers.update(_load_external_workflow_index(settings)) + return sorted(identifiers) + + +def iter_workflow_graphs( + settings: EwoksSettings, + root: json_backend.ResourceUrlType, + worker_options: dict | None = None, +) -> Iterator[dict]: + """Yield `graph` attributes of local or external workflows.""" + shadowed = set() + for identifier in json_backend.resource_identifiers(root): + shadowed.add(identifier) + yield json_backend.load_resource(root, identifier).get("graph", {}) + + index = _load_external_workflow_index(settings) + for identifier, queue in index.items(): + if identifier in shadowed: + continue + graph = _load_external_workflow( + settings, identifier, queue=queue, worker_options=worker_options + ) + if graph is None: + continue + graph.setdefault("graph", {})["id"] = identifier + yield graph["graph"] + + +def register_external_workflows( + settings: EwoksSettings, + root: json_backend.ResourceUrlType, + identifier_to_queue: dict[str, str | None], + worker_options: dict | None = None, +) -> None: + """Register discovered external workflows, skipping ones already + shadowed locally. Persists a shadow right away if `cache_workflows` + is enabled. + """ + create_shadow = settings.ewoks_discovery.cache_workflows + index = _load_external_workflow_index(settings) + index_changed = False + for identifier, queue in identifier_to_queue.items(): + if json_backend.resource_exists(root, identifier): + continue + + if create_shadow: + graph = _load_external_workflow( + settings, identifier, queue=queue, worker_options=worker_options + ) + if graph is None: + continue + graph.setdefault("graph", {})["id"] = identifier + json_backend.save_resource(root, identifier, graph) + + if identifier not in index or index[identifier] != queue: + index[identifier] = queue + index_changed = True + + if index_changed: + _save_external_workflow_index(settings, index) + + +def is_external_workflow(settings: EwoksSettings, identifier: str) -> bool: + """Whether a workflow is registered as a external workflow.""" + return identifier in _load_external_workflow_index(settings) + + +_EXTERNAL_WORKFLOW_INDEX = "external_workflow_index.json" + + +def _external_workflow_index_path(settings: EwoksSettings) -> Path: + return settings.resource_directory / _EXTERNAL_WORKFLOW_INDEX + + +def _load_external_workflow_index(settings: EwoksSettings) -> dict[str, Any]: + """The external workflow index: identifier -> discovery queue.""" + try: + with open(_external_workflow_index_path(settings)) as f: + return json.load(f) + except FileNotFoundError: + return {} + + +def _save_external_workflow_index( + settings: EwoksSettings, index: dict[str, Any] +) -> None: + """The external workflow index: identifier -> discovery queue.""" + path = _external_workflow_index_path(settings) + path.parent.mkdir(parents=True, exist_ok=True) + with open(path, "w") as f: + json.dump(index, f, indent=2) + + +def _load_external_workflow( + settings: EwoksSettings, + identifier: str, + queue: str | None = None, + worker_options: dict | None = None, +) -> dict | None: + """Load a external workflow identified by its fully qualified + module identifier, e.g.``"mypackage.subpackage.myworkflow"``. + + :returns: `None` when the workflow could not be loaded. + """ + package, _, _ = identifier.rpartition(".") + if not package: + return None + + if worker_options is None: + kwargs = dict() + else: + kwargs = dict(worker_options) + kwargs["args"] = (identifier, None) + kwargs["kwargs"] = { + "load_options": {"representation": "json_module", "root_module": package} + } + + timeout = settings.ewoks_discovery.timeout + try: + if settings.ewoks_scheduling.type == EwoksSchedulingType.Local: + future = convert_graph_local(**kwargs) + else: + future = convert_graph(**kwargs, queue=queue) + graph = future.result(timeout=timeout) + except Exception as ex: + logger.warning("Failed to load external workflow %r: %s", identifier, ex) + return None + + if not isinstance(graph, dict): + return None + return graph diff --git a/src/ewoksserver/app/routes/workflows/descriptions.py b/src/ewoksserver/app/routes/workflows/descriptions.py index 6727a5b..7663e47 100644 --- a/src/ewoksserver/app/routes/workflows/descriptions.py +++ b/src/ewoksserver/app/routes/workflows/descriptions.py @@ -1,6 +1,8 @@ from typing import Iterator from ...backends import json_backend +from ...config import EwoksSettings +from . import backend _WORKFLOW_KEYWORDS = ( "id", @@ -13,10 +15,11 @@ def workflow_descriptions( - root: json_backend.ResourceUrlType, keywords: dict | None = None + settings: EwoksSettings, + root: json_backend.ResourceUrlType, + keywords: dict | None = None, ) -> Iterator[dict]: - for res in json_backend.resources(root): - description = res["graph"] + for description in backend.iter_workflow_graphs(settings, root): if not _include_resource(description.get("keywords", dict()), keywords): continue yield { diff --git a/src/ewoksserver/app/routes/workflows/router.py b/src/ewoksserver/app/routes/workflows/router.py index 142a444..b589778 100644 --- a/src/ewoksserver/app/routes/workflows/router.py +++ b/src/ewoksserver/app/routes/workflows/router.py @@ -12,6 +12,7 @@ from .. import status from ..common import discovery from ..common import models as common_models +from . import backend from . import descriptions from . import models @@ -47,8 +48,8 @@ def get_workflow( settings: EwoksSettingsType, ) -> json_backend.ResourceContentType: try: - return json_backend.load_resource( - settings.resource_directory / "workflows", identifier + return backend.load_workflow( + settings, settings.resource_directory / "workflows", identifier ) except PermissionError: return JSONResponse( @@ -85,10 +86,12 @@ def get_workflow_identifiers( if keywords: identifiers = [ desc["id"] - for desc in descriptions.workflow_descriptions(root, keywords=keywords) + for desc in descriptions.workflow_descriptions( + settings, root, keywords=keywords + ) ] else: - identifiers = list(json_backend.resource_identifiers(root)) + identifiers = backend.workflow_identifiers(settings, root) return {"identifiers": identifiers} @@ -121,7 +124,7 @@ def get_workflows( return { "items": list( descriptions.workflow_descriptions( - settings.resource_directory / "workflows", keywords=keywords + settings, settings.resource_directory / "workflows", keywords=keywords ) ) } @@ -174,8 +177,8 @@ def update_workflow( status_code=status.HTTP_400_BAD_REQUEST, ) - exists = json_backend.resource_exists( - settings.resource_directory / "workflows", identifier + exists = backend.workflow_exists( + settings, settings.resource_directory / "workflows", identifier ) if not exists: return JSONResponse( @@ -188,7 +191,8 @@ def update_workflow( ) try: - json_backend.save_resource( + backend.save_workflow( + settings, settings.resource_directory / "workflows", identifier, workflow.model_dump(exclude_none=True), @@ -259,8 +263,8 @@ def create_workflow( ) try: - exists = json_backend.resource_exists( - settings.resource_directory / "workflows", ridentifier + exists = backend.workflow_exists( + settings, settings.resource_directory / "workflows", ridentifier ) except ValueError: return JSONResponse( @@ -282,7 +286,8 @@ def create_workflow( ) try: - json_backend.save_resource( + backend.save_workflow( + settings, settings.resource_directory / "workflows", ridentifier, workflow.model_dump(exclude_none=True), @@ -308,7 +313,8 @@ def create_workflow( status_code=200, responses={ status.HTTP_403_FORBIDDEN: { - "description": "No permission to read workflow", + "description": "No permission to delete workflow, or the workflow " + "is not managed by ewoksserver and cannot be deleted", "model": common_models.ResourceIdentifierError, }, status.HTTP_404_NOT_FOUND: { @@ -327,10 +333,19 @@ def delete_workflow( ], settings: EwoksSettingsType, ) -> dict[str, str]: - try: - json_backend.delete_resource( - settings.resource_directory / "workflows", identifier + if backend.is_external_workflow(settings, identifier): + return JSONResponse( + { + "message": f"Workflow '{identifier}' is not managed by ewoksserver " + "and cannot be deleted.", + "type": "workflow", + "identifier": identifier, + }, + status_code=status.HTTP_403_FORBIDDEN, ) + + try: + backend.delete_workflow(settings.resource_directory / "workflows", identifier) except PermissionError: return JSONResponse( { @@ -380,7 +395,9 @@ def discover_workflows( else: discover_options = dict() try: - identifiers = discovery.discover_workflows(settings, **discover_options) + identifiers, identifier_to_queue = discovery.discover_workflows( + settings, **discover_options + ) except ModuleNotFoundError as e: return JSONResponse( { @@ -390,6 +407,11 @@ def discover_workflows( status_code=status.HTTP_404_NOT_FOUND, ) - print("Discovered workflows not used yet:", identifiers) + backend.register_external_workflows( + settings, + settings.resource_directory / "workflows", + identifier_to_queue, + worker_options=discover_options.get("worker_options"), + ) return {"identifiers": identifiers} diff --git a/src/ewoksserver/tests/conftest.py b/src/ewoksserver/tests/conftest.py index 95d5b19..093192b 100644 --- a/src/ewoksserver/tests/conftest.py +++ b/src/ewoksserver/tests/conftest.py @@ -49,6 +49,31 @@ def get_ewoks_settings_for_tests(): yield client +@pytest.fixture +def rest_client_no_discover_cache(tmp_path): + """Client to the REST server (no execution) without workflow + caching on discovery.""" + app = newserver.create_app() + + @lru_cache() + def get_ewoks_settings_for_tests(): + return serverconfig.EwoksSettings( + configured=True, + resource_directory=str(tmp_path), + # Disable discovery since this client is used to test manual discovery + ewoks_discovery=EwoksDiscoverySettings( + on_start_up=False, cache_workflows=False + ), + ) + + app.dependency_overrides[serverconfig.get_ewoks_settings] = ( + get_ewoks_settings_for_tests + ) + + with TestClient(app) as client: + yield client + + @pytest.fixture def local_patched_ewoks_worker(monkeypatch): """Only works for local (in-process) ewoksjob worker.""" diff --git a/src/ewoksserver/tests/test_workflow_discover.py b/src/ewoksserver/tests/test_workflow_discover.py index 9ca2167..c7f0037 100644 --- a/src/ewoksserver/tests/test_workflow_discover.py +++ b/src/ewoksserver/tests/test_workflow_discover.py @@ -1,3 +1,5 @@ +import json + import pytest from ewoksjob.client.futures import TimeoutError @@ -60,3 +62,213 @@ 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") + + +@api_version_bounds(min_version="2.1.0") +def test_cache_on_discovery(rest_client, api_root, tmp_path): + module_pattern = "ewoksserver.tests._loadtest.*" + + response = rest_client.post( + f"{api_root}/workflows/discover", json={"modules": [module_pattern]} + ) + assert response.status_code == 200, response.json() + + # Discovered workflows are shadowed right away. + assert _workflow_file(tmp_path, "ewoksserver.tests._loadtest.subgraph").exists() + assert _workflow_file(tmp_path, "ewoksserver.tests._loadtest.graph").exists() + + with open(tmp_path / "external_workflow_index.json") as f: + index = json.load(f) + assert index == { + "ewoksserver.tests._loadtest.graph": None, + "ewoksserver.tests._loadtest.subgraph": None, + } + + response = rest_client.get(f"{api_root}/workflows") + data = response.json() + assert "ewoksserver.tests._loadtest.graph" in data["identifiers"] + assert "ewoksserver.tests._loadtest.subgraph" in data["identifiers"] + + # Standalone workflow (no subgraph reference): content is preserved as-is. + identifier = "ewoksserver.tests._loadtest.subgraph" + response = rest_client.get(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 200, data + assert data["graph"]["id"] == identifier + assert [node["id"] for node in data["nodes"]] == ["subnode1"] + + # Workflow referencing a subgraph: converting flattens it into one graph. + identifier = "ewoksserver.tests._loadtest.graph" + response = rest_client.get(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 200, data + assert data["graph"]["id"] == identifier + + +@api_version_bounds(min_version="2.1.0") +def test_no_cache_on_discovery(rest_client_no_discover_cache, api_root, tmp_path): + module_pattern = "ewoksserver.tests._loadtest.*" + + response = rest_client_no_discover_cache.post( + f"{api_root}/workflows/discover", json={"modules": [module_pattern]} + ) + assert response.status_code == 200, response.json() + + # Discovered workflows are not shadowed right away. + assert not _workflow_file(tmp_path, "ewoksserver.tests._loadtest.subgraph").exists() + assert not _workflow_file(tmp_path, "ewoksserver.tests._loadtest.graph").exists() + + with open(tmp_path / "external_workflow_index.json") as f: + index = json.load(f) + assert index == { + "ewoksserver.tests._loadtest.graph": None, + "ewoksserver.tests._loadtest.subgraph": None, + } + + # They are still exposed through the REST API. + response = rest_client_no_discover_cache.get(f"{api_root}/workflows") + data = response.json() + assert "ewoksserver.tests._loadtest.graph" in data["identifiers"] + assert "ewoksserver.tests._loadtest.subgraph" in data["identifiers"] + + # Standalone workflow (no subgraph reference): content is preserved as-is. + identifier = "ewoksserver.tests._loadtest.subgraph" + response = rest_client_no_discover_cache.get(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 200, data + assert data["graph"]["id"] == identifier + assert [node["id"] for node in data["nodes"]] == ["subnode1"] + + # Workflow referencing a subgraph: converting flattens it into one graph. + identifier = "ewoksserver.tests._loadtest.graph" + response = rest_client_no_discover_cache.get(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 200, data + assert data["graph"]["id"] == identifier + + # Loading a workflow converts it on the fly; it still does not persist it. + assert not _workflow_file(tmp_path, "ewoksserver.tests._loadtest.subgraph").exists() + assert not _workflow_file(tmp_path, "ewoksserver.tests._loadtest.graph").exists() + + +@api_version_bounds(min_version="2.1.0") +def test_discover_does_not_override_local_copy(rest_client, api_root, tmp_path): + identifier = "ewoksserver.tests._loadtest.subgraph" + custom_workflow = { + "graph": {"id": identifier, "label": "custom"}, + "nodes": [], + } + response = rest_client.post(f"{api_root}/workflows", json=custom_workflow) + assert response.status_code == 200, response.json() + + response = rest_client.post( + f"{api_root}/workflows/discover", + json={"modules": ["ewoksserver.tests._loadtest.*"]}, + ) + assert response.status_code == 200, response.json() + + response = rest_client.get(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 200, data + assert data == custom_workflow + + # The local identifier is not registered as a external workflow, + # so it stays a normal, deletable, locally owned workflow. + with open(tmp_path / "external_workflow_index.json") as f: + index = json.load(f) + assert identifier not in index + + +@api_version_bounds(min_version="2.1.0") +def test_delete_external_workflow_is_not_allowed(rest_client, api_root): + identifier = "ewoksserver.tests._loadtest.subgraph" + response = rest_client.post( + f"{api_root}/workflows/discover", + json={"modules": ["ewoksserver.tests._loadtest.*"]}, + ) + assert response.status_code == 200, response.json() + + response = rest_client.delete(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 403, data + assert data["identifier"] == identifier + + response = rest_client.get(f"{api_root}/workflow/{identifier}") + assert response.status_code == 200 + + +@api_version_bounds(min_version="2.1.0") +def test_delete_external_workflow_after_edit_is_allowed( + rest_client_no_discover_cache, api_root, tmp_path +): + identifier = "ewoksserver.tests._loadtest.subgraph" + response = rest_client_no_discover_cache.post( + f"{api_root}/workflows/discover", + json={"modules": ["ewoksserver.tests._loadtest.*"]}, + ) + assert response.status_code == 200, response.json() + assert not _workflow_file(tmp_path, identifier).exists() + + edited_workflow = {"graph": {"id": identifier, "label": "edited"}, "nodes": []} + response = rest_client_no_discover_cache.put( + f"{api_root}/workflow/{identifier}", json=edited_workflow + ) + assert response.status_code == 200, response.json() + + # The first edit creates the local shadow. + assert _workflow_file(tmp_path, identifier).exists() + + response = rest_client_no_discover_cache.delete(f"{api_root}/workflow/{identifier}") + data = response.json() + assert response.status_code == 200, data + + response = rest_client_no_discover_cache.get(f"{api_root}/workflow/{identifier}") + assert response.status_code == 404 + + +@api_version_bounds(min_version="2.1.0") +def test_edit_external_workflow_creates_shadow( + rest_client_no_discover_cache, api_root, tmp_path +): + identifier = "ewoksserver.tests._loadtest.subgraph" + response = rest_client_no_discover_cache.post( + f"{api_root}/workflows/discover", + json={"modules": ["ewoksserver.tests._loadtest.*"]}, + ) + assert response.status_code == 200, response.json() + assert not _workflow_file(tmp_path, identifier).exists() + + # Editing creates a local shadow. + edited_workflow = {"graph": {"id": identifier, "label": "edited"}, "nodes": []} + response = rest_client_no_discover_cache.put( + f"{api_root}/workflow/{identifier}", json=edited_workflow + ) + assert response.status_code == 200, response.json() + assert _workflow_file(tmp_path, identifier).exists() + + with open(tmp_path / "external_workflow_index.json") as f: + index = json.load(f) + assert identifier not in index + + response = rest_client_no_discover_cache.get(f"{api_root}/workflow/{identifier}") + assert response.status_code == 200 + assert response.json() == edited_workflow + + +@api_version_bounds(min_version="2.1.0") +def test_create_workflow_conflicts_with_external_workflow(rest_client, api_root): + identifier = "ewoksserver.tests._loadtest.subgraph" + response = rest_client.post( + f"{api_root}/workflows/discover", + json={"modules": ["ewoksserver.tests._loadtest.*"]}, + ) + assert response.status_code == 200, response.json() + + new_workflow = {"graph": {"id": identifier, "label": "new"}, "nodes": []} + response = rest_client.post(f"{api_root}/workflows", json=new_workflow) + data = response.json() + assert response.status_code == 409, data + + +def _workflow_file(tmp_path, identifier): + return tmp_path / "workflows" / f"{identifier}.json"