Build custom container images / build (map[base_image:php:8-fpm-alpine build_args:PHP_VERSION=8
context:php8-pgsql fingerprint_command:{ apk info -v | LC_ALL=C sort; find /usr/local/lib/php/extensions /usr/local/etc/php/conf.d -type f -exec sha256sum {} + | LC_ALL=C sort; }
name:ph… (push) Successful in 43s
Build custom container images / build (map[base_image:postgres:18 build_args:PG_VERSION=18
POSTGIS_VERSION=3
VCHORD_VERSION=0.5.3
context:postgres fingerprint_command:{ dpkg-query -W -f='${binary:Package}=${Version}\n' | LC_ALL=C sort; find /usr/lib/postgresql -type f -exec sha256su… (push) Successful in 54s
Build custom container images / build (map[base_image:python:3 build_args:PYTHON_VERSION=3
context:python-tools fingerprint_command:{ dpkg-query -W -f='${binary:Package}=${Version}\n' | LC_ALL=C sort; pip freeze | LC_ALL=C sort; }
name:python-tools oci_labels:org.opencontainers.ima… (push) Successful in 1m18s
341 lines
11 KiB
Python
341 lines
11 KiB
Python
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()
|