venv script improv
All checks were successful
Build and Deploy MITM Webserver / build (push) Successful in 8s

This commit is contained in:
2026-01-28 18:43:24 +01:00
parent 90087039f7
commit eded798271
2 changed files with 228 additions and 160 deletions

View File

@@ -1,39 +1,40 @@
# script_router_with_venv.py
# script_router_named.py
"""
APIRouter for uploading scripts (with optional requirements.txt),
creating per-script venvs (when requirements provided), and managing
systemd services that run scripts with a queue number argument.
APIRouter: upload scripts with a supplied name, optional requirements -> create per-script venv.
If venv install fails, response includes pip output and the router deletes the uploaded files and venv.
Usage:
from fastapi import FastAPI
from script_router_with_venv import router, register_lifecycle
app = FastAPI()
app.include_router(router)
register_lifecycle(app) # optional: stops/removes manager-created units on shutdown
Endpoints:
- POST /scripts -> upload script (multipart): script file, optional requirements file, required 'name' form field
- GET /scripts -> list scripts
- GET /scripts/{name} -> download script
- POST /scripts/{name}/enable -> enable systemd service for script on given qnum
- POST /scripts/{name}/disable -> disable service for script on qnum
- GET /scripts/status -> status of all fw-script units
- GET /scripts/{name}/status -> status of units for that script
Notes:
- The router does NOT interact with nft. You should create/delete nft queue rules with your separate API.
- Script files are stored under SCRIPT_DIR.
- Virtualenvs (if created) are stored under VENV_BASE/<sid>.
- Systemd units are created under /etc/systemd/system with names: fw-script-<sid>-q<qnum>.service
- This code must run with permissions to create venvs, write unit files and call systemctl (typically root).
- This relies on systemd and writes units to /etc/systemd/system
- Script files are stored at SCRIPT_DIR/<name>.py
- Venv stored at VENV_BASE/<name> (if requirements provided)
- 'name' must match regex [A-Za-z0-9_.-]+ (no path separators)
"""
import os
import sys
import re
import uuid
import json
import shutil
import subprocess
import time
import re
import logging
from typing import Optional, List, Dict
from fastapi import APIRouter, UploadFile, File, HTTPException
from fastapi import APIRouter, UploadFile, File, Form, HTTPException
from fastapi.responses import FileResponse
from pydantic import BaseModel
# ---------- Config ----------
# ---------- Configuration ----------
SCRIPT_DIR = "/srv/fw-scripts"
VENV_BASE = "/srv/fw-scripts/venvs"
UNIT_DIR = "/etc/systemd/system"
@@ -45,87 +46,91 @@ os.makedirs(VENV_BASE, exist_ok=True)
# ---------- Logging ----------
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s [%(name)s] %(message)s")
logger = logging.getLogger("script-router-venv")
logger = logging.getLogger("script-router-named")
# ---------- Router ----------
router = APIRouter(prefix="/scripts", tags=["scripts"])
# ---------- Models ----------
class ScriptInfo(BaseModel):
id: str
name: str
path: str
# ---------- Name validation ----------
# Accept only safe file-name characters to avoid path traversal: letters, digits, dot, underscore, hyphen
_NAME_RE = re.compile(r'^[A-Za-z0-9_.-]+$')
class EnableRequest(BaseModel):
qnum: int
service_name: Optional[str] = None
extra_args: Optional[str] = None
enable_at_boot: Optional[bool] = False
def validate_name(name: str) -> None:
if not name:
raise ValueError("name must be provided")
if not _NAME_RE.match(name):
raise ValueError("invalid name; allowed characters: letters, digits, dot, underscore, hyphen")
# prevent reserved names or dots-only
if name in (".", ".."):
raise ValueError("invalid name")
# ---------- Utilities: service names / unit paths ----------
def make_service_name(sid: str, qnum: int) -> str:
return f"{UNIT_PREFIX}-{sid}-q{qnum}"
# ---------- Utility paths ----------
def script_path_for(name: str) -> str:
return os.path.join(SCRIPT_DIR, f"{name}.py")
def requirements_path_for(name: str) -> str:
return os.path.join(SCRIPT_DIR, f"{name}-requirements.txt")
def venv_path_for(name: str) -> str:
return os.path.join(VENV_BASE, name)
def venv_python_for(name: str) -> str:
vpy = os.path.join(venv_path_for(name), "bin", "python")
if os.path.exists(vpy):
return vpy
return "/usr/bin/python3"
def make_service_name(name: str, qnum: int) -> str:
return f"{UNIT_PREFIX}-{name}-q{qnum}"
def unit_path_for(service_name: str) -> str:
return os.path.join(UNIT_DIR, service_name + ".service")
# ---------- Venv helpers ----------
def venv_path_for(sid: str) -> str:
return os.path.join(VENV_BASE, sid)
def venv_python_for(sid: str) -> str:
vpy = os.path.join(venv_path_for(sid), "bin", "python")
if os.path.exists(vpy):
return vpy
# fallback to system python
return "/usr/bin/python3"
def create_venv(sid: str, timeout: int = 60):
venv_dir = venv_path_for(sid)
def create_venv(name: str, timeout: int = 60) -> str:
venv_dir = venv_path_for(name)
if os.path.exists(venv_dir):
logger.debug("Venv already exists for sid=%s: %s", sid, venv_dir)
logger.debug("Venv already exists for %s", name)
return venv_dir
os.makedirs(os.path.dirname(venv_dir), exist_ok=True)
logger.info("Creating venv for sid=%s at %s", sid, venv_dir)
logger.info("Creating venv for %s at %s", name, venv_dir)
try:
subprocess.run([sys.executable, "-m", "venv", venv_dir], check=True, timeout=timeout)
subprocess.run([sys.executable, "-m", "venv", venv_dir], check=True, timeout=timeout, capture_output=True, text=True)
except subprocess.CalledProcessError as e:
raise RuntimeError(f"venv creation failed: {e}")
logger.exception("venv creation failed for %s: %s", name, e.stderr if hasattr(e, "stderr") else str(e))
raise RuntimeError("venv creation failed: " + (e.stderr or str(e)))
except subprocess.TimeoutExpired:
logger.exception("venv creation timed out for %s", name)
raise RuntimeError("venv creation timed out")
return venv_dir
def pip_install_requirements(sid: str, requirements_path: str, timeout: int = 600) -> Dict[str, str]:
def pip_install_requirements(name: str, requirements_path: str, timeout: int = 600) -> Dict[str, str]:
"""
Install requirements into the venv for sid from requirements_path.
Returns dict: { "stdout": "...", "stderr": "..." }
Raises on failure with details in exception message.
Install requirements into the venv for name from requirements_path.
Returns dict with stdout/stderr. Raises RuntimeError on failure including outputs.
"""
venv_dir = create_venv(sid)
venv_dir = create_venv(name)
pip_path = os.path.join(venv_dir, "bin", "pip")
# ensure pip exists and upgrade
# ensure pip exists and attempt to upgrade
try:
subprocess.run([pip_path, "install", "--upgrade", "pip"], check=True, capture_output=True, text=True, timeout=300)
except subprocess.CalledProcessError as e:
# continue but warn
logger.warning("pip upgrade failed for sid=%s: %s", sid, e.stderr if hasattr(e, "stderr") else str(e))
# install requirements
logger.warning("pip upgrade warning for %s: %s", name, getattr(e, "stderr", str(e)))
# run install
try:
p = subprocess.run([pip_path, "install", "-r", requirements_path, "--no-cache-dir"],
check=True, capture_output=True, text=True, timeout=timeout)
logger.info("pip install success for sid=%s", sid)
return {"stdout": p.stdout, "stderr": p.stderr}
logger.info("pip install succeeded for %s", name)
return {"stdout": p.stdout or "", "stderr": p.stderr or ""}
except subprocess.CalledProcessError as e:
logger.error("pip install failed for sid=%s: %s", sid, e.stderr if hasattr(e, "stderr") else str(e))
# return output for debugging
out = {"stdout": getattr(e, "stdout", "") or "", "stderr": getattr(e, "stderr", "") or str(e)}
logger.error("pip install failed for %s: %s", name, out["stderr"][:4000])
raise RuntimeError(jsonify_cmd_output(out))
except subprocess.TimeoutExpired:
logger.error("pip install timeout for sid=%s", sid)
logger.error("pip install timed out for %s", name)
raise RuntimeError("pip install timed out")
def jsonify_cmd_output(out: Dict[str, str]) -> str:
# helper to pack stdout/stderr into a single string message
s = ""
if out.get("stdout"):
s += "STDOUT:\n" + out["stdout"] + "\n"
@@ -153,7 +158,6 @@ WantedBy=multi-user.target
"""
with open(unit_path, "w") as fh:
fh.write(unit_text)
# reload systemd
subprocess.run(["systemctl", "daemon-reload"], check=True)
logger.info("Wrote unit %s", unit_path)
if enable_at_boot:
@@ -192,18 +196,18 @@ def is_unit_active(service_name: str) -> bool:
return p.returncode == 0
def list_fw_units() -> List[str]:
"""Return list of fw unit names without .service suffix."""
"""List our fw-script units (without .service suffix)."""
units = []
try:
for fn in os.listdir(UNIT_DIR):
if fn.startswith(UNIT_PREFIX + "-") and fn.endswith(".service"):
units.append(fn[:-8]) # remove .service
units.append(fn[:-8])
except FileNotFoundError:
logger.warning("Unit dir %s not found", UNIT_DIR)
return units
# ExecStart parse pattern: python path + script path + qnum + optional extra
_RE_EXECSTART = re.compile(r'(?P<py>/\S*python\S*)\s+(?P<script>/\S*?/srv/fw-scripts/(?P<sid>[0-9a-fA-F]+)\.py)\s+(?P<qnum>\d+)(?:\s+(?P<extra>.*))?')
# ExecStart parse pattern: python + /srv/fw-scripts/<name>.py + qnum + optional extra
_RE_EXECSTART = re.compile(r'(?P<py>/\S*python\S*)\s+(?P<script>/\S*?/srv/fw-scripts/(?P<name>[A-Za-z0-9_.-]+)\.py)\s+(?P<qnum>\d+)(?:\s+(?P<extra>.*))?')
def parse_unit_execstart(service_name: str) -> Optional[Dict]:
unit_path = unit_path_for(service_name)
@@ -218,59 +222,122 @@ def parse_unit_execstart(service_name: str) -> Optional[Dict]:
match = _RE_EXECSTART.search(exec_start)
if match:
sd = match.groupdict()
return {"service": service_name, "exec_start": exec_start, "sid": sd["sid"], "script_path": sd["script"], "qnum": int(sd["qnum"]), "extra": sd.get("extra") or ""}
return {"service": service_name, "exec_start": exec_start, "sid": None, "script_path": None, "qnum": None, "extra": None}
return {"service": service_name, "exec_start": exec_start, "name": sd["name"], "script_path": sd["script"], "qnum": int(sd["qnum"]), "extra": sd.get("extra") or ""}
return {"service": service_name, "exec_start": exec_start, "name": None, "script_path": None, "qnum": None, "extra": None}
# ---------- End utilities ----------
# ---------- Models ----------
class ScriptInfo(BaseModel):
name: str
path: str
class EnableRequest(BaseModel):
qnum: int
service_name: Optional[str] = None
extra_args: Optional[str] = None
enable_at_boot: Optional[bool] = False
# ---------- Endpoints ----------
@router.post("", response_model=ScriptInfo)
async def upload_script(script: UploadFile = File(...), requirements: Optional[UploadFile] = File(None)):
async def upload_script(
script: UploadFile = File(...),
name: str = Form(...),
requirements: Optional[UploadFile] = File(None),
):
"""
Upload a script and optional requirements.txt.
If requirements is supplied, a venv for the script will be created and pip will install the requirements.
Returns sid, path, and pip output if any.
Upload a script with a supplied name (Form field 'name'), optional requirements file.
If requirements provided, create venv and install; on failure delete files and venv and return error details.
"""
# validate name
try:
validate_name(name)
except ValueError as e:
logger.warning("Invalid name provided: %s", name)
raise HTTPException(status_code=400, detail=str(e))
# ensure script file extension is .py
if not script.filename.endswith(".py"):
logger.warning("Rejected upload with invalid extension: %s", script.filename)
logger.warning("Upload rejected: script not .py (name=%s original=%s)", name, script.filename)
raise HTTPException(status_code=400, detail="only .py scripts allowed")
# ensure uniqueness
spath = script_path_for(name)
if os.path.exists(spath):
logger.warning("Upload rejected: script with name already exists: %s", name)
raise HTTPException(status_code=409, detail="script with that name already exists")
# write script
data = await script.read()
if len(data) > 2_000_000:
logger.warning("Rejected upload too large: %s size=%d", script.filename, len(data))
raise HTTPException(status_code=400, detail="script too large")
sid = uuid.uuid4().hex
fname = f"{sid}.py"
path = os.path.join(SCRIPT_DIR, fname)
with open(path, "wb") as fh:
fh.write(data)
os.chmod(path, 0o700)
logger.info("Uploaded script %s as sid=%s path=%s", script.filename, sid, path)
try:
with open(spath, "wb") as fh:
fh.write(data)
os.chmod(spath, 0o700)
logger.info("Saved script for name=%s at %s", name, spath)
except Exception as e:
logger.exception("Failed to write script file for %s: %s", name, e)
raise HTTPException(status_code=500, detail="failed to save script")
req_path = None
pip_output = None
if requirements is not None:
# save requirements to a temp path
req_data = await requirements.read()
req_path = os.path.join(SCRIPT_DIR, f"{sid}-requirements.txt")
with open(req_path, "wb") as fh:
fh.write(req_data)
logger.info("Saved requirements for sid=%s at %s (size=%d)", sid, req_path, len(req_data))
# create venv and install
try:
create_venv(sid)
res = pip_install_requirements(sid, req_path)
pip_output = {"stdout": res.get("stdout", ""), "stderr": res.get("stderr", "")}
logger.info("Installed requirements for sid=%s", sid)
except Exception as e:
# cleanup venv on failure (optional)
logger.exception("pip install failed for sid=%s: %s", sid, e)
# Provide useful error to client but keep script stored for inspection
raise HTTPException(status_code=500, detail=f"pip install failed: {e}")
venv_created = False
# return pip_output in headers/body? We'll include in body (if present).
resp = {"id": sid, "name": script.filename, "path": path}
try:
if requirements is not None:
# save requirements file
req_data = await requirements.read()
req_path = requirements_path_for(name)
with open(req_path, "wb") as fh:
fh.write(req_data)
logger.info("Saved requirements for %s at %s", name, req_path)
# create venv and install
try:
# create venv and pip install
create_venv(name)
venv_created = True
res = pip_install_requirements(name, req_path)
pip_output = {"stdout": res.get("stdout", ""), "stderr": res.get("stderr", "")}
logger.info("pip install completed for %s", name)
except Exception as pip_exc:
# pip install failed: cleanup and report
err_msg = str(pip_exc)
logger.error("pip install error for %s: %s", name, err_msg[:2000])
# cleanup: remove script file, req file, venv dir if created
try:
if os.path.exists(spath):
os.remove(spath)
if req_path and os.path.exists(req_path):
os.remove(req_path)
if venv_created and os.path.isdir(venv_path_for(name)):
shutil.rmtree(venv_path_for(name), ignore_errors=True)
except Exception:
logger.exception("Cleanup after pip failure partially failed for %s", name)
# respond with failure detail (include pip output if available)
raise HTTPException(status_code=500, detail=f"pip install failed: {err_msg}")
except HTTPException:
# pass-through to ensure upstream returns after cleanup
raise
except Exception as e:
# unexpected failure: attempt cleanup of script + req + venv
logger.exception("Unexpected error during upload for %s: %s", name, e)
try:
if os.path.exists(spath):
os.remove(spath)
except Exception:
pass
try:
if req_path and os.path.exists(req_path):
os.remove(req_path)
except Exception:
pass
try:
if os.path.isdir(venv_path_for(name)):
shutil.rmtree(venv_path_for(name), ignore_errors=True)
except Exception:
pass
raise HTTPException(status_code=500, detail="internal error during upload")
# success
resp = {"name": name, "path": spath}
if pip_output is not None:
resp["pip"] = pip_output
return resp
@@ -281,90 +348,94 @@ def list_scripts():
for fn in os.listdir(SCRIPT_DIR):
if not fn.endswith(".py"):
continue
sid = fn.rsplit(".", 1)[0]
p = os.path.join(SCRIPT_DIR, fn)
out.append({"id": sid, "name": fn, "path": p})
logger.debug("Listed scripts: %d", len(out))
name = fn.rsplit(".", 1)[0]
out.append({"name": name, "path": os.path.join(SCRIPT_DIR, fn)})
logger.debug("Listed %d scripts", len(out))
return out
@router.get("/{sid}")
def download_script(sid: str):
path = os.path.join(SCRIPT_DIR, f"{sid}.py")
@router.get("/{name}")
def download_script(name: str):
try:
validate_name(name)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
path = script_path_for(name)
if not os.path.exists(path):
logger.warning("Download requested for missing sid=%s", sid)
logger.warning("Download requested for missing script %s", name)
raise HTTPException(status_code=404, detail="not found")
logger.info("Download script sid=%s path=%s", sid, path)
return FileResponse(path, media_type="text/x-python", filename=f"{sid}.py")
logger.info("Download script %s", name)
return FileResponse(path, media_type="text/x-python", filename=f"{name}.py")
@router.post("/{sid}/enable")
def enable_script(sid: str, req: EnableRequest):
"""
Create and start systemd service that runs the script with the queue number.
If a venv was created at upload time, ExecStart will use that venv's python.
"""
script_path = os.path.join(SCRIPT_DIR, f"{sid}.py")
@router.post("/{name}/enable")
def enable_script(name: str, req: EnableRequest):
try:
validate_name(name)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
script_path = script_path_for(name)
if not os.path.exists(script_path):
logger.warning("Enable requested for missing sid=%s", sid)
raise HTTPException(status_code=404, detail="script not found")
qnum = req.qnum
service_name = req.service_name or make_service_name(sid, qnum)
# prefer venv python if exists
python_path = venv_python_for(sid)
service_name = req.service_name or make_service_name(name, qnum)
python_path = venv_python_for(name)
exec_start = f"{python_path} {script_path} {qnum}"
if req.extra_args:
exec_start += " " + req.extra_args
try:
write_unit(service_name, exec_start, description=f"FW script {sid} queue {qnum}", enable_at_boot=req.enable_at_boot)
# slight delay then start
write_unit(service_name, exec_start, description=f"FW script {name} queue {qnum}", enable_at_boot=req.enable_at_boot)
time.sleep(0.05)
start_unit(service_name)
except subprocess.CalledProcessError as e:
logger.exception("Failed to start service %s for sid=%s: %s", service_name, sid, e)
logger.exception("Failed to start service %s for %s: %s", service_name, name, e)
try:
remove_unit(service_name)
except Exception:
pass
raise HTTPException(status_code=500, detail=f"systemd start failed: {e}")
except Exception as e:
logger.exception("Unknown error starting service %s for sid=%s: %s", service_name, sid, e)
logger.exception("Unknown error starting service %s for %s: %s", service_name, name, e)
try:
remove_unit(service_name)
except Exception:
pass
raise HTTPException(status_code=500, detail=str(e))
logger.info("Enabled script sid=%s on qnum=%d as service=%s (python=%s)", sid, qnum, service_name, python_path)
return {"status": "ok", "sid": sid, "qnum": qnum, "service": service_name, "python": python_path}
logger.info("Enabled script %s on qnum=%d as service=%s (python=%s)", name, qnum, service_name, python_path)
return {"status": "ok", "name": name, "qnum": qnum, "service": service_name, "python": python_path}
@router.post("/{sid}/disable")
def disable_script(sid: str, qnum: int):
service_name = make_service_name(sid, qnum)
@router.post("/{name}/disable")
def disable_script(name: str, qnum: int):
try:
validate_name(name)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
service_name = make_service_name(name, qnum)
unit_p = unit_path_for(service_name)
if not os.path.exists(unit_p):
# try stopping anyway in case only registered with systemd without unit file
# attempt to stop anyway
try:
subprocess.run(["systemctl", "stop", service_name], check=False)
except Exception:
pass
logger.warning("Disable requested but unit missing for sid=%s qnum=%s", sid, qnum)
raise HTTPException(status_code=404, detail="service/unit not found")
try:
stop_unit(service_name)
except subprocess.CalledProcessError:
except Exception:
logger.exception("Failed stopping service %s", service_name)
try:
remove_unit(service_name)
except Exception:
logger.exception("Failed removing unit %s", service_name)
logger.info("Disabled script sid=%s on qnum=%s (removed service %s)", sid, qnum, service_name)
return {"status": "ok", "sid": sid, "qnum": qnum}
logger.info("Disabled script %s on qnum=%s (removed service %s)", name, qnum, service_name)
return {"status": "ok", "name": name, "qnum": qnum}
@router.get("/status")
def status_all():
units = list_fw_units()
results: Dict[str, Dict] = {}
results = {}
for svc in units:
parsed = parse_unit_execstart(svc)
try:
@@ -372,37 +443,36 @@ def status_all():
except Exception:
active = False
results[svc] = {"parsed": parsed, "active": active}
logger.debug("Status queried: found %d fw units", len(results))
logger.debug("Status queried: found %d units", len(results))
return results
@router.get("/{sid}/status")
def status_for_sid(sid: str):
@router.get("/{name}/status")
def status_for_name(name: str):
try:
validate_name(name)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
units = list_fw_units()
matches = []
prefix = f"{UNIT_PREFIX}-{name}-q"
for svc in units:
if svc.startswith(f"{UNIT_PREFIX}-{sid}-q"):
if svc.startswith(prefix):
parsed = parse_unit_execstart(svc)
try:
active = is_unit_active(svc)
except Exception:
active = False
matches.append({"service": svc, "parsed": parsed, "active": active})
logger.debug("Status for sid=%s -> %d matches", sid, len(matches))
return {"sid": sid, "mappings": matches}
logger.debug("Status for %s -> %d matches", name, len(matches))
return {"name": name, "mappings": matches}
# ---------- Lifecycle helper (optional) ----------
# ---------- Lifecycle helper ----------
def register_lifecycle(app):
"""
Optionally call this in your main FastAPI app to stop & remove any manager-created units at shutdown.
This will scan for units with our prefix and remove them -- WARNING: if you want systemd to keep services after manager stops,
do NOT register this lifecycle handler.
"""
@app.on_event("shutdown")
def _shutdown_event():
logger.info("Shutdown: removing manager-created systemd units with prefix %s", UNIT_PREFIX)
logger.info("Shutdown: stopping/removing manager-created units with prefix %s", UNIT_PREFIX)
units = list_fw_units()
for svc in units:
# only remove units that match our naming convention (fw-script-<sid>-q<qnum>)
if svc.startswith(UNIT_PREFIX + "-"):
try:
subprocess.run(["systemctl", "stop", svc], check=False)
@@ -413,5 +483,3 @@ def register_lifecycle(app):
except Exception:
logger.exception("Failed to remove unit %s on shutdown", svc)
logger.info("Shutdown cleanup complete")
# End of router

View File

@@ -24,7 +24,7 @@ PYTHON_VERSION="3" # aktuelle Python 3 Version
# -----------------------------
echo "==> Update & Upgrade"
apt update && apt upgrade -y
apt install -y git curl build-essential nginx python3 python3-pip python3-venv unzip wget
apt install -y git curl build-essential nginx python3 python3-pip python3-venv unzip wget python3-dev libnetfilter-queue-dev libnfnetlink-dev libpcap-dev
# -----------------------------
# Node.js installieren (LTS)