from __future__ import annotations import time from dataclasses import replace from io import BytesIO from pathlib import Path from app.discovery import TaskCatalog from app.runner import RunManager from app.inputs import PendingUpload from app.scheduler import BackgroundScheduler from .conftest import write_task def wait_for(manager: RunManager, run_id: int, statuses: set[str], timeout: float = 3) -> str: deadline = time.monotonic() + timeout while time.monotonic() < deadline: run = manager.store.get(run_id) assert run is not None if run.status in statuses: return run.status time.sleep(0.01) raise AssertionError(f"run {run_id} did not reach {statuses}") def wait_for_output(path: Path, expected: str, timeout: float = 3) -> None: deadline = time.monotonic() + timeout while time.monotonic() < deadline: if path.exists() and expected in path.read_text(): return time.sleep(0.01) raise AssertionError(f"{expected!r} was not written to {path}") SLEEP_TASK = '''"""Sleeps after writing a marker.""" import os import time from pathlib import Path state = Path(os.environ["SERVER_MAINTENANCE_STATE"]) state.mkdir(parents=True, exist_ok=True) (state / "events.txt").open("a").write(f"start {Path(__file__).name}\\n") time.sleep(0.25) (state / "events.txt").open("a").write(f"end {Path(__file__).name}\\n") ''' def test_tasks_in_one_group_are_serialized_and_groups_overlap(manager: RunManager) -> None: root = manager.settings.task_root write_task(root, "first_group", "first.py", SLEEP_TASK) write_task(root, "first_group", "second.py", SLEEP_TASK) write_task(root, "second_group", "other.py", SLEEP_TASK) catalog = TaskCatalog(root) first_task = catalog.task("first_group/first.py") second_task = catalog.task("first_group/second.py") other_task = catalog.task("second_group/other.py") assert first_task is not None and second_task is not None and other_task is not None first, _ = manager.run_task(first_task) second, _ = manager.run_task(second_task) other, _ = manager.run_task(other_task) wait_for(manager, first.id, {"running"}) wait_for(manager, other.id, {"running"}) assert manager.store.get(second.id).status == "queued" assert wait_for(manager, first.id, {"succeeded"}) == "succeeded" assert wait_for(manager, second.id, {"succeeded"}) == "succeeded" assert wait_for(manager, other.id, {"succeeded"}) == "succeeded" events = (manager.settings.state_root / "tasks" / "first_group" / "events.txt").read_text() assert events.splitlines() == [ "start first.py", "end first.py", "start second.py", "end second.py", ] def test_cleanup_removes_state_for_a_removed_group_after_a_grace_period(manager: RunManager) -> None: write_task(manager.settings.task_root, "active", "task.py", 'print("active")\n') removed = manager.settings.state_root / "tasks" / "removed" removed.mkdir(parents=True) (removed / "state.json").write_text("{}\n", encoding="utf-8") manager.cleanup_transient_files(orphaned_task_state_retention_seconds=60) assert removed.is_dir() manager.cleanup_transient_files(orphaned_task_state_retention_seconds=0) assert not removed.exists() def test_group_stops_after_failed_task(manager: RunManager) -> None: root = manager.settings.task_root write_task(root, "workflow", "01_ok.py", '"""Works."""\nprint("ok")\n') write_task(root, "workflow", "02_fail.py", '"""Fails."""\nraise SystemExit(3)\n') write_task(root, "workflow", "03_never.py", '"""Must not run."""\nraise SystemExit(4)\n') group = TaskCatalog(root).group("workflow") assert group is not None parent, _ = manager.run_group(group) assert wait_for(manager, parent.id, {"failed"}) == "failed" children = manager.store.children(parent.id) assert [(child.target_path, child.status, child.exit_code) for child in children] == [ ("workflow/01_ok.py", "succeeded", 0), ("workflow/02_fail.py", "failed", 3), ] def test_duplicate_task_request_returns_the_existing_run(manager: RunManager) -> None: root = manager.settings.task_root write_task( root, "long_task", "wait.py", '"""Waits."""\nimport time\ntime.sleep(30)\n', ) task = TaskCatalog(root).task("long_task/wait.py") assert task is not None first, first_created = manager.run_task(task) second, second_created = manager.run_task(task) assert first_created is True assert second_created is False assert second.id == first.id manager.request_cancel(first.id) assert wait_for(manager, first.id, {"cancelled"}) == "cancelled" def test_cancellation_signals_the_task_process_group(manager: RunManager) -> None: root = manager.settings.task_root write_task( root, "long_task", "wait.py", '"""Waits."""\nimport time\nprint("started", flush=True)\ntime.sleep(30)\n', ) task = TaskCatalog(root).task("long_task/wait.py") assert task is not None run, _ = manager.run_task(task) wait_for(manager, run.id, {"running"}) wait_for_output(manager.settings.state_root / "runs" / f"{run.id}.log", "started") manager.request_cancel(run.id) assert wait_for(manager, run.id, {"cancelled"}) == "cancelled" completed = manager.store.get(run.id) assert completed is not None assert completed.cancel_requested_at is not None assert completed.log_path is not None assert "started" in Path(completed.log_path).read_text() def test_task_receives_validated_input_as_json_file(manager: RunManager) -> None: root = manager.settings.task_root write_task( root, "input_group", "read_input.py", '''"""Reads supplied input.""" import json import os from pathlib import Path state = Path(os.environ["SERVER_MAINTENANCE_STATE"]) data = json.loads(Path(os.environ["SERVER_MAINTENANCE_INPUT"]).read_text()) (state / "input.json").write_text(json.dumps(data, sort_keys=True)) ''', ) task = TaskCatalog(root).task("input_group/read_input.py") assert task is not None run, _ = manager.run_task(task, {"identifier": "AB123", "attempts": 2}) assert wait_for(manager, run.id, {"succeeded"}) == "succeeded" assert (manager.settings.state_root / "tasks" / "input_group" / "input.json").read_text() == ( '{"attempts": 2, "identifier": "AB123"}' ) stored = manager.store.get(run.id) assert stored is not None assert stored.input_json == '{"attempts":2,"identifier":"AB123"}' def test_task_does_not_receive_parent_environment( manager: RunManager, monkeypatch, ) -> None: monkeypatch.setenv("TASK_SETTING", "available") monkeypatch.setenv("UNRELATED_SETTING", "hidden") root = manager.settings.task_root write_task( root, "environment", "read_environment.py", '''"""Reads the task environment.""" import os from pathlib import Path state = Path(os.environ["SERVER_MAINTENANCE_STATE"]) (state / "environment.txt").write_text( f'{os.environ.get("TASK_SETTING", "")}:{os.environ.get("UNRELATED_SETTING", "")}') ''', ) task = TaskCatalog(root).task("environment/read_environment.py") assert task is not None run, _ = manager.run_task(task) assert wait_for(manager, run.id, {"succeeded"}) == "succeeded" assert (manager.settings.state_root / "tasks" / "environment" / "environment.txt").read_text() == ( ":" ) def test_task_receives_application_state_root(manager: RunManager) -> None: root = manager.settings.task_root write_task( root, "environment", "read_state_root.py", '''"""Reads the application state root.""" import os from pathlib import Path Path(os.environ["SERVER_MAINTENANCE_STATE"]).joinpath("state-root.txt").write_text( os.environ["SERVER_MAINTENANCE_STATE_ROOT"]) ''', ) task = TaskCatalog(root).task("environment/read_state_root.py") assert task is not None run, _ = manager.run_task(task) assert wait_for(manager, run.id, {"succeeded"}) == "succeeded" assert (manager.settings.state_root / "tasks" / "environment" / "state-root.txt").read_text() == str( manager.settings.state_root ) def test_sensitive_input_is_redacted_and_removed_after_execution(manager: RunManager) -> None: root = manager.settings.task_root task_path = write_task( root, "input_group", "read_secret.py", '''"""Reads a sensitive input.""" import json import os from pathlib import Path state = Path(os.environ["SERVER_MAINTENANCE_STATE"]) data = json.loads(Path(os.environ["SERVER_MAINTENANCE_INPUT"]).read_text()) (state / "secret.txt").write_text(data["source_url"]) ''', ) (task_path.parent / "task-inputs.json").write_text( '''{ "version": 1, "tasks": { "read_secret.py": { "inputs": [ {"name": "source_url", "label": "Source URL", "type": "text", "sensitive": true} ] } } }''', encoding="utf-8", ) task = TaskCatalog(root).task("input_group/read_secret.py") assert task is not None run, _ = manager.run_task(task, {"source_url": "https://example.test/?passkey=secret"}) assert wait_for(manager, run.id, {"succeeded"}) == "succeeded" assert (manager.settings.state_root / "tasks" / "input_group" / "secret.txt").read_text() == ( "https://example.test/?passkey=secret" ) stored = manager.store.get(run.id) assert stored is not None assert stored.input_json == '{"source_url":"[redacted]"}' assert not (manager.settings.state_root / "runs" / f"{run.id}.execution-input.json").exists() def test_declared_background_task_runs_without_creating_run_history(manager: RunManager) -> None: root = manager.settings.task_root task = write_task( root, "background", "internal/poll.py", '"""Internal."""\nprint("background complete")\n', ) (task.parent.parent / "background-tasks.json").write_text( '''{ "version": 1, "tasks": [ {"id": "poll", "script": "internal/poll.py", "interval_seconds": 60} ] }''', encoding="utf-8", ) scheduler = BackgroundScheduler(manager.catalog, manager) scheduler.start() log_path = manager.settings.state_root / "background" / "background" / "poll.log" wait_for_output(log_path, "background complete") scheduler.stop() assert manager.store.history("task", "background/poll") == [] def test_uploaded_file_is_available_to_task_then_removed(manager: RunManager) -> None: root = manager.settings.task_root write_task( root, "imports", "read_csv.py", '''"""Reads an uploaded CSV.""" import json import os from pathlib import Path state = Path(os.environ["SERVER_MAINTENANCE_STATE"]) values = json.loads(Path(os.environ["SERVER_MAINTENANCE_INPUT"]).read_text()) (state / "csv.txt").write_bytes(Path(values["review_csv"]).read_bytes()) ''', ) task = TaskCatalog(root).task("imports/read_csv.py") assert task is not None run, created = manager.run_task( task, uploads={ "review_csv": PendingUpload( field_name="review_csv", extension=".csv", maximum_bytes=1024, stream=BytesIO(b"answer,left_cluster_id\nsame,cluster\n"), ) }, ) assert created is True assert wait_for(manager, run.id, {"succeeded"}) == "succeeded" assert (manager.settings.state_root / "tasks" / "imports" / "csv.txt").read_bytes() == ( b"answer,left_cluster_id\nsame,cluster\n" ) assert not (manager.settings.state_root / "runs" / f"{run.id}.uploads").exists() def test_upload_is_removed_when_task_process_cannot_start(manager: RunManager) -> None: root = manager.settings.task_root write_task(root, "imports", "apply.py", '\"\"\"Applies a CSV.\"\"\"\n') task = TaskCatalog(root).task("imports/apply.py") assert task is not None broken_manager = RunManager( replace(manager.settings, python_executable="/missing/python"), TaskCatalog(root), manager.store, ) run, created = broken_manager.run_task( task, uploads={ "review_csv": PendingUpload( field_name="review_csv", extension=".csv", maximum_bytes=1024, stream=BytesIO(b"answer,left_cluster_id\n"), ) }, ) assert created is True assert wait_for(broken_manager, run.id, {"failed"}) == "failed" assert not (broken_manager.settings.state_root / "runs" / f"{run.id}.uploads").exists()