- A12 — `og_fetcher` : GET sans `follow_redirects`, `_is_public_host` revérifié à chaque saut (max 5) ; `POST /board/api/og/metadata` → 400 sur hôte privé/loopback - A13 — router automations sous `Depends(_require_session)` (CRUD, run, press-button) + `created_by` sans fallback ; action `webhook` validée par `_is_public_host` avant POST (SSRF) - A15 — webhooks sortants : `_require_admin` sur GET/POST/DELETE + `_is_public_host` sur l'URL en création - A17 — router legacy `/api` sous `Depends(_require_session_or_bearer)` (session ou Bearer `/api/v1`), allowlist explicite `/api/health` + `/api/frontend-error` - A22 — les 2 uploads locales : session exigée (`_require_user_id`) + `validate_upload` branché (taille + extension) + `FLOWDECK_DATA_DIR` au lieu de `/data` codé en dur - A23 — N+1 : COUNT→`GROUP BY` (dashboard), cards→`executemany` (board sync), duplicata de propriétés→`executemany` + remap des ids par SELECT (collections) - A24 — 2 routes écrasées supprimées : `GET /api/projects` (api.py) et `GET /workspace` (workspace.py) + test « aucun doublon méthode+chemin » - Tests : +9 dans `tests/test_audit_p0_fixes.py` (SSRF, 401s, validate_upload, doublons de routes) ; tests OG sur hôtes résolubles (la garde fait du DNS) - suite **1025/1025** · `ruff check app tests` OK
836 lines
35 KiB
Python
836 lines
35 KiB
Python
"""FlowDeck — Automations engine (v5.1.0).
|
|
|
|
Implements the "if-this-then-that" rule engine: automations match an event (or a
|
|
cron schedule, or a clickable button), optionally guard on a condition, then run
|
|
a list of actions.
|
|
|
|
Condition clauses (``condition_json``), all combined with AND:
|
|
{"property": "Status", "op": "eq", "value": "Done"}
|
|
{"property": "Priority", "op": "not_contains", "value": "Low"}
|
|
{"property": "Assignee", "op": "is_empty"}
|
|
{"property": "Estimate", "op": "changed"} (only event triggers)
|
|
Flags:
|
|
op in {eq, neq, contains, not_contains, is_empty, is_not_empty, changed}
|
|
|
|
Actions (``actions_json``), executed sequentially:
|
|
{"type": "webhook", "url": "...", "secret": "..."}
|
|
{"type": "set_property", "property": "Status", "value": "Done"}
|
|
{"type": "create_page", "collection_id": 3, "title": "...", "properties": {...}}
|
|
{"type": "notify", "message": "Automation fired"}
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import time
|
|
from datetime import datetime, timedelta
|
|
|
|
import httpx
|
|
|
|
from app.db import get_conn
|
|
from app.services import notifications
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
COND_OPS = {"eq", "neq", "contains", "not_contains", "is_empty", "is_not_empty", "changed"}
|
|
# Events that fire on collection pages (payload carries a `properties` dict).
|
|
PAGE_PROP_EVENTS = {"page.created", "page.updated", "page.deleted"}
|
|
|
|
|
|
def _prop_value(props: dict, key) -> tuple[bool, object]:
|
|
"""Resolve a property value by id or name. Returns ``(found, value)``.
|
|
|
|
``props`` may be keyed by property id (FlowDeckDB UI) or name (agent / API).
|
|
"""
|
|
if props is None:
|
|
return False, None
|
|
if key is None:
|
|
return True, None
|
|
skey = str(key)
|
|
if skey in props:
|
|
return True, props[skey]
|
|
if isinstance(key, int) and str(key) in props:
|
|
return True, props[str(key)]
|
|
return False, None
|
|
|
|
|
|
def match_condition_props(props: dict, before_props: dict | None, clause: dict) -> bool:
|
|
"""Evaluate a single condition clause against page property values."""
|
|
op = clause.get("op", "eq")
|
|
if op not in COND_OPS:
|
|
return False
|
|
if op == "changed":
|
|
key = clause.get("property")
|
|
if before_props is None:
|
|
return False
|
|
found_before, before_val = _prop_value(before_props, key)
|
|
found_after, after_val = _prop_value(props, key)
|
|
return found_before and found_after and before_val != after_val
|
|
|
|
found, val = _prop_value(props, clause.get("property"))
|
|
|
|
if op == "is_empty":
|
|
if not found:
|
|
return True
|
|
return val is None or str(val).strip() == ""
|
|
if op == "is_not_empty":
|
|
return found and val is not None and str(val).strip() != ""
|
|
|
|
if not found:
|
|
return False
|
|
want = clause.get("value")
|
|
if op == "eq":
|
|
return _norm(val) == _norm(want)
|
|
if op == "neq":
|
|
return _norm(val) != _norm(want)
|
|
if op == "contains":
|
|
return _norm(want) in _norm(val) if _norm(val) else False
|
|
if op == "not_contains":
|
|
return _norm(want) not in _norm(val) if _norm(val) else True
|
|
return False
|
|
|
|
|
|
def _norm(v) -> str:
|
|
if v is None:
|
|
return ""
|
|
if isinstance(v, (list, dict)):
|
|
return json.dumps(v)
|
|
return str(v)
|
|
|
|
|
|
def evaluate_conditions(condition_json, props: dict | None, before_props: dict | None = None) -> bool:
|
|
"""Evaluate the stored condition list (AND of all clauses). Empty list → True."""
|
|
try:
|
|
clauses = json.loads(condition_json) if isinstance(condition_json, str) else (condition_json or [])
|
|
except (TypeError, json.JSONDecodeError):
|
|
clauses = []
|
|
for clause in clauses or []:
|
|
if not match_condition_props(props, before_props, clause):
|
|
return False
|
|
return True
|
|
|
|
|
|
def get_page_context(page_id: int, collection_id: int) -> dict:
|
|
"""Load a collection page's property values for condition evaluation."""
|
|
with get_conn() as conn:
|
|
row = conn.execute(
|
|
"SELECT id, title, icon, property_values_json FROM collection_pages WHERE id=?",
|
|
(page_id,),
|
|
).fetchone()
|
|
if not row:
|
|
return {"page_id": page_id, "collection_id": collection_id,
|
|
"title": "", "properties": {}, "icon": "file"}
|
|
try:
|
|
props = json.loads(row["property_values_json"])
|
|
except (TypeError, json.JSONDecodeError):
|
|
props = {}
|
|
return {"page_id": page_id, "collection_id": collection_id,
|
|
"title": row["title"], "icon": row["icon"], "properties": props}
|
|
|
|
|
|
def _save_run(automation_id: int, trigger_source: str, status: str, detail: str,
|
|
collection_id: int | None = None, page_id: int | None = None) -> None:
|
|
with get_conn() as conn:
|
|
conn.execute(
|
|
"""INSERT INTO automation_runs
|
|
(automation_id, trigger_source, status, detail, collection_id, page_id)
|
|
VALUES (?,?,?,?,?,?)""",
|
|
(automation_id, trigger_source, status, detail, collection_id, page_id),
|
|
)
|
|
conn.execute(
|
|
"UPDATE automations SET run_count=run_count+1, last_run_at=CURRENT_TIMESTAMP WHERE id=?",
|
|
(automation_id,),
|
|
)
|
|
conn.commit()
|
|
|
|
|
|
def _maybe_convert_prediction(value, props: dict) -> tuple[bool, object]:
|
|
"""Allow action values to interpolate other page properties: e.g. [[Assignee]] or {{title}}."""
|
|
if not isinstance(value, str):
|
|
return True, value
|
|
replaced = value
|
|
for key in props:
|
|
if "[[" + str(key) + "]]" in replaced:
|
|
replaced = replaced.replace("[[" + str(key) + "]]", str(props[key]))
|
|
if "{{title}}" in replaced:
|
|
replaced = replaced.replace("{{title}}", str(props.get("title", "")))
|
|
if "{{id}}" in replaced:
|
|
replaced = replaced.replace("{{id}}", str(props.get("page_id", "")))
|
|
return True, replaced
|
|
|
|
|
|
async def _run_action(action: dict, context: dict, trigger_source: str) -> str:
|
|
"""Execute a single action. Returns a human summary. Raises on failure."""
|
|
atype = action.get("type")
|
|
|
|
if atype == "webhook":
|
|
url = action.get("url", "").strip()
|
|
if not url:
|
|
raise ValueError("webhook action requires a url")
|
|
# A13 : SSRF — même garde que l'importer URL (loopback/privé refusé).
|
|
from urllib.parse import urlparse as _urlparse
|
|
|
|
from app.services.importers.url_fetch import _is_public_host
|
|
_parsed = _urlparse(url)
|
|
if _parsed.scheme not in ("http", "https") or not _parsed.hostname or not _is_public_host(_parsed.hostname):
|
|
raise ValueError(f"webhook url non autorisée: {_parsed.hostname!r}")
|
|
secret = action.get("secret", "")
|
|
headers = {"Content-Type": "application/json", "X-FlowDeck-Event": context.get("event", "")}
|
|
if secret:
|
|
headers["X-FlowDeck-Secret"] = secret
|
|
async with httpx.AsyncClient(timeout=10) as client:
|
|
resp = await client.post(url, json=context, headers=headers)
|
|
if resp.status_code >= 400:
|
|
raise RuntimeError(f"webhook returned HTTP {resp.status_code}")
|
|
return f"webhook → {url} ({resp.status_code})"
|
|
|
|
if atype == "set_property":
|
|
prop = action.get("property")
|
|
value = action.get("value")
|
|
page_id = context.get("page_id")
|
|
if not prop or not page_id:
|
|
raise ValueError("set_property requires property + page context")
|
|
_, resolved = _maybe_convert_prediction(value, context)
|
|
with get_conn() as conn:
|
|
row = conn.execute(
|
|
"SELECT property_values_json, collection_id FROM collection_pages WHERE id=?",
|
|
(page_id,),
|
|
).fetchone()
|
|
if not row:
|
|
raise ValueError(f"page {page_id} not found")
|
|
try:
|
|
props = json.loads(row["property_values_json"])
|
|
except (TypeError, json.JSONDecodeError):
|
|
props = {}
|
|
props[prop] = resolved
|
|
from app.routers.collections import _validate_page_properties
|
|
_validate_page_properties(conn, row["collection_id"], props, exclude_page_id=page_id)
|
|
conn.execute(
|
|
"UPDATE collection_pages SET property_values_json=?, updated_at=CURRENT_TIMESTAMP WHERE id=?",
|
|
(json.dumps(props), page_id),
|
|
)
|
|
conn.commit()
|
|
return f"set property {prop} = {resolved}"
|
|
|
|
if atype == "create_page":
|
|
coll_id = action.get("collection_id")
|
|
title = action.get("title", "Automation page")
|
|
properties = action.get("properties", {}) or {}
|
|
if not coll_id:
|
|
raise ValueError("create_page requires a collection_id")
|
|
_, resolved_title = _maybe_convert_prediction(title, context)
|
|
resolved_props = {}
|
|
for k, v in properties.items():
|
|
_, pv = _maybe_convert_prediction(v, context)
|
|
resolved_props[k] = pv
|
|
with get_conn() as conn:
|
|
max_pos = conn.execute(
|
|
"SELECT COALESCE(MAX(position), -1) + 1 FROM collection_pages WHERE collection_id=?",
|
|
(coll_id,),
|
|
).fetchone()[0]
|
|
cur = conn.execute(
|
|
"INSERT INTO collection_pages (collection_id, title, position, property_values_json) VALUES (?,?,?,?)",
|
|
(coll_id, resolved_title, max_pos, json.dumps(resolved_props)),
|
|
)
|
|
conn.commit()
|
|
return f"created page {cur.lastrowid} in collection {coll_id}"
|
|
|
|
if atype == "notify":
|
|
message = action.get("message", "Automation fired")
|
|
user_id = action.get("user_id")
|
|
if not user_id:
|
|
user_id = context.get("created_by") or 1
|
|
_, resolved = _maybe_convert_prediction(message, context)
|
|
notifications.create_notification(
|
|
user_id=user_id,
|
|
actor_id=context.get("created_by") or 1,
|
|
ntype="page",
|
|
title=context.get("automation_name", "Automation"),
|
|
message=resolved,
|
|
resource_type="collection_page" if context.get("page_id") else "page",
|
|
resource_id=context.get("page_id") or context.get("collection_id") or 0,
|
|
url=context.get("url", ""),
|
|
)
|
|
return f"notified user {user_id}"
|
|
|
|
if atype == "slack":
|
|
url = _secret_value(action.get("webhook_url") or action.get("url") or "")
|
|
if not url:
|
|
raise ValueError("slack action requires a webhook_url")
|
|
_, text = _maybe_convert_prediction(
|
|
action.get("text") or action.get("message") or "Automation fired", context)
|
|
return await _post_slack(url, text)
|
|
|
|
if atype == "email":
|
|
to = action.get("to", "")
|
|
_, subject = _maybe_convert_prediction(action.get("subject", "FlowDeck automation"), context)
|
|
_, body = _maybe_convert_prediction(action.get("body", action.get("message", "")), context)
|
|
return await _send_email_action(to, subject, body, context)
|
|
|
|
if atype == "forge_issue":
|
|
provider = (action.get("provider") or "gitea").lower()
|
|
owner = action.get("owner", "")
|
|
repo = action.get("repo", "")
|
|
if not owner or not repo:
|
|
raise ValueError("forge_issue requires owner + repo")
|
|
_, title = _maybe_convert_prediction(action.get("title", "Automation issue"), context)
|
|
_, body = _maybe_convert_prediction(action.get("body", ""), context)
|
|
return await _create_forge_issue(
|
|
provider, owner, repo, title, body,
|
|
labels=action.get("labels") or [],
|
|
user_id=context.get("created_by"),
|
|
)
|
|
|
|
if atype == "agent_trigger":
|
|
agent_id = action.get("agent_id")
|
|
if not agent_id:
|
|
raise ValueError("agent_trigger requires an agent_id")
|
|
_, message = _maybe_convert_prediction(action.get("message", ""), context)
|
|
return await _run_linked_agent(
|
|
int(agent_id), context.get("created_by") or 1,
|
|
context.get("workspace_id"), message, context)
|
|
|
|
raise ValueError(f"unknown action type: {atype!r}")
|
|
|
|
|
|
async def run_automation(automation_id: int, trigger_source: str, context: dict) -> dict:
|
|
"""Load, condition-check and execute an automation. Records a run row."""
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT * FROM automations WHERE id=?", (automation_id,)).fetchone()
|
|
if not row:
|
|
return {"status": "skipped", "detail": "automation not found"}
|
|
auto = dict(row)
|
|
|
|
if not auto["enabled"]:
|
|
return {"status": "skipped", "detail": "automation disabled"}
|
|
|
|
# v7.0.0: chained steps take over when present (legacy path otherwise).
|
|
stepped = await _maybe_run_stepped(auto, trigger_source, context)
|
|
if stepped is not None:
|
|
return stepped
|
|
|
|
props = context.get("properties")
|
|
before = context.get("before_properties")
|
|
if not evaluate_conditions(auto["condition_json"], props, before):
|
|
_save_run(automation_id, trigger_source, "skipped", "condition not met",
|
|
context.get("collection_id"), context.get("page_id"))
|
|
return {"status": "skipped", "detail": "condition not met"}
|
|
|
|
try:
|
|
actions = json.loads(auto["actions_json"]) if auto["actions_json"] else []
|
|
except (TypeError, json.JSONDecodeError):
|
|
actions = []
|
|
|
|
ctx = dict(context)
|
|
ctx["automation_name"] = auto["name"]
|
|
ctx["created_by"] = auto["created_by"] or ctx.get("created_by")
|
|
|
|
results = []
|
|
try:
|
|
for action in actions or []:
|
|
results.append(await _run_action(action, ctx, trigger_source))
|
|
detail = "; ".join(results)
|
|
_save_run(automation_id, trigger_source, "fired", detail,
|
|
ctx.get("collection_id"), ctx.get("page_id"))
|
|
# v6.4.0: emit automation.fired (goes through fire_event → outbound
|
|
# webhooks, but NOT back through automations to avoid recursion).
|
|
try:
|
|
from app.services.webhook_outbound import fire_event as _fire_wh
|
|
await _fire_wh("automation.fired", {
|
|
"automation_id": automation_id,
|
|
"name": auto["name"],
|
|
"trigger": trigger_source,
|
|
"collection_id": ctx.get("collection_id"),
|
|
"page_id": ctx.get("page_id"),
|
|
"detail": detail,
|
|
})
|
|
except Exception: # noqa: BLE001
|
|
logger.debug("automation.fired webhook dispatch failed")
|
|
return {"status": "fired", "detail": detail}
|
|
except Exception as exc: # noqa: BLE001 — record every failure in history
|
|
logger.warning("Automation %s failed: %s", automation_id, exc)
|
|
_save_run(automation_id, trigger_source, "error", str(exc),
|
|
ctx.get("collection_id"), ctx.get("page_id"))
|
|
return {"status": "error", "detail": str(exc)}
|
|
|
|
|
|
async def fire_event(event: str, payload: dict):
|
|
"""Dispatch an event to outbound webhooks and matching automations."""
|
|
# v7.3.0: page.updated → in-app notification to followers (throttled).
|
|
if event == "page.updated":
|
|
try:
|
|
from app.services.wiki import notify_followers_of_page_update
|
|
notify_followers_of_page_update(
|
|
payload.get("page_id"), payload.get("actor_id"),
|
|
payload.get("title") or "")
|
|
except Exception: # noqa: BLE001 — notifications are best-effort
|
|
logger.debug("followers notification failed for page.updated")
|
|
# Outbound webhooks (v2.1.0 machinery, previously called nowhere).
|
|
try:
|
|
from app.services.webhook_outbound import fire_event as fire_webhooks
|
|
await fire_webhooks(event, payload)
|
|
except Exception: # noqa: BLE001
|
|
logger.debug("Webhook dispatch failed for %s", event)
|
|
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"""SELECT * FROM automations
|
|
WHERE trigger_type='event' AND event=? AND enabled=1""",
|
|
(event,),
|
|
).fetchall()
|
|
stepped_ids: set[int] = set()
|
|
try:
|
|
with get_conn() as _c:
|
|
stepped_ids = {r[0] for r in _c.execute(
|
|
"SELECT DISTINCT automation_id FROM automation_steps").fetchall()}
|
|
except Exception: # noqa: BLE001 — table missing on very old DBs
|
|
pass
|
|
for row in rows:
|
|
auto = dict(row)
|
|
if auto["id"] in stepped_ids:
|
|
continue # v7.0.0: handled by fire_stepped_event below (no double run)
|
|
if auto["collection_id"] and payload.get("collection_id") != auto["collection_id"]:
|
|
continue
|
|
context = dict(payload)
|
|
context["event"] = event
|
|
await run_automation(auto["id"], "event", context)
|
|
|
|
# v7.0.0: step-based automations (multi-trigger any/all, chains).
|
|
try:
|
|
await fire_stepped_event(event, payload)
|
|
except Exception: # noqa: BLE001
|
|
logger.debug("stepped dispatch failed for %s", event)
|
|
|
|
|
|
# ═══════════ Cron scheduling (trigger_type='cron') ═══════════
|
|
|
|
_SUPPORTED_CRON = {
|
|
"*/1": 1, "*/5": 5, "*/10": 10, "*/15": 15, "*/30": 30,
|
|
"*/2": 2, "*/3": 3, "*/6": 6, "*/12": 12, "*/20": 20, "*/45": 45,
|
|
}
|
|
|
|
|
|
def cron_due(expression: str, last_run_at: str | None, now: datetime | None = None) -> bool:
|
|
"""True when a ``*/N`-style or fixed-minute cron expression is due.
|
|
|
|
Supports ``*/15 * * * *`` (every N minutes) and ``*/N`` alone, plus exact
|
|
``H * * * *`` at minute H of every hour. ``@hourly`` / ``@daily`` also work.
|
|
"""
|
|
expr = (expression or "").strip().lower()
|
|
if not expr:
|
|
return False
|
|
now = now or datetime.utcnow()
|
|
minute = now.minute
|
|
fields = expr.split()
|
|
|
|
if expr in ("@hourly", "hourly"):
|
|
if last_run_at is None:
|
|
return True
|
|
try:
|
|
last = datetime.fromisoformat(str(last_run_at).replace("Z", ""))
|
|
except Exception:
|
|
return True
|
|
return (now - last.replace(tzinfo=None)) >= timedelta(minutes=60)
|
|
|
|
if expr in ("@daily", "daily"):
|
|
if last_run_at is None:
|
|
return True
|
|
try:
|
|
last = datetime.fromisoformat(str(last_run_at).replace("Z", ""))
|
|
except Exception:
|
|
return True
|
|
return (now - last.replace(tzinfo=None)) >= timedelta(hours=24)
|
|
|
|
# "*/N * * * *" → every N minutes
|
|
if fields and fields[0].startswith("*/"):
|
|
val = fields[0][2:]
|
|
if not val.isdigit() or int(val) not in _SUPPORTED_CRON.values():
|
|
return False
|
|
n = int(val)
|
|
if last_run_at is None:
|
|
return True
|
|
try:
|
|
last = datetime.fromisoformat(str(last_run_at).replace("Z", ""))
|
|
except Exception:
|
|
return True
|
|
return (now - last.replace(tzinfo=None)) >= timedelta(minutes=n)
|
|
|
|
# "H * * * *" → at a fixed minute of each hour
|
|
if len(fields) == 5 and fields[0].isdigit():
|
|
return int(fields[0]) == minute
|
|
|
|
return False
|
|
|
|
|
|
async def automation_scheduler():
|
|
"""Background loop: fire due cron automations (checked every 60s)."""
|
|
while True:
|
|
try:
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"SELECT * FROM automations WHERE trigger_type='cron' AND enabled=1"
|
|
).fetchall()
|
|
for row in rows:
|
|
auto = dict(row)
|
|
try:
|
|
if cron_due(auto["cron_expression"], auto["last_run_at"]):
|
|
context = {
|
|
"collection_id": auto["collection_id"] or 0,
|
|
"page_id": None,
|
|
"properties": None,
|
|
}
|
|
await run_automation(auto["id"], "cron", context)
|
|
except Exception: # noqa: BLE001
|
|
logger.warning("Cron automation %s errored", auto["id"])
|
|
# v7.0.0: workers on a cron schedule share the same 60s loop.
|
|
try:
|
|
from app.services.workers import run_due_workers
|
|
await run_due_workers()
|
|
except Exception: # noqa: BLE001
|
|
logger.debug("worker cron iteration failed")
|
|
except Exception: # noqa: BLE001
|
|
logger.warning("automation_scheduler iteration failed")
|
|
await asyncio.sleep(60)
|
|
|
|
|
|
# ═══════════ v7.0.0 — multi-step automations (triggers/conditions/delay) ══
|
|
|
|
STEP_KINDS = ("trigger", "condition", "delay", "action")
|
|
STEP_ACTION_TYPES = ("webhook", "set_property", "create_page", "notify",
|
|
"slack", "email", "forge_issue", "agent_trigger")
|
|
ALL_MODE_WINDOW_S = 300.0
|
|
|
|
# mode=all bookkeeping (single-process): automation_id -> {event: timestamp}.
|
|
_ALL_PENDING: dict[int, dict[str, float]] = {}
|
|
|
|
|
|
def reset_all_pending() -> None:
|
|
"""Test helper: clear the mode=all arrival window."""
|
|
_ALL_PENDING.clear()
|
|
|
|
|
|
def _secret_value(stored: str | None) -> str:
|
|
"""Decrypt a Fernet secret, falling back to raw plaintext (legacy/tests)."""
|
|
if not stored:
|
|
return ""
|
|
try:
|
|
from app.services.sso_provisioning import decrypt_secret
|
|
decrypted = decrypt_secret(stored)
|
|
if decrypted:
|
|
return decrypted
|
|
except Exception: # noqa: BLE001
|
|
pass
|
|
if isinstance(stored, str) and not stored.startswith("gAAAAA"):
|
|
return stored
|
|
return ""
|
|
|
|
|
|
def validate_step(kind: str, config: dict) -> None:
|
|
"""Validate a step payload. Raises ValueError with a human message."""
|
|
from fastapi import HTTPException
|
|
if kind not in STEP_KINDS:
|
|
raise HTTPException(400, f"invalid kind: {kind!r} (want trigger|condition|delay|action)")
|
|
config = config or {}
|
|
if kind == "trigger":
|
|
if not config.get("event"):
|
|
raise HTTPException(400, "trigger step requires an event")
|
|
elif kind == "condition":
|
|
if config.get("op", "eq") not in COND_OPS:
|
|
raise HTTPException(400, f"invalid op: {config.get('op')!r}")
|
|
elif kind == "delay":
|
|
try:
|
|
seconds = int(config.get("seconds", 0))
|
|
except (TypeError, ValueError):
|
|
raise HTTPException(400, "delay step requires integer seconds") from None
|
|
if seconds < 0 or seconds > 86400:
|
|
raise HTTPException(400, "delay seconds must be 0..86400")
|
|
elif kind == "action":
|
|
if config.get("type") not in STEP_ACTION_TYPES:
|
|
raise HTTPException(400, f"invalid action type: {config.get('type')!r}")
|
|
|
|
|
|
def get_steps(automation_id: int) -> list[dict]:
|
|
"""Ordered steps of an automation (empty when legacy single-mode)."""
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"SELECT * FROM automation_steps WHERE automation_id=? ORDER BY position, id",
|
|
(automation_id,),
|
|
).fetchall()
|
|
out = []
|
|
for r in rows:
|
|
d = dict(r)
|
|
try:
|
|
d["config"] = json.loads(d.get("config_json") or "{}")
|
|
except (TypeError, json.JSONDecodeError):
|
|
d["config"] = {}
|
|
out.append(d)
|
|
return out
|
|
|
|
|
|
def _steps_by_kind(steps: list[dict]) -> dict[str, list[dict]]:
|
|
grouped: dict[str, list[dict]] = {"trigger": [], "condition": [],
|
|
"delay": [], "action": []}
|
|
for s in steps:
|
|
if s.get("kind") in grouped:
|
|
grouped[s["kind"]].append(s)
|
|
return grouped
|
|
|
|
|
|
def _step_trigger_matches(step_cfg: dict, event: str, payload: dict,
|
|
automation_collection_id: int | None) -> bool:
|
|
if step_cfg.get("event") != event:
|
|
return False
|
|
want_coll = step_cfg.get("collection_id") or automation_collection_id
|
|
if want_coll and payload.get("collection_id") != want_coll:
|
|
return False
|
|
return True
|
|
|
|
|
|
async def _run_with_steps(auto: dict, steps: list[dict], trigger_source: str,
|
|
context: dict) -> dict:
|
|
"""Execute a chained automation. Records one run row with per-step detail."""
|
|
grouped = _steps_by_kind(steps)
|
|
props = context.get("properties")
|
|
before = context.get("before_properties")
|
|
# Legacy single condition still applies on top of step conditions.
|
|
if not evaluate_conditions(auto.get("condition_json") or "[]", props, before):
|
|
_save_run(auto["id"], trigger_source, "skipped", "condition not met",
|
|
context.get("collection_id"), context.get("page_id"))
|
|
return {"status": "skipped", "detail": "condition not met"}
|
|
for cond in grouped["condition"]:
|
|
cfg = cond.get("config") or {}
|
|
if not match_condition_props(props, before, {
|
|
"property": cfg.get("property"), "op": cfg.get("op", "eq"),
|
|
"value": cfg.get("value")}):
|
|
_save_run(auto["id"], trigger_source, "skipped",
|
|
f"step condition not met: {cfg.get('property')}",
|
|
context.get("collection_id"), context.get("page_id"))
|
|
return {"status": "skipped", "detail": "step condition not met"}
|
|
|
|
ctx = dict(context)
|
|
ctx["automation_name"] = auto["name"]
|
|
ctx["created_by"] = auto["created_by"] or ctx.get("created_by")
|
|
ordered = sorted(steps, key=lambda s: (s.get("position", 0), s.get("id", 0)))
|
|
results = []
|
|
try:
|
|
for step in ordered:
|
|
kind = step.get("kind")
|
|
cfg = step.get("config") or {}
|
|
if kind in ("trigger", "condition"):
|
|
continue
|
|
if kind == "delay":
|
|
seconds = max(0, min(int(cfg.get("seconds", 0)), 300))
|
|
if seconds:
|
|
await asyncio.sleep(seconds)
|
|
results.append(f"delay {cfg.get('seconds', 0)}s")
|
|
elif kind == "action":
|
|
summary = await _run_action({"type": cfg.get("type"), **cfg}, ctx,
|
|
trigger_source)
|
|
results.append(summary)
|
|
detail = "; ".join(results) or "no steps executed"
|
|
_save_run(auto["id"], trigger_source, "fired", detail,
|
|
ctx.get("collection_id"), ctx.get("page_id"))
|
|
try:
|
|
from app.services.webhook_outbound import fire_event as _fire_wh
|
|
await _fire_wh("automation.fired", {
|
|
"automation_id": auto["id"], "name": auto["name"],
|
|
"trigger": trigger_source, "collection_id": ctx.get("collection_id"),
|
|
"page_id": ctx.get("page_id"), "detail": detail})
|
|
except Exception: # noqa: BLE001
|
|
logger.debug("automation.fired webhook dispatch failed")
|
|
return {"status": "fired", "detail": detail}
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("Automation %s (steps) failed: %s", auto["id"], exc)
|
|
_save_run(auto["id"], trigger_source, "error", str(exc),
|
|
ctx.get("collection_id"), ctx.get("page_id"))
|
|
return {"status": "error", "detail": str(exc)}
|
|
|
|
|
|
async def _maybe_run_stepped(auto: dict, trigger_source: str, context: dict) -> dict | None:
|
|
"""Run via steps when the automation has any; None → use legacy path."""
|
|
steps = get_steps(auto["id"])
|
|
if not steps:
|
|
return None
|
|
return await _run_with_steps(auto, steps, trigger_source, context)
|
|
|
|
|
|
def _match_stepped_automations(event: str, payload: dict) -> list[tuple[dict, list[dict]]]:
|
|
"""Automations (enabled) whose trigger steps match ``event`` + collection."""
|
|
with get_conn() as conn:
|
|
rows = conn.execute(
|
|
"""SELECT a.* FROM automations a
|
|
JOIN automation_steps s ON s.automation_id = a.id
|
|
WHERE a.enabled=1 AND s.kind='trigger' GROUP BY a.id"""
|
|
).fetchall()
|
|
matched = []
|
|
for row in rows:
|
|
auto = dict(row)
|
|
steps = get_steps(auto["id"])
|
|
triggers = [s for s in steps if s.get("kind") == "trigger"]
|
|
if any(_step_trigger_matches(t.get("config") or {}, event, payload,
|
|
auto.get("collection_id")) for t in triggers):
|
|
matched.append((auto, triggers))
|
|
return matched
|
|
|
|
|
|
async def fire_stepped_event(event: str, payload: dict) -> None:
|
|
"""Dispatch ``event`` to step-based automations (mode any/all).
|
|
|
|
Called from :func:`fire_event` after the legacy matcher. Unknown events
|
|
(not in the webhook catalogue) still work here — steps are independent
|
|
from outbound webhooks.
|
|
"""
|
|
now = time.time()
|
|
for auto, triggers in _match_stepped_automations(event, payload):
|
|
# Skip automations already handled by the legacy matcher to avoid
|
|
# double runs (legacy = trigger_type event + no steps).
|
|
if not get_steps(auto["id"]):
|
|
continue
|
|
mode = (auto.get("trigger_mode") or "any").lower()
|
|
if mode == "all":
|
|
pending = _ALL_PENDING.setdefault(auto["id"], {})
|
|
pending[event] = now
|
|
# Expire arrivals outside the window.
|
|
for ev in [e for e, ts in pending.items() if now - ts > ALL_MODE_WINDOW_S]:
|
|
del pending[ev]
|
|
wanted = {t.get("config", {}).get("event") for t in triggers}
|
|
if not wanted <= set(pending):
|
|
continue
|
|
_ALL_PENDING.pop(auto["id"], None)
|
|
context = dict(payload)
|
|
context["event"] = event
|
|
await run_automation(auto["id"], "event", context)
|
|
|
|
|
|
# ── v7.0.0 action backends (module-level = monkeypatchable in tests) ───────
|
|
|
|
async def _post_slack(webhook_url: str, text: str) -> str:
|
|
async with httpx.AsyncClient(timeout=10) as client:
|
|
resp = await client.post(webhook_url, json={"text": text})
|
|
if resp.status_code >= 400:
|
|
raise RuntimeError(f"slack webhook returned HTTP {resp.status_code}")
|
|
return f"slack → ({resp.status_code})"
|
|
|
|
|
|
async def _send_email_action(to: str, subject: str, body: str, context: dict) -> str:
|
|
from app.services import mailer
|
|
address = (to or "").strip()
|
|
if address.startswith("user:"):
|
|
try:
|
|
uid = int(address.split(":", 1)[1])
|
|
except ValueError:
|
|
raise ValueError(f"bad email target: {to!r}") from None
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT email FROM users WHERE id=?", (uid,)).fetchone()
|
|
address = (row["email"] if row and row["email"] else "")
|
|
if not address:
|
|
raise ValueError(f"user {uid} has no email")
|
|
if not address:
|
|
address = None
|
|
with get_conn() as conn:
|
|
row = conn.execute("SELECT email FROM users WHERE id=?",
|
|
(context.get("created_by") or 1,)).fetchone()
|
|
if row and row["email"]:
|
|
address = row["email"]
|
|
if not address:
|
|
return "email skipped (no recipient)"
|
|
ok = mailer.send_email(address, subject or "FlowDeck automation", body or "")
|
|
return f"email → {address}" if ok else "email skipped (SMTP not configured)"
|
|
|
|
|
|
async def _create_forge_issue(provider: str, owner: str, repo: str, title: str,
|
|
body: str, labels: list | None = None,
|
|
user_id: int | None = None) -> str:
|
|
token = ""
|
|
if user_id:
|
|
with get_conn() as conn:
|
|
row = conn.execute(
|
|
"SELECT access_token FROM user_oauth_tokens WHERE user_id=? AND provider=?",
|
|
(user_id, provider)).fetchone()
|
|
token = (row["access_token"] if row else "") or ""
|
|
if provider == "github":
|
|
if not token:
|
|
raise ValueError("github action needs a linked GitHub account (token)")
|
|
async with httpx.AsyncClient(timeout=15) as client:
|
|
resp = await client.post(
|
|
f"https://api.github.com/repos/{owner}/{repo}/issues",
|
|
headers={"Authorization": f"Bearer {token}",
|
|
"Accept": "application/vnd.github+json"},
|
|
json={"title": title, "body": body,
|
|
"labels": labels or []} if labels else {"title": title, "body": body},
|
|
)
|
|
if resp.status_code >= 400:
|
|
raise RuntimeError(f"github returned HTTP {resp.status_code}")
|
|
return f"github issue #{resp.json().get('number')} in {owner}/{repo}"
|
|
# gitea (default)
|
|
from app.services.gitea_client import GiteaClient
|
|
gitea = GiteaClient(user_token=token or None)
|
|
issue = await gitea.create_issue(owner, repo, title, body)
|
|
return f"gitea issue #{issue.get('number')} in {owner}/{repo}"
|
|
|
|
|
|
async def _run_linked_agent(agent_id: int, user_id: int, workspace_id: int | None,
|
|
message: str, context: dict) -> str:
|
|
from app.services.agent_engine import AgentEngine
|
|
with get_conn() as conn:
|
|
agent = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone()
|
|
if not agent:
|
|
raise ValueError(f"agent {agent_id} not found")
|
|
cur = conn.execute(
|
|
"""INSERT INTO agent_conversations (agent_id, user_id, title, context_json)
|
|
VALUES (?,?,?,?)""",
|
|
(agent_id, user_id,
|
|
f"Automation: {context.get('automation_name', 'run')}",
|
|
json.dumps({"workspace_id": workspace_id})),
|
|
)
|
|
conv_id = cur.lastrowid
|
|
conn.commit()
|
|
objective = ((agent["system_instructions"] or "").strip()
|
|
or f"Exécute l'agent « {agent['name']} ».")
|
|
if message:
|
|
objective = f"{objective}\n\n{message}"
|
|
engine = AgentEngine(user_id, workspace_id, agent["model"] or None)
|
|
final = ""
|
|
async for _ev in engine.run(conv_id, objective, model=agent["model"]):
|
|
pass
|
|
with get_conn() as conn:
|
|
row = conn.execute(
|
|
"SELECT content FROM agent_messages WHERE conversation_id=? AND role='assistant'"
|
|
" ORDER BY id DESC LIMIT 1", (conv_id,)).fetchone()
|
|
final = (row["content"][:300] if row and row["content"] else "")
|
|
return f"agent « {agent['name']} » ran (conversation {conv_id})" + (f": {final}" if final else "")
|
|
|
|
|
|
# ── v7.0.0 native DB button ────────────────────────────────────────────────
|
|
|
|
async def press_button(collection_id: int, row_id: int, prop_ref: str | int,
|
|
user_id: int) -> dict:
|
|
"""Run the automation linked to a ``button`` property cell."""
|
|
with get_conn() as conn:
|
|
if isinstance(prop_ref, int) or str(prop_ref).isdigit():
|
|
prop = conn.execute(
|
|
"SELECT * FROM collection_properties WHERE id=? AND collection_id=?",
|
|
(int(prop_ref), collection_id)).fetchone()
|
|
else:
|
|
prop = conn.execute(
|
|
"SELECT * FROM collection_properties WHERE collection_id=? AND name=?",
|
|
(collection_id, prop_ref)).fetchone()
|
|
if not prop:
|
|
raise ValueError("button property not found")
|
|
prop = dict(prop)
|
|
if prop.get("prop_type") != "button":
|
|
raise ValueError("property is not a button")
|
|
auto_id = prop.get("button_automation_id")
|
|
if not auto_id:
|
|
raise ValueError("button has no linked automation")
|
|
row = conn.execute(
|
|
"SELECT id FROM collection_pages WHERE id=? AND collection_id=?",
|
|
(row_id, collection_id)).fetchone()
|
|
if not row:
|
|
raise ValueError("row not found")
|
|
context = get_page_context(row_id, collection_id)
|
|
context["created_by"] = user_id
|
|
return await run_automation(auto_id, "button", context)
|
|
|