- Conversion `async def` → `def` de TOUTES les routes dont le corps ne contient ni `await`, ni `async with`, ni `async for`, ni `asyncio` (scan automatique corps par corps sur app/ : 352 converties, 0 dangereuses, vérifié `asyncio`/`run_coroutine`/`.result()` absents). FastAPI exécute ces handlers dans son threadpool → tout leur SQLite (`get_conn()` + `conn.execute`) quitte l'event loop, sans changer une ligne de logique. - Répartition : api_v2 60, dashboard 40, collections 25, board 23, workspace 19, wiki 17, permissions 14, api 14, main.py 6, + 35 fichiers. - Les 4 routers prioritaires de l'audit sont couverts par ce lot : api_v2 60 + dashboard 40 + collections 25 + board 23 = 148 conversions (le reste de leurs routes attend la phase 2 : elles ont de vrais `await`). - Reste (phase 2) : les 311 routes avec de vrais `await` → enrouler les blocs DB dans `await anyio.to_thread.run_sync(...)` ; pas de wrapper partagé livré (rien ne l'appellerait — YAGNI jusqu'au premier usage). suite **1037/1037** (233 s) · `ruff check app tests` OK · docs à jour
217 lines
8.8 KiB
Python
217 lines
8.8 KiB
Python
"""FlowDeck — Workers API (v7.0.0): CRUD, manual run, history, fork, usage."""
|
|
from __future__ import annotations
|
|
|
|
from fastapi import APIRouter, HTTPException, Request
|
|
from fastapi.responses import JSONResponse
|
|
|
|
from app.auth.session import SessionManager
|
|
from app.db import get_conn
|
|
from app.services import workers as worker_service
|
|
from app.services.api_v2_helpers import (
|
|
audit_log,
|
|
has_scope,
|
|
paginate_headers,
|
|
parse_pagination,
|
|
resolve_bearer_token,
|
|
row_to_dict,
|
|
)
|
|
|
|
router = APIRouter(tags=["workers"])
|
|
|
|
|
|
def _auth_user(request: Request, *, require_write: bool = False) -> dict:
|
|
sess = SessionManager.decode_session(request.cookies.get("flowdeck_session", ""))
|
|
if sess:
|
|
return sess
|
|
auth = request.headers.get("authorization") or request.headers.get("Authorization") or ""
|
|
if auth.lower().startswith("bearer "):
|
|
user = resolve_bearer_token(auth[7:].strip())
|
|
if not user:
|
|
raise HTTPException(401, "Invalid or expired API token")
|
|
if require_write and not has_scope(user.get("_token_scopes") or "read", "write"):
|
|
raise HTTPException(403, "Insufficient scope. Required: write")
|
|
return user
|
|
raise HTTPException(401, "Authentication required")
|
|
|
|
|
|
def _row_to_api(row) -> dict:
|
|
d = row_to_dict(row)
|
|
d.pop("code_py", None) # code only via ?include_code=1 or owner fetch
|
|
return d
|
|
|
|
|
|
@router.post("/api/v2/workers")
|
|
async def create_worker(request: Request):
|
|
user = _auth_user(request, require_write=True)
|
|
try:
|
|
body = await request.json()
|
|
except Exception:
|
|
body = {}
|
|
name = (body.get("name") or "Untitled worker").strip()[:200]
|
|
code = body.get("code_py") or ""
|
|
try:
|
|
worker_service.validate_code(code)
|
|
except worker_service.WorkerRejected as exc:
|
|
raise HTTPException(400, f"code rejected: {exc}") from None
|
|
slug = worker_service.unique_slug(body.get("slug") or name)
|
|
with get_conn() as conn:
|
|
cur = conn.execute(
|
|
"""INSERT INTO workers (slug, workspace_id, name, code_py, schedule_cron,
|
|
shared, daily_budget_s, created_by)
|
|
VALUES (?,?,?,?,?,?,?,?)""",
|
|
(slug, body.get("workspace_id"), name, code,
|
|
(body.get("schedule_cron") or "")[:60],
|
|
1 if body.get("shared") else 0,
|
|
max(1, min(int(body.get("daily_budget_s") or 60), 3600)),
|
|
user["id"]))
|
|
conn.commit()
|
|
wid = cur.lastrowid
|
|
row = conn.execute("SELECT * FROM workers WHERE id=?", (wid,)).fetchone()
|
|
audit_log(user, "worker.create", "worker", wid, slug, request)
|
|
return JSONResponse(status_code=201, content={**_row_to_api(row), "code_py": code})
|
|
|
|
|
|
@router.get("/api/v2/workers")
|
|
def list_workers(request: Request):
|
|
_auth_user(request)
|
|
limit, offset = parse_pagination(request)
|
|
with get_conn() as conn:
|
|
total = conn.execute("SELECT COUNT(*) FROM workers").fetchone()[0]
|
|
rows = conn.execute(
|
|
"SELECT * FROM workers ORDER BY id DESC LIMIT ? OFFSET ?",
|
|
(limit, offset)).fetchall()
|
|
resp = JSONResponse([_row_to_api(r) for r in rows])
|
|
for k, v in paginate_headers(total).items():
|
|
resp.headers[k] = v
|
|
return resp
|
|
|
|
|
|
@router.get("/api/v2/workers/{worker_id}")
|
|
def get_worker(worker_id: int, request: Request):
|
|
user = _auth_user(request)
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, "Worker not found")
|
|
out = _row_to_api(row)
|
|
if (request.query_params.get("include_code") == "1" or row["created_by"] == user["id"]
|
|
or user.get("is_admin")):
|
|
out["code_py"] = row["code_py"]
|
|
return out
|
|
|
|
|
|
@router.patch("/api/v2/workers/{worker_id}")
|
|
async def update_worker(worker_id: int, request: Request):
|
|
user = _auth_user(request, require_write=True)
|
|
try:
|
|
body = await request.json()
|
|
except Exception:
|
|
body = {}
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, "Worker not found")
|
|
if row["created_by"] != user["id"] and not user.get("is_admin"):
|
|
raise HTTPException(403, "Only the owner can update this worker")
|
|
updates: dict = {}
|
|
if "name" in body:
|
|
updates["name"] = str(body["name"] or "")[:200]
|
|
if "code_py" in body:
|
|
try:
|
|
worker_service.validate_code(body["code_py"] or "")
|
|
except worker_service.WorkerRejected as exc:
|
|
raise HTTPException(400, f"code rejected: {exc}") from None
|
|
updates["code_py"] = body["code_py"] or ""
|
|
if "schedule_cron" in body:
|
|
updates["schedule_cron"] = str(body["schedule_cron"] or "")[:60]
|
|
if "shared" in body:
|
|
updates["shared"] = 1 if body["shared"] else 0
|
|
if "daily_budget_s" in body:
|
|
updates["daily_budget_s"] = max(1, min(int(body["daily_budget_s"] or 60), 3600))
|
|
if "slug" in body and body["slug"] != row["slug"]:
|
|
if not worker_service._SLUG_RE.match(str(body["slug"] or "")):
|
|
raise HTTPException(400, "Invalid slug")
|
|
if conn.execute("SELECT id FROM workers WHERE slug=? AND id!=?",
|
|
(body["slug"], worker_id)).fetchone():
|
|
raise HTTPException(409, "Slug already taken")
|
|
updates["slug"] = body["slug"]
|
|
if updates:
|
|
sets = ", ".join(f"{k}=?" for k in updates)
|
|
conn.execute(f"UPDATE workers SET {sets}, updated_at=CURRENT_TIMESTAMP WHERE id=?",
|
|
(*updates.values(), worker_id))
|
|
conn.commit()
|
|
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
|
|
audit_log(user, "worker.update", "worker", worker_id, ",".join(updates), request)
|
|
return _row_to_api(row)
|
|
|
|
|
|
@router.delete("/api/v2/workers/{worker_id}")
|
|
def delete_worker(worker_id: int, request: Request):
|
|
user = _auth_user(request, require_write=True)
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, "Worker not found")
|
|
if row["created_by"] != user["id"] and not user.get("is_admin"):
|
|
raise HTTPException(403, "Only the owner can delete this worker")
|
|
conn.execute("DELETE FROM workers WHERE id=?", (worker_id,))
|
|
conn.commit()
|
|
audit_log(user, "worker.delete", "worker", worker_id, "", request)
|
|
return {"status": "deleted", "id": worker_id}
|
|
|
|
|
|
@router.post("/api/v2/workers/{worker_id}/run")
|
|
async def run_worker_endpoint(worker_id: int, request: Request):
|
|
user = _auth_user(request, require_write=True)
|
|
try:
|
|
body = await request.json() if request.headers.get("content-type") else {}
|
|
except Exception:
|
|
body = {}
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, "Worker not found")
|
|
if (row["created_by"] != user["id"] and not row["shared"]
|
|
and not user.get("is_admin")):
|
|
raise HTTPException(403, "Worker is private")
|
|
import asyncio
|
|
loop = asyncio.get_running_loop()
|
|
try:
|
|
out = await loop.run_in_executor(
|
|
None, worker_service.run_worker, worker_id, body.get("ctx") or {})
|
|
except HTTPException:
|
|
raise
|
|
audit_log(user, "worker.run", "worker", worker_id, out.get("status", ""), request)
|
|
return out
|
|
|
|
|
|
@router.get("/api/v2/workers/{worker_id}/runs")
|
|
def worker_runs(worker_id: int, request: Request):
|
|
_auth_user(request)
|
|
limit, _offset = parse_pagination(request, default_limit=20)
|
|
with get_conn() as conn:
|
|
if not conn.execute("SELECT id FROM workers WHERE id=?", (worker_id,)).fetchone():
|
|
raise HTTPException(404, "Worker not found")
|
|
rows = conn.execute(
|
|
"SELECT * FROM worker_runs WHERE worker_id=? ORDER BY id DESC LIMIT ?",
|
|
(worker_id, limit)).fetchall()
|
|
return {"worker_id": worker_id, "runs": [row_to_dict(r) for r in rows]}
|
|
|
|
|
|
@router.post("/api/v2/workers/{worker_id}/fork")
|
|
def fork_worker_endpoint(worker_id: int, request: Request):
|
|
user = _auth_user(request, require_write=True)
|
|
out = worker_service.fork_worker(worker_id, user["id"])
|
|
audit_log(user, "worker.fork", "worker", worker_id, "", request)
|
|
return JSONResponse(status_code=201, content=out)
|
|
|
|
|
|
@router.get("/api/v2/workers-usage")
|
|
def workers_usage(request: Request):
|
|
user = _auth_user(request)
|
|
ws_raw = request.query_params.get("workspace_id")
|
|
wid = int(ws_raw) if ws_raw and str(ws_raw).isdigit() else None
|
|
return {"workspace_id": wid,
|
|
"used_seconds_today": round(worker_service.daily_usage_s(wid), 2),
|
|
"user_id": user.get("id")}
|