Files
bruno 45e59009c3
FlowDeck CI / lint (push) Successful in 1m55s
FlowDeck CI / test (push) Successful in 15m23s
FlowDeck CI / docker (push) Canceled after 0s
fix: A21 phase 2c — 190 routes hors loop, 86 % total (v7.26.0)
4 passes (283 → 93 routes async sur 667 = 86 % hors loop, avant 61 %) :

A. RACINE AUTH — `get_current_user` (auth/session.py) était `async def`
   SANS aucun await (cookie decode = synchrone) ; idem ses clones :
   `agent._current_user_id/_workspace_id/_current_admin` (34 sites) et
   `sso._require_admin` (corps 0 await, 6 sites) → `def` +
   47 `await` supprimés. Piège : 3 call sites passaient par l'alias `gcu`
   (grep littéral aveugle) — 8 tests en échec → corrigés.

B. Re-scan : 19 routes devenues SANS await → `def` (agent 8, sso 5,
   web_clipper 3, projects 2, auth 1…).

C/D. 155 routes dont les seuls awaits = `request.json()` / événements :
   - try/except `body = {}` → `Body(default={})` (même tolérance)
   - try/except `raise HTTPException(400)` → `Body(...)` REQUIS
     (422 FastAPI — aucun test ne couvrait le 400)
   - forme conditionnelle `request.json() if content-type else {}`
     (54 sites) → défaut `{}` (sans corps = `{}` dans les 2 cas)
   - `await fire_*` → `run_event_sync(...)` ; imports `Body` /
     `run_event_sync` ajoutés aux routers convertis

Reste async (93, justifié) : form/upload/file (22), réseau gitea/llm/oidc,
`_json_body` (9), 2 JSON inline en argument, 1 fallback logique
(capture_frontend_error), 1 lecture conditionnelle (web_clipper), mixtes.

suite **1089/1089** · ruff OK · docs à jour
2026-10-01 22:17:48 -04:00

319 lines
12 KiB
Python

"""FlowDeck — Automations API (v5.1.0): rules CRUD, manual/button run, history."""
from __future__ import annotations
import json
import logging
from fastapi import APIRouter, Body, Depends, HTTPException, Request
from app.auth.session import SessionManager
from app.db import get_conn
from app.services.automations import (
get_page_context,
get_steps,
press_button,
run_automation,
validate_step,
)
logger = logging.getLogger(__name__)
def _require_session(request: Request) -> None:
"""A13 : toute la route (CRUD, run, press-button) exige une session."""
if not SessionManager.decode_session(request.cookies.get("flowdeck_session", "")):
raise HTTPException(status_code=401, detail="Authentication required")
router = APIRouter(tags=["automations"], dependencies=[Depends(_require_session)])
TRIGGER_TYPES = ("event", "cron", "button")
def _json_or_dumps(val, default="[]"):
"""Store JSON string columns without double-encoding."""
if val is None:
return default
if isinstance(val, str):
try:
json.loads(val)
return val
except (TypeError, json.JSONDecodeError):
return json.dumps(val)
return json.dumps(val)
def _current_user(request: Request) -> dict:
user = SessionManager.decode_session(request.cookies.get("flowdeck_session", ""))
return user if user and user.get("id") else {}
def _validate_payload(body: dict) -> None:
name = (body.get("name") or "").strip()
if not name:
raise HTTPException(status_code=400, detail="name required")
trigger_type = body.get("trigger_type", "event")
if trigger_type not in TRIGGER_TYPES:
raise HTTPException(status_code=400, detail="invalid trigger_type")
if trigger_type == "event" and not body.get("event"):
raise HTTPException(status_code=400, detail="event required for event trigger")
if trigger_type == "cron" and not (body.get("cron_expression") or "").strip():
raise HTTPException(status_code=400, detail="cron_expression required for cron trigger")
for key in ("condition_json", "actions_json"):
val = body.get(key, "[]")
try:
if isinstance(val, str):
json.loads(val)
else:
json.dumps(val)
except (TypeError, json.JSONDecodeError):
raise HTTPException(status_code=400, detail=f"{key} must be valid JSON") from None
@router.get("/workspace/automations")
def list_automations(request: Request):
with get_conn() as conn:
rows = conn.execute("SELECT * FROM automations ORDER BY created_at DESC").fetchall()
items = [dict(r) for r in rows]
return {"automations": items}
@router.post("/workspace/automations")
def create_automation(request: Request, body: dict = Body(default={})):
_validate_payload(body)
user = _current_user(request)
by = user["id"]
with get_conn() as conn:
cur = conn.execute(
"""INSERT INTO automations
(workspace, name, trigger_type, event, cron_expression, collection_id,
condition_json, actions_json, enabled, created_by)
VALUES (?,?,?,?,?,?,?,?,?,?)""",
(
body.get("workspace", "") or "",
(body.get("name") or "").strip(),
body.get("trigger_type", "event"),
body.get("event", "page.created"),
body.get("cron_expression", "") or "",
body.get("collection_id") or None,
_json_or_dumps(body.get("condition", body.get("condition_json", []))),
_json_or_dumps(body.get("actions", body.get("actions_json", []))),
int(body.get("enabled", True)),
by,
),
)
conn.commit()
new_id = cur.lastrowid
return {"id": new_id, "status": "created"}
@router.get("/workspace/automations/{auto_id}")
def get_automation(request: Request, auto_id: int):
with get_conn() as conn:
row = conn.execute("SELECT * FROM automations WHERE id=?", (auto_id,)).fetchone()
if not row:
raise HTTPException(status_code=404, detail="Automation not found")
return dict(row)
@router.put("/workspace/automations/{auto_id}")
def update_automation(request: Request, auto_id: int, body: dict = Body(default={})):
_validate_payload(body)
with get_conn() as conn:
row = conn.execute("SELECT id FROM automations WHERE id=?", (auto_id,)).fetchone()
if not row:
raise HTTPException(status_code=404, detail="Automation not found")
conn.execute(
"""UPDATE automations SET
name=?, trigger_type=?, event=?, cron_expression=?, collection_id=?,
condition_json=?, actions_json=?, enabled=?, updated_at=CURRENT_TIMESTAMP
WHERE id=?""",
(
(body.get("name") or "").strip(),
body.get("trigger_type", "event"),
body.get("event", "page.created"),
body.get("cron_expression", "") or "",
body.get("collection_id") or None,
_json_or_dumps(body.get("condition", body.get("condition_json", []))),
_json_or_dumps(body.get("actions", body.get("actions_json", []))),
int(body.get("enabled", True)),
auto_id,
),
)
conn.commit()
return {"id": auto_id, "status": "updated"}
@router.delete("/workspace/automations/{auto_id}")
def delete_automation(request: Request, auto_id: int):
with get_conn() as conn:
conn.execute("DELETE FROM automations WHERE id=?", (auto_id,))
conn.commit()
return {"id": auto_id, "status": "deleted"}
async def _execute(automation_id: int, trigger_source: str, body: dict) -> dict:
page_id = body.get("page_id") if isinstance(body, dict) else None
collection_id = body.get("collection_id") if isinstance(body, dict) else None
context = {"collection_id": collection_id, "page_id": page_id}
if page_id:
context.update(get_page_context(int(page_id), collection_id or 0))
result = await run_automation(automation_id, trigger_source, context)
result["automation_id"] = automation_id
return result
@router.post("/workspace/automations/{auto_id}/run")
async def run_automation_endpoint(request: Request, auto_id: int):
body = await request.json() if request.headers.get("content-type") else {}
return await _execute(auto_id, "manual", body)
@router.post("/api/automations/{auto_id}/run")
async def run_automation_button(request: Request, auto_id: int):
body = await request.json() if request.headers.get("content-type") else {}
return await _execute(auto_id, "button", body)
@router.get("/workspace/automations/{auto_id}/runs")
def automation_runs_history(request: Request, auto_id: int, limit: int = 50):
with get_conn() as conn:
rows = conn.execute(
"""SELECT * FROM automation_runs WHERE automation_id=?
ORDER BY created_at DESC, id DESC LIMIT ?""",
(auto_id, limit),
).fetchall()
return {"runs": [dict(r) for r in rows]}
# ── v7.0.0 — chained steps (trigger/condition/delay/action) ───────────────
STEP_SECRET_FIELDS = {"webhook_url"}
def _require_session(request: Request) -> dict:
user = SessionManager.decode_session(request.cookies.get("flowdeck_session", ""))
if not user or not user.get("id"):
raise HTTPException(status_code=401, detail="Authentication required")
return user
def _get_auto(auto_id: int) -> dict | None:
with get_conn() as conn:
row = conn.execute("SELECT * FROM automations WHERE id=?", (auto_id,)).fetchone()
return dict(row) if row else None
def _auto_404():
# NOTE: return (not raise) — the global 404 handler redirects non-/api
# paths to /workspaces, which TestClient follows into a 200.
from fastapi.responses import JSONResponse
return JSONResponse({"detail": "Automation not found"}, status_code=404)
def _encrypt_step_config(config: dict) -> dict:
"""Encrypt secret fields at rest (empty = keep existing, like sso_config)."""
from app.services.sso_provisioning import encrypt_secret
cfg = dict(config or {})
for field in STEP_SECRET_FIELDS:
if field in cfg and cfg[field]:
val = str(cfg[field])
if not val.startswith("gAAAAA"):
cfg[field] = encrypt_secret(val)
return cfg
@router.get("/workspace/automations/{auto_id}/steps")
def list_steps(request: Request, auto_id: int):
if _get_auto(auto_id) is None:
return _auto_404()
return {"automation_id": auto_id, "steps": get_steps(auto_id)}
@router.post("/workspace/automations/{auto_id}/steps")
def create_step(request: Request, auto_id: int, body: dict = Body(default={})):
_require_session(request)
if _get_auto(auto_id) is None:
return _auto_404()
kind = body.get("kind", "")
config = body.get("config", {}) or {}
validate_step(kind, config)
with get_conn() as conn:
pos = conn.execute(
"SELECT COALESCE(MAX(position), -1)+1 FROM automation_steps WHERE automation_id=?",
(auto_id,)).fetchone()[0]
cur = conn.execute(
"INSERT INTO automation_steps (automation_id, kind, position, config_json)"
" VALUES (?,?,?,?)",
(auto_id, kind, int(body.get("position", pos)),
json.dumps(_encrypt_step_config(config))))
conn.commit()
step_id = cur.lastrowid
return {"id": step_id, "status": "created"}
@router.put("/workspace/automations/steps/{step_id}")
def update_step(request: Request, step_id: int, body: dict = Body(default={})):
_require_session(request)
with get_conn() as conn:
row = conn.execute("SELECT * FROM automation_steps WHERE id=?", (step_id,)).fetchone()
if not row:
from fastapi.responses import JSONResponse
return JSONResponse({"detail": "Step not found"}, status_code=404)
kind = body.get("kind", row["kind"])
try:
config = body.get("config", json.loads(row["config_json"] or "{}"))
except (TypeError, json.JSONDecodeError):
config = {}
validate_step(kind, config if isinstance(config, dict) else {})
conn.execute(
"UPDATE automation_steps SET kind=?, position=?, config_json=? WHERE id=?",
(kind, int(body.get("position", row["position"])),
json.dumps(_encrypt_step_config(config)), step_id))
conn.commit()
return {"id": step_id, "status": "updated"}
@router.delete("/workspace/automations/steps/{step_id}")
def delete_step(request: Request, step_id: int):
_require_session(request)
with get_conn() as conn:
conn.execute("DELETE FROM automation_steps WHERE id=?", (step_id,))
conn.commit()
return {"id": step_id, "status": "deleted"}
@router.put("/workspace/automations/{auto_id}/mode")
def set_trigger_mode(request: Request, auto_id: int, body: dict = Body(default={})):
"""Set multi-trigger mode: any (default) or all (5-minute window)."""
_require_session(request)
if _get_auto(auto_id) is None:
return _auto_404()
mode = (body.get("mode") or "any").lower()
if mode not in ("any", "all"):
raise HTTPException(status_code=400, detail="mode must be any or all")
with get_conn() as conn:
conn.execute("UPDATE automations SET trigger_mode=? WHERE id=?", (mode, auto_id))
conn.commit()
return {"id": auto_id, "trigger_mode": mode}
@router.post("/api/automations/press-button")
async def press_button_endpoint(request: Request):
"""Run the automation linked to a native DB button cell (CSRF-exempt)."""
body = await request.json() if request.headers.get("content-type") else {}
try:
collection_id = int(body.get("collection_id", 0))
row_id = int(body.get("row_id", 0))
except (TypeError, ValueError):
raise HTTPException(status_code=400, detail="collection_id + row_id required") from None
prop_ref = body.get("property", body.get("property_id", ""))
if not prop_ref:
raise HTTPException(status_code=400, detail="property required")
user = _current_user(request)
try:
result = await press_button(collection_id, row_id, prop_ref, user.get("id") or 1)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from None
return result