from fastapi.testclient import TestClient from app.core.config import settings PREFIX = f"{settings.API_V1_STR}/flows" WORKING_NODE = """ def process(params): return {"reading": 21.5} """ BROKEN_NODE = """ def process(params): raise RuntimeError("boom") """ def a_flow(name: str = "demo") -> dict: return { "name": name, "title": "Demo", "nodes": [ { "id": "sensor", "type": "python", "position": {"x": 0, "y": 0}, "provides": [{"name": "reading", "dtype": "float"}], }, { "id": "logger", "type": "python", "position": {"x": 240, "y": 0}, "requires": [{"name": "reading", "dtype": "float"}], }, ], } def test_flows_require_authentication(client: TestClient) -> None: assert client.get(f"{PREFIX}/").status_code == 401 def test_save_then_read_flow( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: response = client.put( f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow() ) assert response.status_code == 200 response = client.get(f"{PREFIX}/demo", headers=superuser_token_headers) assert response.status_code == 200 body = response.json() assert body["definition"]["title"] == "Demo" assert {node["id"] for node in body["definition"]["nodes"]} == {"sensor", "logger"} def test_name_mismatch_is_rejected( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: response = client.put( f"{PREFIX}/other", headers=superuser_token_headers, json=a_flow("demo") ) assert response.status_code == 400 def test_broken_node_is_reported_and_siblings_stay_active( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow()) client.put( f"{PREFIX}/demo/nodes/sensor/source", headers=superuser_token_headers, json={"code": WORKING_NODE}, ) response = client.put( f"{PREFIX}/demo/nodes/logger/source", headers=superuser_token_headers, json={"code": "def process(reading, params:\n"}, ) assert response.status_code == 200 assert response.json()["status"] == "error" statuses = client.get(f"{PREFIX}/demo", headers=superuser_token_headers).json() by_id = {node["id"]: node for node in statuses["nodes"]} assert by_id["demo.sensor"]["status"] == "active" assert by_id["demo.logger"]["status"] == "error" def test_running_a_flow_produces_values( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow()) client.put( f"{PREFIX}/demo/nodes/sensor/source", headers=superuser_token_headers, json={"code": WORKING_NODE}, ) client.put( f"{PREFIX}/demo/nodes/logger/source", headers=superuser_token_headers, json={"code": "def process(reading, params):\n return {}\n"}, ) response = client.post( f"{PREFIX}/demo/run", headers=superuser_token_headers, json={"inputs": {}} ) assert response.status_code == 200 assert response.json()["values"]["demo.reading"]["value"] == 21.5 def test_editing_does_not_deploy_until_published( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: saved = client.put( f"{PREFIX}/staged", headers=superuser_token_headers, json=a_flow("staged") ).json() assert saved["has_draft"] is True # Nothing is running yet, so the engine knows no nodes of this flow. assert client.get(f"{PREFIX}/staged/state", headers=superuser_token_headers).json()[ "nodes" ] == [] published = client.post( f"{PREFIX}/staged/publish", headers=superuser_token_headers, json={"version": saved["definition"]["version"]}, ) assert published.status_code == 200 assert published.json()["has_draft"] is False assert { node["id"] for node in client.get( f"{PREFIX}/staged/state", headers=superuser_token_headers ).json()["nodes"] } == {"staged.sensor", "staged.logger"} client.delete(f"{PREFIX}/staged", headers=superuser_token_headers) def test_a_stale_save_is_refused( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: saved = client.put( f"{PREFIX}/contested", headers=superuser_token_headers, json=a_flow("contested") ).json() stale = saved["definition"] client.put( f"{PREFIX}/contested", headers=superuser_token_headers, json={**stale, "title": "Mine"}, ) response = client.put( f"{PREFIX}/contested", headers=superuser_token_headers, json={**stale, "title": "Theirs"}, ) assert response.status_code == 409 assert response.json()["detail"]["current_version"] == stale["version"] + 1 client.delete(f"{PREFIX}/contested", headers=superuser_token_headers) def test_discarding_a_draft_restores_what_is_running( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: saved = client.put( f"{PREFIX}/reverted", headers=superuser_token_headers, json=a_flow("reverted") ).json() client.post( f"{PREFIX}/reverted/publish", headers=superuser_token_headers, json={"version": saved["definition"]["version"]}, ) published = client.get(f"{PREFIX}/reverted", headers=superuser_token_headers).json() client.put( f"{PREFIX}/reverted", headers=superuser_token_headers, json={**published["definition"], "title": "Scratch that"}, ) response = client.post( f"{PREFIX}/reverted/discard-draft", headers=superuser_token_headers ) assert response.status_code == 200 assert response.json()["has_draft"] is False assert response.json()["definition"]["title"] == "Demo" client.delete(f"{PREFIX}/reverted", headers=superuser_token_headers) def test_stopping_a_flow_takes_it_off_the_engine( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: saved = client.put( f"{PREFIX}/halted", headers=superuser_token_headers, json=a_flow("halted") ).json() client.post( f"{PREFIX}/halted/publish", headers=superuser_token_headers, json={"version": saved["definition"]["version"]}, ) stopped = client.post(f"{PREFIX}/halted/stop", headers=superuser_token_headers) assert stopped.status_code == 200 assert stopped.json()["enabled"] is False # Nothing of it is loaded, so there is nothing to run. assert ( client.post( f"{PREFIX}/halted/run", headers=superuser_token_headers, json={"inputs": {}} ).status_code == 409 ) started = client.post(f"{PREFIX}/halted/start", headers=superuser_token_headers) assert started.json()["enabled"] is True assert ( client.post( f"{PREFIX}/halted/run", headers=superuser_token_headers, json={"inputs": {}} ).status_code == 200 ) client.delete(f"{PREFIX}/halted", headers=superuser_token_headers) def test_pausing_is_reported_back( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: saved = client.put( f"{PREFIX}/held", headers=superuser_token_headers, json=a_flow("held") ).json() client.post( f"{PREFIX}/held/publish", headers=superuser_token_headers, json={"version": saved["definition"]["version"]}, ) assert ( client.post(f"{PREFIX}/held/pause", headers=superuser_token_headers).status_code == 200 ) assert ( client.get(f"{PREFIX}/held", headers=superuser_token_headers).json()["paused"] is True ) client.post(f"{PREFIX}/held/resume", headers=superuser_token_headers) assert ( client.get(f"{PREFIX}/held", headers=superuser_token_headers).json()["paused"] is False ) client.delete(f"{PREFIX}/held", headers=superuser_token_headers) def test_unconnected_input_is_surfaced( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: flow = a_flow() flow["nodes"][0]["provides"] = [] # nothing produces "reading" any more client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=flow) issues = client.post( f"{PREFIX}/demo/validate", headers=superuser_token_headers ).json()["issues"] assert any(issue["code"] == "unconnected_input" for issue in issues) def test_delete_flow( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow()) assert ( client.delete(f"{PREFIX}/demo", headers=superuser_token_headers).status_code == 200 ) assert ( client.get(f"{PREFIX}/demo", headers=superuser_token_headers).status_code == 404 ) def test_node_types_are_listed( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: types = client.get(f"{PREFIX}/node-types", headers=superuser_token_headers).json() by_type = {entry["type"]: entry for entry in types} assert by_type["python"]["has_source"] is True assert "properties" in by_type["mqtt"]["params_schema"] def test_rename_flow( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow()) response = client.post( f"{PREFIX}/demo/rename", headers=superuser_token_headers, json={"new_name": "demo_renamed"}, ) assert response.status_code == 200 assert response.json()["definition"]["name"] == "demo_renamed" assert ( client.get(f"{PREFIX}/demo", headers=superuser_token_headers).status_code == 404 ) assert ( client.get( f"{PREFIX}/demo_renamed", headers=superuser_token_headers ).status_code == 200 ) # Put it back so the tests that follow find the flow they expect. client.post( f"{PREFIX}/demo_renamed/rename", headers=superuser_token_headers, json={"new_name": "demo"}, ) def test_rename_onto_a_taken_name_is_refused( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow()) client.put( f"{PREFIX}/occupied", headers=superuser_token_headers, json=a_flow("occupied") ) response = client.post( f"{PREFIX}/demo/rename", headers=superuser_token_headers, json={"new_name": "occupied"}, ) assert response.status_code == 409 client.delete(f"{PREFIX}/occupied", headers=superuser_token_headers) def test_rename_rejects_an_invalid_name( client: TestClient, superuser_token_headers: dict[str, str] ) -> None: client.put(f"{PREFIX}/demo", headers=superuser_token_headers, json=a_flow()) response = client.post( f"{PREFIX}/demo/rename", headers=superuser_token_headers, json={"new_name": "Not A Flow Name"}, ) assert response.status_code == 400