Files
flowdeck/app/services/workers.py
T
bruno 1706ad1ee9
FlowDeck CI / lint (push) Successful in 1m48s
FlowDeck CI / test (push) Failing after 21m19s
FlowDeck CI / docker (push) Skipped
feat: v7.3.0 — cycle v6.8.0→v7.3.0 (Sites, Search, Automations, Calendar, SCIM, Wiki) + audit A9
- v6.8.0 Sites & Forms publics (migrations 24)
- v6.9.0 Recherche sémantique hybride + Ask AI (migration 25)
- v7.0.0 Automations v2 multi-étapes + Workers sandboxés (migration 26)
- v7.1.0 Calendar sync Google/CalDAV + Meeting Notes (migration 27)
- v7.2.0 Enterprise : SCIM 2.0, 2FA TOTP/passkeys, audit UI, agent approvals (migration 28)
- v7.3.0 Wiki/Teamspaces, verified pages, collab polish, charts, unfurl (migration 29)
- docs V68→V73, ROADMAP/CHANGELOG/WORKLOAD à jour, VERSION 7.3.0
- A9 : flowdeck.db, flowdeck_dev.db, test-commit.md, upload_test.txt et e2e/{node_modules,shots,test-results} désindexés + ignorés (.gitignore/.dockerignore)
2026-09-30 20:02:57 -04:00

235 lines
9.4 KiB
Python

"""FlowDeck — Workers lite (v7.0.0).
Custom Python snippets run on FlowDeck infrastructure: manual, on a cron
schedule, or shared across the team (fork). Parité Notion Workers (07/2026),
sans facturation : un budget journalier secondes/workspace fait office de
« credits dashboard ».
Sandbox (documenté, best-effort single-process) :
- AST blacklist : ``import os/sys/subprocess/socket``, ``open()``,
``exec/eval/compile``, attributs dunder.
- Pas de réseau, pas de FS ; builtins restreints (pas de ``__import__``).
- Timeout 30 s (thread + join), budget journalier ``daily_budget_s``.
- Seules API exposées : ``log()``, ``ctx`` (dict), ``result`` (dict out).
Voir ``docs/V70_Automations_Workers.md``.
"""
from __future__ import annotations
import ast
import asyncio
import io
import logging
import re
import time
from contextlib import redirect_stdout
from datetime import UTC, datetime
from app.db import get_conn
logger = logging.getLogger(__name__)
RUN_TIMEOUT_S = 30
MAX_CODE_CHARS = 20_000
MAX_LOG_CHARS = 10_000
_SLUG_RE = re.compile(r"^[a-z0-9-]{3,60}$")
_FORBIDDEN_IMPORTS = {"os", "sys", "subprocess", "socket", "shutil",
"pathlib", "io", "asyncio", "threading", "multiprocessing"}
_FORBIDDEN_CALLS = {"open", "exec", "eval", "compile", "__import__"}
class WorkerRejected(ValueError):
"""Raised when worker code violates the sandbox policy."""
def validate_code(code: str) -> None:
"""AST lint of worker code. Raises WorkerRejected on violation."""
code = code or ""
if len(code) > MAX_CODE_CHARS:
raise WorkerRejected(f"code too long ({len(code)} > {MAX_CODE_CHARS})")
try:
tree = ast.parse(code)
except SyntaxError as exc:
raise WorkerRejected(f"syntax error: {exc}") from None
for node in ast.walk(tree):
if isinstance(node, (ast.Import, ast.ImportFrom)):
names = [a.name.split(".")[0] for a in node.names]
if getattr(node, "module", None):
names.append(str(node.module).split(".")[0])
for name in names:
if name in _FORBIDDEN_IMPORTS:
raise WorkerRejected(f"import forbidden: {name}")
elif isinstance(node, ast.Call):
func = node.func
if isinstance(func, ast.Name) and func.id in _FORBIDDEN_CALLS:
raise WorkerRejected(f"call forbidden: {func.id}()")
elif isinstance(node, ast.Attribute):
if isinstance(node.attr, str) and node.attr.startswith("__"):
raise WorkerRejected(f"dunder access forbidden: {node.attr}")
def _slugify(name: str) -> str:
import unicodedata
slug = unicodedata.normalize("NFKD", name or "").encode("ascii", "ignore").decode("ascii")
slug = re.sub(r"[^\w\s-]", "", slug.lower())
return re.sub(r"[-\s]+", "-", slug).strip("-") or "worker"
def unique_slug(base: str, ignore_id: int | None = None) -> str:
slug, i = _slugify(base)[:60] or "worker", 1
with get_conn() as conn:
while conn.execute(
"SELECT id FROM workers WHERE slug=? AND id != COALESCE(?, -1)",
(slug, ignore_id)).fetchone():
i += 1
slug = f"{_slugify(base)[:55]}-{i}"
return slug
_SAFE_BUILTINS = {
"abs": abs, "all": all, "any": any, "bool": bool, "dict": dict,
"enumerate": enumerate, "filter": filter, "float": float, "format": format,
"frozenset": frozenset, "int": int, "len": len, "list": list, "map": map,
"max": max, "min": min, "range": range, "reversed": reversed, "round": round,
"set": set, "sorted": sorted, "str": str, "sum": sum, "tuple": tuple,
"zip": zip, "print": print, "isinstance": isinstance, "type": type,
}
def _exec_code(code: str, ctx: dict) -> tuple[dict, str]:
"""Run validated code in a thread. Returns (result_dict, logs)."""
logs: list[str] = []
def _log(*args) -> None:
logs.append(" ".join(str(a) for a in args))
namespace = {"__builtins__": dict(_SAFE_BUILTINS),
"log": _log, "ctx": dict(ctx or {}), "result": {}}
buf = io.StringIO()
with redirect_stdout(buf):
exec(compile(code, "<worker>", "exec"), namespace) # noqa: S102 — sandboxed
printed = buf.getvalue()
if printed:
logs.append(printed)
result = namespace.get("result")
return result if isinstance(result, dict) else {}, "\n".join(logs)[:MAX_LOG_CHARS]
def daily_usage_s(workspace_id: int | None) -> float:
"""CPU seconds consumed today (UTC) by a workspace's workers."""
day = datetime.now(UTC).strftime("%Y-%m-%d")
with get_conn() as conn:
row = conn.execute(
"""SELECT COALESCE(SUM(wr.duration_ms), 0) FROM worker_runs wr
JOIN workers w ON w.id = wr.worker_id
WHERE date(wr.created_at) = date(?)
AND COALESCE(w.workspace_id, -1) = COALESCE(?, -1)""",
(day, workspace_id)).fetchone()
return (row[0] or 0) / 1000.0
def _save_run(worker_id: int, status: str, logs: str, duration_ms: int) -> int:
with get_conn() as conn:
cur = conn.execute(
"INSERT INTO worker_runs (worker_id, status, logs, duration_ms)"
" VALUES (?,?,?,?)", (worker_id, status, logs[:MAX_LOG_CHARS], duration_ms))
conn.commit()
return cur.lastrowid
def run_worker(worker_id: int, ctx: dict | None = None) -> dict:
"""Execute a worker synchronously (used by the router + cron loop).
Returns {status, run_id, duration_ms}. Never raises for user-code errors
(they become ``error`` runs); raises only when the worker is missing or
over budget (caller maps to 404/429).
"""
from fastapi import HTTPException
with get_conn() as conn:
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
if not row:
raise HTTPException(404, "Worker not found")
worker = dict(row)
validate_code(worker.get("code_py") or "")
used = daily_usage_s(worker.get("workspace_id"))
if used >= (worker.get("daily_budget_s") or 60):
_save_run(worker_id, "over_budget", f"daily budget exceeded ({used:.1f}s used)", 0)
raise HTTPException(429, "Worker daily budget exceeded")
outcome: dict = {}
def _target() -> None:
try:
result, logs = _exec_code(worker.get("code_py") or "", ctx or {})
outcome["result"] = result
outcome["logs"] = logs
except Exception as exc: # noqa: BLE001 — user code, recorded
outcome["error"] = f"{type(exc).__name__}: {exc}"
import threading
started = time.time()
thread = threading.Thread(target=_target, daemon=True)
thread.start()
thread.join(timeout=RUN_TIMEOUT_S)
duration_ms = int((time.time() - started) * 1000)
if thread.is_alive():
run_id = _save_run(worker_id, "timeout",
f"exceeded {RUN_TIMEOUT_S}s timeout", duration_ms)
return {"status": "timeout", "run_id": run_id, "duration_ms": duration_ms}
if "error" in outcome:
run_id = _save_run(worker_id, "error", outcome["error"], duration_ms)
return {"status": "error", "run_id": run_id,
"duration_ms": duration_ms, "error": outcome["error"]}
run_id = _save_run(worker_id, "ok", outcome.get("logs", ""), duration_ms)
return {"status": "ok", "run_id": run_id, "duration_ms": duration_ms,
"result": outcome.get("result", {})}
async def run_due_workers() -> int:
"""Fire workers whose ``schedule_cron`` is due (called from the 60s loop)."""
from app.services.automations import cron_due
with get_conn() as conn:
rows = conn.execute(
"SELECT * FROM workers WHERE schedule_cron IS NOT NULL AND schedule_cron != ''"
).fetchall()
fired = 0
for row in rows:
worker = dict(row)
with get_conn() as conn:
last = conn.execute(
"SELECT MAX(created_at) FROM worker_runs WHERE worker_id=?",
(worker["id"],)).fetchone()[0]
try:
if cron_due(worker["schedule_cron"] or "", last):
loop = asyncio.get_running_loop()
await loop.run_in_executor(None, run_worker, worker["id"], {})
fired += 1
except Exception as exc: # noqa: BLE001 — one worker must not kill the loop
logger.debug("worker %s cron failed: %s", worker["id"], exc)
return fired
def fork_worker(worker_id: int, user_id: int) -> dict:
"""Duplicate a shared worker for another user (Notion-style sharing)."""
with get_conn() as conn:
row = conn.execute("SELECT * FROM workers WHERE id=?", (worker_id,)).fetchone()
if not row:
from fastapi import HTTPException
raise HTTPException(404, "Worker not found")
src = dict(row)
if not src.get("shared") and src.get("created_by") != user_id:
from fastapi import HTTPException
raise HTTPException(403, "Worker is not shared")
slug = unique_slug(f"{src['slug']}-fork")
with get_conn() as conn:
cur = conn.execute(
"""INSERT INTO workers (slug, workspace_id, name, code_py, schedule_cron,
shared, daily_budget_s, created_by)
VALUES (?,?,?,?,?,?,?,?)""",
(slug, src["workspace_id"], f"{src['name']} (fork)", src["code_py"],
"", 0, src["daily_budget_s"], user_id))
conn.commit()
new_id = cur.lastrowid
return {"id": new_id, "slug": slug, "status": "forked", "from": worker_id}