Files
flowdeck/app/services/automations.py
T
bruno 16f73fe39e
FlowDeck CI / lint (push) Successful in 2m14s
FlowDeck CI / docker (push) Canceled after 0s
FlowDeck CI / test (push) Canceled after 2h7m52s
feat: Add plugins — catalogue on/off à effet réel (v7.58.0, phase 8/8)
- app/services/plugins.py + migration 35 : table plugins (slug, name,
  description, enabled) pré-remplie avec 3 modules câblés — web-tools,
  web-clipper, automations ; ligne absente = activé (défaut sûr)
- automations OFF → dépendance FastAPI posée à l'include_router dans main.py
  (aucun router touché) → toutes les routes /workspace/automations* refusées +
  garde de tick du scheduler en arrière-plan
- web-clipper OFF → GET /extensions et tout /api/v2/web-clipper/* refusés
- web-tools OFF → web_search et fetch_url retirés du schéma ET de execute()
  via ToolRegistry._all() : le LLM ne les voit plus
- UI rendue côté serveur : global Jinja plugin_enabled(slug) — nav
  « Extensions » / « Automations » en {% if %} (absentes du DOM), sections
  conditionnées en x-show dans settings.html
- menu + : l'entrée « Add plugins » devient vivante (fini disabled:true) —
  liste des 3 plugins avec bascule, GET/PATCH /api/agent/plugins[/slug]
  (slug inconnu → 404, 401 sans session)
- tests : tests/test_v758_plugins.py (10 tests) — routes refusées (302 hors
  /api, 404 JSON pour /api*), outils retirés, nav disparue, persistance,
  câblage ; assertions disabled:true == 0 dans les tests des phases 1/3/4/5/7
- livraison : VERSION + app/main = 7.58.0, OpenAPI 525 chemins, CHANGELOG,
  ROADMAP phase 8 cochée (menu + complet), avenant phase 8 (docs)
2026-10-07 10:19:25 -04:00

851 lines
36 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 UTC, datetime, timedelta
from app.db import get_conn
from app.services import notifications, plugins
from app.services.http_client import shared_client
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 shared_client(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)}
def run_event_sync(coro, timeout: float = 60.0):
"""A21 phase 2b : exécute une coroutine d'événement depuis un handler synchrone.
Bloque le worker threadpool (jamais la boucle d'event) et attend la fin —
déterministe, exactement ce que faisait l'await avant la conversion des
routes en `def`.
ponytail: les clients httpx des services sont créés à chaque appel (aucun
lien de boucle) ; si un jour un client/queue est lié à la boucle de l'app,
passer à `asyncio.run_coroutine_threadsafe` + boucle capturée au lifespan.
"""
return asyncio.run(asyncio.wait_for(coro, timeout))
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.now(UTC).replace(tzinfo=None)
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:
if not plugins.is_enabled("automations"): # plugin OFF = rien de planifié
await asyncio.sleep(60)
continue
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.warning("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 shared_client(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 shared_client(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)