"""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, "", "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}