From b0c707d62461d714f0bb7d8be093c3c3533fa5f2 Mon Sep 17 00:00:00 2001 From: ajp_anton Date: Sun, 9 Aug 2026 01:24:06 +0000 Subject: [PATCH] Update Python tools --- README.md | 5 +- python-tools/app/discovery.py | 122 +++++++++++++++++++++++++++--- python-tools/app/runner.py | 51 ++++++++++--- python-tools/app/scheduler.py | 67 ++++++++++++++++ python-tools/app/web.py | 15 +++- python-tools/static/app.js | 85 +++++++++++++++++++++ python-tools/static/style.css | 36 +++++++++ python-tools/templates/index.html | 10 ++- 8 files changed, 365 insertions(+), 26 deletions(-) create mode 100644 python-tools/app/scheduler.py diff --git a/README.md b/README.md index 2d6c352..dcca7a0 100644 --- a/README.md +++ b/README.md @@ -85,8 +85,9 @@ Pinned Python packages: Included application: -- A server-rendered web control panel for discovering and running reviewed, - zero-input Python scripts mounted at runtime. +- A server-rendered web control panel for discovering and running reviewed + Python scripts mounted at runtime, with typed form inputs and optional + internal recurring jobs. - Per-group task queues, SQLite run history, captured output, and task-process cancellation. - The default command starts the control panel with `python -m app`. diff --git a/python-tools/app/discovery.py b/python-tools/app/discovery.py index 2188301..aff2590 100644 --- a/python-tools/app/discovery.py +++ b/python-tools/app/discovery.py @@ -27,6 +27,17 @@ class Task: display_name: str description: str | None inputs: tuple["InputField", ...] + wait_for_result: bool = False + + +@dataclass(frozen=True) +class BackgroundTask: + """A declared internal task that is never exposed through the web UI.""" + + id: str + group_id: str + path: Path + interval_seconds: int @dataclass(frozen=True) @@ -56,6 +67,10 @@ class Group: display_name: str tasks: tuple[Task, ...] + @property + def can_run_as_group(self) -> bool: + return all(not task.inputs for task in self.tasks) + class TaskCatalog: """Discovers only direct Python children of non-hidden group directories.""" @@ -96,6 +111,21 @@ class TaskCatalog: return task return None + def background_tasks(self) -> tuple[BackgroundTask, ...]: + if not self.task_root.is_dir(): + return () + + tasks = [] + for group_path in sorted(self.task_root.iterdir(), key=lambda path: path.name): + if ( + group_path.name.startswith((".", "_")) + or group_path.is_symlink() + or not group_path.is_dir() + ): + continue + tasks.extend(load_background_tasks(group_path)) + return tuple(tasks) + def _tasks_in_group(self, group_path: Path) -> list[Task]: task_paths = [ path @@ -107,12 +137,13 @@ class TaskCatalog: or path.suffix != ".py" ) ] - inputs_by_filename = load_input_manifest( + declarations_by_filename = load_input_manifest( group_path, {path.name for path in task_paths}, ) tasks = [] for path in task_paths: + declaration = declarations_by_filename.get(path.name) tasks.append( Task( id=f"{group_path.name}/{path.name}", @@ -121,12 +152,76 @@ class TaskCatalog: path=path, display_name=display_name(path.stem), description=read_description(path), - inputs=inputs_by_filename.get(path.name, ()), + inputs=declaration.inputs if declaration else (), + wait_for_result=declaration.wait_for_result if declaration else False, ) ) return tasks +def load_background_tasks(group_path: Path) -> list[BackgroundTask]: + """Read optional internal recurring-task declarations for one group.""" + + manifest_path = group_path / "background-tasks.json" + if not manifest_path.exists(): + return [] + if manifest_path.is_symlink(): + raise ValueError(f"Background task manifest cannot be a symlink: {manifest_path}") + try: + data = json.loads(manifest_path.read_text(encoding="utf-8")) + except (OSError, UnicodeDecodeError, json.JSONDecodeError) as error: + raise ValueError(f"Could not read background task manifest {manifest_path}: {error}") from error + if not isinstance(data, dict) or set(data) != {"version", "tasks"}: + raise ValueError( + f"Background task manifest {manifest_path} must contain only version and tasks." + ) + if data["version"] != 1 or not isinstance(data["tasks"], list): + raise ValueError(f"Background task manifest {manifest_path} has an unsupported structure.") + + tasks = [] + seen_ids = set() + for declaration in data["tasks"]: + if not isinstance(declaration, dict) or set(declaration) != {"id", "script", "interval_seconds"}: + raise ValueError(f"Background task manifest {manifest_path} has an invalid task declaration.") + task_id = declaration["id"] + script = declaration["script"] + interval_seconds = declaration["interval_seconds"] + if not isinstance(task_id, str) or not FIELD_NAME.fullmatch(task_id) or task_id in seen_ids: + raise ValueError(f"Background task manifest {manifest_path} has an invalid or duplicate task id.") + if ( + not isinstance(script, str) + or not script.startswith("internal/") + or Path(script).suffix != ".py" + or Path(script).is_absolute() + or ".." in Path(script).parts + ): + raise ValueError(f"Background task manifest {manifest_path} task {task_id} has an invalid script.") + if not isinstance(interval_seconds, int) or isinstance(interval_seconds, bool) or interval_seconds < 1: + raise ValueError(f"Background task manifest {manifest_path} task {task_id} has an invalid interval.") + path = group_path / script + try: + path.relative_to(group_path) + except ValueError as error: + raise ValueError(f"Background task manifest {manifest_path} task {task_id} escapes its group.") from error + if path.is_symlink() or not path.is_file(): + raise ValueError(f"Background task manifest {manifest_path} task {task_id} script is unavailable.") + parent = path.parent + while parent != group_path: + if parent.is_symlink(): + raise ValueError(f"Background task manifest {manifest_path} task {task_id} script uses a symlink.") + parent = parent.parent + tasks.append( + BackgroundTask( + id=f"{group_path.name}/{task_id}", + group_id=group_path.name, + path=path, + interval_seconds=interval_seconds, + ) + ) + seen_ids.add(task_id) + return tasks + + def read_description(path: Path) -> str | None: """Return the first docstring line without importing or executing a task.""" @@ -146,7 +241,7 @@ def read_description(path: Path) -> str | None: def load_input_manifest( group_path: Path, task_filenames: set[str], -) -> dict[str, tuple[InputField, ...]]: +) -> dict[str, "TaskDeclaration"]: """Read the optional data-only input declaration for one task group.""" manifest_path = group_path / "task-inputs.json" @@ -171,18 +266,24 @@ def load_input_manifest( raise ValueError(f"Input manifest {manifest_path} names unknown tasks: {names}.") return { - filename: parse_task_inputs(manifest_path, filename, declaration) + filename: parse_task_declaration(manifest_path, filename, declaration) for filename, declaration in data["tasks"].items() } -def parse_task_inputs( +@dataclass(frozen=True) +class TaskDeclaration: + inputs: tuple[InputField, ...] + wait_for_result: bool + + +def parse_task_declaration( manifest_path: Path, filename: str, declaration: Any, -) -> tuple[InputField, ...]: - if not isinstance(declaration, dict) or set(declaration) != {"inputs"}: - raise ValueError(f"{manifest_path} task {filename} must contain only inputs.") +) -> TaskDeclaration: + if not isinstance(declaration, dict) or not {"inputs"} <= set(declaration) or set(declaration) - {"inputs", "wait_for_result"}: + raise ValueError(f"{manifest_path} task {filename} has an invalid declaration.") raw_inputs = declaration["inputs"] if not isinstance(raw_inputs, list): raise ValueError(f"{manifest_path} task {filename} inputs must be a list.") @@ -190,7 +291,10 @@ def parse_task_inputs( names = [field.name for field in fields] if len(names) != len(set(names)): raise ValueError(f"{manifest_path} task {filename} has duplicate input names.") - return fields + wait_for_result = declaration.get("wait_for_result", False) + if not isinstance(wait_for_result, bool): + raise ValueError(f"{manifest_path} task {filename} wait_for_result must be true or false.") + return TaskDeclaration(inputs=fields, wait_for_result=wait_for_result) def parse_input_field(manifest_path: Path, filename: str, raw: Any) -> InputField: diff --git a/python-tools/app/runner.py b/python-tools/app/runner.py index c7195de..e3501ed 100644 --- a/python-tools/app/runner.py +++ b/python-tools/app/runner.py @@ -10,7 +10,7 @@ import threading from pathlib import Path from .config import Settings -from .discovery import Group, Task, TaskCatalog +from .discovery import BackgroundTask, Group, Task, TaskCatalog from .store import Run, RunStore @@ -60,6 +60,29 @@ class RunManager: self._terminate_active_process(run.id) return self.store.get(run.id) + def run_background_task(self, task: BackgroundTask) -> None: + """Run a declared internal task without creating high-volume run history.""" + + log_path = self.settings.state_root / "background" / f"{task.id}.log" + log_path.parent.mkdir(parents=True, exist_ok=True) + environment = self._environment(task.group_id) + try: + with log_path.open("wb") as log_file: + process = subprocess.Popen( + [self.settings.python_executable, str(task.path)], + cwd=task.path.parent, + stdin=subprocess.DEVNULL, + stdout=log_file, + stderr=subprocess.STDOUT, + env=environment, + start_new_session=True, + ) + exit_code = process.wait() + if exit_code: + log_file.write(f"Background task exited with status {exit_code}.\n".encode()) + except OSError as error: + log_path.write_text(f"Could not start background task: {error}\n", encoding="utf-8") + def _wake_worker(self, group_id: str) -> None: with self._lock: event = self._wake_events.setdefault(group_id, threading.Event()) @@ -145,8 +168,6 @@ class RunManager: ): return self.store.finish(run.id, status="cancelled", message="Cancelled before execution.") - task_state = self.settings.state_root / "tasks" / task.group_id - task_state.mkdir(parents=True, exist_ok=True) log_path = self.settings.state_root / "runs" / f"{run.id}.log" log_path.parent.mkdir(parents=True, exist_ok=True) input_path = self.settings.state_root / "runs" / f"{run.id}.input.json" @@ -159,15 +180,8 @@ class RunManager: status="failed", message=f"Could not prepare task input: {error}", ) - environment = { - "PATH": "/usr/local/bin:/usr/bin:/bin", - "PYTHONUNBUFFERED": "1", - "SERVER_MAINTENANCE_CREDENTIALS": str(self.settings.credentials_root), - "SERVER_MAINTENANCE_STATE": str(task_state), - "SERVER_MAINTENANCE_INPUT": str(input_path), - } - if self.settings.timezone_name: - environment["TZ"] = self.settings.timezone_name + environment = self._environment(task.group_id) + environment["SERVER_MAINTENANCE_INPUT"] = str(input_path) try: with log_path.open("wb") as log_file: @@ -210,6 +224,19 @@ class RunManager: log_path=str(log_path), ) + def _environment(self, group_id: str) -> dict[str, str]: + task_state = self.settings.state_root / "tasks" / group_id + task_state.mkdir(parents=True, exist_ok=True) + environment = { + "PATH": "/usr/local/bin:/usr/bin:/bin", + "PYTHONUNBUFFERED": "1", + "SERVER_MAINTENANCE_CREDENTIALS": str(self.settings.credentials_root), + "SERVER_MAINTENANCE_STATE": str(task_state), + } + if self.settings.timezone_name: + environment["TZ"] = self.settings.timezone_name + return environment + def _terminate_active_process(self, run_id: int) -> None: with self._lock: process = self._processes.get(run_id) diff --git a/python-tools/app/scheduler.py b/python-tools/app/scheduler.py new file mode 100644 index 0000000..0a8a4fd --- /dev/null +++ b/python-tools/app/scheduler.py @@ -0,0 +1,67 @@ +"""Internal recurring-task scheduler.""" + +from __future__ import annotations + +import threading +import time + +from .discovery import BackgroundTask, TaskCatalog +from .runner import RunManager + + +class BackgroundScheduler: + """Starts each declared background task at most once per configured interval.""" + + def __init__(self, catalog: TaskCatalog, manager: RunManager) -> None: + self.catalog = catalog + self.manager = manager + self.tasks = catalog.background_tasks() + self._stopped = threading.Event() + self._lock = threading.Lock() + self._next_due: dict[str, float] = {} + self._running: set[str] = set() + self._thread: threading.Thread | None = None + + def start(self) -> None: + if self._thread is not None: + return + self._thread = threading.Thread( + target=self._loop, + name="background-task-scheduler", + daemon=True, + ) + self._thread.start() + + def stop(self) -> None: + self._stopped.set() + if self._thread is not None: + self._thread.join(timeout=1) + + def _loop(self) -> None: + while not self._stopped.is_set(): + current = time.monotonic() + for task in self.tasks: + due = self._next_due.setdefault(task.id, current) + if current >= due: + self._next_due[task.id] = current + task.interval_seconds + self._start_task(task) + self._stopped.wait(0.25) + + def _start_task(self, task: BackgroundTask) -> None: + with self._lock: + if task.id in self._running: + return + self._running.add(task.id) + threading.Thread( + target=self._run_task, + args=(task,), + name=f"background-task-{task.id}", + daemon=True, + ).start() + + def _run_task(self, task: BackgroundTask) -> None: + try: + self.manager.run_background_task(task) + finally: + with self._lock: + self._running.discard(task.id) diff --git a/python-tools/app/web.py b/python-tools/app/web.py index e61b859..e8b94d3 100644 --- a/python-tools/app/web.py +++ b/python-tools/app/web.py @@ -12,6 +12,7 @@ from .config import Settings from .discovery import TaskCatalog from .inputs import validate_inputs from .runner import RunManager +from .scheduler import BackgroundScheduler from .store import RunStore @@ -22,6 +23,8 @@ def create_app(settings: Settings | None = None) -> Flask: store.initialize() store.mark_interrupted() manager = RunManager(settings, catalog, store) + scheduler = BackgroundScheduler(catalog, manager) + scheduler.start() project_root = Path(__file__).resolve().parent.parent app = Flask( @@ -33,6 +36,7 @@ def create_app(settings: Settings | None = None) -> Flask: app.extensions["catalog"] = catalog app.extensions["run_manager"] = manager app.extensions["run_store"] = store + app.extensions["background_scheduler"] = scheduler @app.before_request def allow_only_proxy() -> None: @@ -92,18 +96,27 @@ def create_app(settings: Settings | None = None) -> Flask: } input_data, errors = validate_inputs(task.inputs, submitted_values) if errors: + if request.accept_mimetypes.best == "application/json": + return {"errors": errors}, 400 return render_index( form_values={task.id: submitted_values}, input_errors={task.id: errors}, status=400, ) run, _ = manager.run_task(task, input_data) + if request.accept_mimetypes.best == "application/json": + return { + "detail_url": url_for("run_detail", run_id=run.id), + "output_url": url_for("run_output", run_id=run.id), + "run_id": run.id, + "state_url": url_for("run_state", run_id=run.id), + }, 202 return redirect(url_for("run_detail", run_id=run.id)) @app.post("/groups//run") def run_group(group_id: str): group = catalog.group(group_id) - if group is None: + if group is None or not group.can_run_as_group: abort(404) run, _ = manager.run_group(group) return redirect(url_for("run_detail", run_id=run.id)) diff --git a/python-tools/static/app.js b/python-tools/static/app.js index 223a2b3..fa6e3f8 100644 --- a/python-tools/static/app.js +++ b/python-tools/static/app.js @@ -17,3 +17,88 @@ if (output) { refresh(); window.setInterval(refresh, 2000); } + +for (const form of document.querySelectorAll("form[data-wait-for-result]")) { + const result = form.querySelector(".task-result"); + const controls = [...form.querySelectorAll("input, select, button")]; + let waiting = false; + + const warnBeforeLeaving = (event) => { + event.preventDefault(); + event.returnValue = ""; + }; + + const setWaiting = (value) => { + waiting = value; + form.classList.toggle("is-waiting", value); + controls.forEach((control) => { + control.disabled = value; + }); + if (value) { + window.addEventListener("beforeunload", warnBeforeLeaving); + } else { + window.removeEventListener("beforeunload", warnBeforeLeaving); + } + }; + + const showResult = (message, kind, detailUrl) => { + result.replaceChildren(); + result.className = `task-result ${kind}`; + result.append(message); + if (detailUrl) { + const link = document.createElement("a"); + link.href = detailUrl; + link.textContent = "View details"; + result.append(" ", link); + } + }; + + const waitForRun = async (run) => { + while (true) { + const response = await fetch(run.state_url, { cache: "no-store" }); + if (!response.ok) { + throw new Error("Could not read task status."); + } + const state = await response.json(); + if (state.status !== "queued" && state.status !== "running") { + if (state.status === "succeeded") { + showResult("Completed successfully.", "succeeded", run.detail_url); + } else { + showResult(`Finished with status: ${state.status}.`, "failed", run.detail_url); + } + return; + } + await new Promise((resolve) => window.setTimeout(resolve, 500)); + } + }; + + form.addEventListener("submit", async (event) => { + if (waiting) { + event.preventDefault(); + return; + } + if (!form.reportValidity()) { + return; + } + event.preventDefault(); + setWaiting(true); + showResult("Running...", "running"); + try { + const response = await fetch(form.action, { + method: "POST", + body: new FormData(form), + headers: { Accept: "application/json" }, + }); + const data = await response.json(); + if (!response.ok) { + const messages = data.errors ? Object.values(data.errors).join(" ") : "Could not start task."; + throw new Error(messages); + } + await waitForRun(data); + } catch (error) { + showResult(error instanceof Error ? error.message : "Could not run task.", "failed"); + } finally { + setWaiting(false); + } + }); +} diff --git a/python-tools/static/style.css b/python-tools/static/style.css index 4bfc3be..a5aa56c 100644 --- a/python-tools/static/style.css +++ b/python-tools/static/style.css @@ -97,6 +97,42 @@ button { padding: .4rem .7rem; } +button:disabled { + cursor: wait; +} + +.spinner { + border: .15em solid currentColor; + border-right-color: transparent; + border-radius: 50%; + display: none; + height: .75em; + margin-left: .35em; + vertical-align: -.05em; + width: .75em; +} + +.is-waiting .spinner { + animation: spin .7s linear infinite; + display: inline-block; +} + +.task-result { + margin: .5rem 0 0; +} + +.task-result.succeeded { + color: #167c3a; +} + +.task-result.failed { + color: #b00020; +} + +@keyframes spin { + to { transform: rotate(360deg); } +} + dt { font-weight: 700; } diff --git a/python-tools/templates/index.html b/python-tools/templates/index.html index 9504605..6c10fc7 100644 --- a/python-tools/templates/index.html +++ b/python-tools/templates/index.html @@ -9,9 +9,11 @@

{{ group.display_name }}

+ {% if group.can_run_as_group %}
+ {% endif %}
{% if group_run %}

@@ -35,7 +37,7 @@

{% endif %} -
+ {% set values = form_values.get(task.id, {}) %} {% set errors = input_errors.get(task.id, {}) %} {% for field in task.inputs %} @@ -79,7 +81,11 @@ {% if errors.get(field.name) %}

{{ errors[field.name] }}

{% endif %} {% endfor %} - + + {% if task.wait_for_result %}

{% endif %}
{% else %}