- 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)
235 lines
9.4 KiB
Python
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}
|