Files
flowdeck/app/services/automations.py
T
bruno e0237e576f
FlowDeck CI / lint (push) Failing after 1m12s
FlowDeck CI / test (push) Failing after 3h3m3s
FlowDeck CI / docker (push) Skipped
Fire automation events across API routers
2026-09-21 20:30:05 -04:00

415 lines
16 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
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")
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}"
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"}
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."""
# 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()
for row in rows:
auto = dict(row)
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)
# ═══════════ 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"])
except Exception: # noqa: BLE001
logger.warning("automation_scheduler iteration failed")
await asyncio.sleep(60)