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_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_receives_only_explicitly_forwarded_environment( manager: RunManager, monkeypatch, ) -> None: monkeypatch.setenv("SERVER_MAINTENANCE_TASK_ENV", "TASK_SETTING") monkeypatch.setenv("TASK_SETTING", "available") monkeypatch.setenv("UNRELATED_SETTING", "hidden") root = manager.settings.task_root write_task( root, "environment", "read_environment.py", '''"""Reads forwarded configuration.""" 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() == ( "available:" ) 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()