119 lines
4.5 KiB
Python
119 lines
4.5 KiB
Python
"""Scheduled tasks endpoints (#170).
|
|
|
|
Tasks reuse the existing mutation/notification services — this router only
|
|
validates, persists and triggers. File-writing actions check vault access
|
|
at creation time; the background tick re-checks nothing (system context) but
|
|
records failures and notifies on ``schedule_failure`` (#168).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from typing import Any
|
|
|
|
from fastapi import APIRouter, Body, Depends, HTTPException
|
|
from pydantic import BaseModel, ConfigDict, Field
|
|
|
|
from backend import scheduler as _scheduler
|
|
from backend.auth.middleware import check_vault_access, require_auth
|
|
from backend.schemas import StatusResponse
|
|
|
|
router = APIRouter(prefix="/api/scheduler", tags=["scheduler"])
|
|
|
|
|
|
class ScheduledTask(BaseModel):
|
|
"""A programmed automatic task."""
|
|
|
|
model_config = ConfigDict(extra="allow")
|
|
id: str = Field(description="Task id")
|
|
name: str = Field(description="Display name")
|
|
action: dict[str, Any] = Field(description="{kind, params}")
|
|
schedule: dict[str, Any] = Field(description="{kind, ...}")
|
|
enabled: bool = Field(description="Whether the tick executes it")
|
|
created_by: str = Field(description="Owner username")
|
|
created_at: str = Field(description="ISO-8601 creation time")
|
|
last_run_at: str | None = Field(default=None)
|
|
last_status: str | None = Field(default=None)
|
|
last_error: str | None = Field(default=None)
|
|
run_count: int = Field(default=0)
|
|
next_run_at: str = Field(description="ISO-8601 next due time")
|
|
|
|
|
|
class TaskRunResult(BaseModel):
|
|
"""Outcome of a manual or due run."""
|
|
|
|
model_config = ConfigDict(extra="allow")
|
|
task_id: str = Field(description="Task id")
|
|
ok: bool = Field(description="True on success")
|
|
result: dict[str, Any] | None = Field(default=None)
|
|
error: str | None = Field(default=None)
|
|
|
|
|
|
def _check_action_vault(action: dict[str, Any], user: dict[str, Any]) -> None:
|
|
from backend.services.errors import ServiceError
|
|
from backend.services.vaults import get_vault_root
|
|
|
|
kind = (action or {}).get("kind")
|
|
params = (action or {}).get("params") or {}
|
|
if kind in ("create_file", "append_to_file"):
|
|
vault = str(params.get("vault") or "")
|
|
if not check_vault_access(vault, user):
|
|
raise HTTPException(403, f"No access to vault '{vault}'")
|
|
try:
|
|
get_vault_root(vault)
|
|
except ServiceError as e:
|
|
raise HTTPException(404, f"Unknown vault '{vault}'") from e
|
|
|
|
|
|
@router.get("/tasks", response_model=list[ScheduledTask])
|
|
async def api_scheduler_list(current_user: dict[str, Any] = Depends(require_auth)):
|
|
"""List scheduled tasks (newest first)."""
|
|
return _scheduler.list_tasks()
|
|
|
|
|
|
@router.post("/tasks", response_model=ScheduledTask)
|
|
async def api_scheduler_create(body: dict = Body(...), current_user: dict[str, Any] = Depends(require_auth)):
|
|
"""Create a task (``{name, action, schedule, enabled?}``)."""
|
|
action = dict(body.get("action") or {})
|
|
_check_action_vault(action, current_user)
|
|
try:
|
|
return _scheduler.create_task(
|
|
str(body.get("name") or ""),
|
|
action,
|
|
dict(body.get("schedule") or {}),
|
|
created_by=str(current_user.get("username", "api")),
|
|
enabled=bool(body.get("enabled", True)),
|
|
)
|
|
except ValueError as e:
|
|
raise HTTPException(400, str(e)) from e
|
|
|
|
|
|
@router.patch("/tasks/{task_id}", response_model=ScheduledTask)
|
|
async def api_scheduler_update(task_id: str, body: dict = Body(...), current_user: dict[str, Any] = Depends(require_auth)):
|
|
"""Update a task (name / enabled / action / schedule)."""
|
|
if "action" in body:
|
|
_check_action_vault(dict(body["action"] or {}), current_user)
|
|
try:
|
|
result = _scheduler.update_task(task_id, body)
|
|
except ValueError as e:
|
|
raise HTTPException(400, str(e)) from e
|
|
if result is None:
|
|
raise HTTPException(404, "Task not found")
|
|
return result
|
|
|
|
|
|
@router.delete("/tasks/{task_id}", response_model=StatusResponse)
|
|
async def api_scheduler_delete(task_id: str, current_user: dict[str, Any] = Depends(require_auth)):
|
|
"""Delete a task."""
|
|
if not _scheduler.delete_task(task_id):
|
|
raise HTTPException(404, "Task not found")
|
|
return {"status": "deleted"}
|
|
|
|
|
|
@router.post("/tasks/{task_id}/run", response_model=TaskRunResult)
|
|
async def api_scheduler_run(task_id: str, current_user: dict[str, Any] = Depends(require_auth)):
|
|
"""Execute a task immediately (manual run)."""
|
|
try:
|
|
return _scheduler.run_task(task_id, manual=True)
|
|
except KeyError:
|
|
raise HTTPException(404, "Task not found") from None
|