Files
flowdeck/app/routers/automations.py
T
bruno ba363eaee9
FlowDeck CI / lint (push) Successful in 43s
FlowDeck CI / test (push) Successful in 4m2s
FlowDeck CI / lint (pull_request) Successful in 42s
FlowDeck CI / test (pull_request) Successful in 4m3s
FlowDeck CI / docker (push) Successful in 1m2s
FlowDeck CI / docker (pull_request) Successful in 35s
feat(v5.2.0): finalize Infrastructure & Polish (tests isolation, xdist, lint, CI)
tests/conftest.py: mutate the settings singleton (instead of rebinding) so DB + backup dir are isolated per test -> pytest-xdist safe.
Real backup tests (snapshot/prune/admin API) and OAuth mock tests (Gitea/GitHub/link) replace the previous skips.
init_db() now also creates webhook_subscriptions (full schema without the FastAPI lifespan).
ruff check is clean; .eslintrc.json migrated to eslint.config.mjs (flat config).
CI: lint job (ruff + eslint), parallel tests (-n auto), run on every branch push.
VERSION 5.11.1.
2026-09-11 23:36:53 -04:00

175 lines
6.8 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, HTTPException, Request
from app.auth.session import SessionManager
from app.db import get_conn
from app.services.automations import get_page_context, run_automation
logger = logging.getLogger(__name__)
router = APIRouter(tags=["automations"])
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")
async 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")
async def create_automation(request: Request):
body = await request.json() if request.headers.get("content-type") else {}
_validate_payload(body)
user = _current_user(request)
by = user.get("id") or 1
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}")
async 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}")
async def update_automation(request: Request, auto_id: int):
body = await request.json() if request.headers.get("content-type") else {}
_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}")
async 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")
async 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]}