Files
ObsiGate/backend/routers/scheduler.py
T
bruno 4ce677902a
CI / lint (push) Successful in 2m53s
CI / security (push) Successful in 2m5s
CI / test (push) Successful in 4m52s
CI / build (push) Successful in 2m5s
CI / e2e (push) Successful in 16m14s
feat: agent IA phase 4 — doublons #166, notifications externes #168 (Discord/Telegram/SMTP/webhook), taches planifiees #170
2026-10-04 13:11:25 -04:00

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