Files
livef12rocks/app/pipeline.py
T
Frank Schwenk 90192cd284 feat: initial live.f12.rocks SFTP → gmic/rembg → web pipeline
Event pep stack with SFTPGo inbox, sequential worker, FastAPI gallery/remix, and Traefik-ready compose.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-16 21:32:10 +02:00

419 lines
15 KiB
Python

"""Image compose pipeline — ported from `make_random.py`.
Core building blocks (gmic filter application, blending, filter-retry with
timeout) are kept as close to the reference script as possible. On top of
that this module adds job-directory bookkeeping, kept intermediates, and a
manifest format the web app can read/append to (for remix).
"""
from __future__ import annotations
import json
import logging
import random
import re
import subprocess
import sys
import time
import uuid
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from . import config
logger = logging.getLogger("livef12.pipeline")
_ANSI_RE = re.compile(r"\x1b\[[0-9;]*m")
class PipelineError(RuntimeError):
"""Raised when a job (or a single variant) cannot be completed."""
class FilterNotFoundError(PipelineError):
"""Raised when no working filter could be picked after retries."""
@dataclass
class FilterAssets:
background_names: list[str]
foreground_names: list[str]
blend_modes: list[str]
commands: dict[str, str]
# --- low-level helpers, ported from make_random.py ------------------------
def strip_ansi(text: str) -> str:
return _ANSI_RE.sub("", text)
def _nice(cmd: list[str]) -> list[str]:
if config.NICE_LEVEL <= 0:
return cmd
return ["nice", "-n", str(config.NICE_LEVEL), *cmd]
def load_lines(path: Path) -> list[str]:
if not path.exists():
raise FileNotFoundError(f"Datei nicht gefunden: {path}")
lines = [line.strip() for line in path.read_text(encoding="utf-8").splitlines() if line.strip()]
if not lines:
raise ValueError(f"Datei ist leer: {path}")
return lines
def load_filter_commands() -> dict[str, str]:
data = json.loads(config.FILTERS_JSON.read_text(encoding="utf-8"))
return {item["plain_name"]: item["full_command"] for item in data}
def load_assets() -> FilterAssets:
"""Load background/foreground/blend-mode lists and cross-check against
filters.json, dropping any names without a known command (mirrors the
warning+filter behaviour in make_random.py's main())."""
bg_names = load_lines(config.BACKGROUND_FILE)
fg_names = load_lines(config.FOREGROUND_FILE)
blend_modes = load_lines(config.BLEND_MODES_FILE)
commands = load_filter_commands()
missing_bg = [name for name in bg_names if name not in commands]
missing_fg = [name for name in fg_names if name not in commands]
if missing_bg:
logger.warning("%d background filter names have no command, dropping them", len(missing_bg))
bg_names = [name for name in bg_names if name in commands]
if missing_fg:
logger.warning("%d foreground filter names have no command, dropping them", len(missing_fg))
fg_names = [name for name in fg_names if name in commands]
if not bg_names or not fg_names:
raise PipelineError("Keine gueltigen Filter in background/foreground Listen")
return FilterAssets(background_names=bg_names, foreground_names=fg_names, blend_modes=blend_modes, commands=commands)
def run_gmic(args: list[str], output_image: Path) -> tuple[bool, str, float]:
output_image.parent.mkdir(parents=True, exist_ok=True)
cmd = _nice([config.GMIC_BIN, *args, "-o", str(output_image)])
start = time.monotonic()
try:
proc = subprocess.run(cmd, capture_output=True, text=True, timeout=config.FILTER_TIMEOUT)
except subprocess.TimeoutExpired:
return False, f"Timeout nach {config.FILTER_TIMEOUT}s", time.monotonic() - start
elapsed = time.monotonic() - start
if proc.returncode != 0:
err = strip_ansi((proc.stderr or proc.stdout or "").strip())
err = err.splitlines()[-1] if err else f"Exit code {proc.returncode}"
return False, err[:500], elapsed
if not output_image.exists() or output_image.stat().st_size == 0:
return False, "Keine Ausgabedatei erzeugt", elapsed
return True, "", elapsed
def gmic_filter_args(image: Path, full_command: str) -> list[str]:
if " " in full_command:
name, args = full_command.split(" ", 1)
return [str(image), name, args]
return [str(image), full_command]
def apply_filter(image: Path, full_command: str, output_image: Path) -> tuple[bool, str]:
ok, err, _ = run_gmic(gmic_filter_args(image, full_command), output_image)
return ok, err
def apply_filter_chain(image: Path, commands: tuple[str, ...], output_image: Path, tmp_dir: Path, prefix: str) -> tuple[bool, str]:
current = image
for i, command in enumerate(commands):
target = output_image if i == len(commands) - 1 else tmp_dir / f"{prefix}_post_{i}.png"
ok, err = apply_filter(current, command, target)
if not ok:
return False, err
current = target
return True, ""
def blend_layers(base: Path, overlay: Path, mode: str, output_image: Path) -> tuple[bool, str]:
ok, err, _ = run_gmic([str(base), str(overlay), "blend", f"{mode},{config.BLEND_OPACITY}"], output_image)
return ok, err
def blend_layers_opacity(base: Path, overlay: Path, mode: str, opacity: str, output_image: Path) -> tuple[bool, str]:
ok, err, _ = run_gmic([str(base), str(overlay), "blend", f"{mode},{opacity}"], output_image)
return ok, err
def alpha_composite(base: Path, overlay: Path, output_image: Path) -> tuple[bool, str]:
ok, err, _ = run_gmic([str(base), str(overlay), "blend", "alpha"], output_image)
return ok, err
def run_rembg(input_image: Path, output_image: Path) -> tuple[bool, str]:
output_image.parent.mkdir(parents=True, exist_ok=True)
cmd = _nice([sys.executable, "-m", "app.rembg_cli", str(input_image), str(output_image)])
try:
proc = subprocess.run(cmd, capture_output=True, text=True, timeout=config.REMBG_TIMEOUT)
except subprocess.TimeoutExpired:
return False, f"rembg Timeout nach {config.REMBG_TIMEOUT}s"
if proc.returncode != 0 or not output_image.exists() or output_image.stat().st_size == 0:
lines = strip_ansi((proc.stderr or proc.stdout or "").strip()).splitlines()
return False, (lines[-1] if lines else f"Exit code {proc.returncode}")
return True, ""
def pick_working_filter(
names: list[str],
commands: dict[str, str],
image: Path,
tmp_dir: Path,
label: str,
rng: random.Random,
) -> tuple[str, str]:
tried: set[str] = set()
for _ in range(config.MAX_FILTER_ATTEMPTS):
candidates = [n for n in names if n not in tried]
if not candidates:
break
name = rng.choice(candidates)
tried.add(name)
command = commands.get(name)
if not command:
continue
probe = tmp_dir / f"probe_{label}_{uuid.uuid4().hex[:8]}.png"
ok, err = apply_filter(image, command, probe)
probe.unlink(missing_ok=True)
if ok:
return name, command
logger.info("skip %s filter %r: %s", label, name, err)
raise FilterNotFoundError(f"Kein funktionierender {label}-Filter nach {config.MAX_FILTER_ATTEMPTS} Versuchen")
# --- job-level orchestration ------------------------------------------------
def now_iso() -> str:
return datetime.now(timezone.utc).isoformat(timespec="seconds")
def new_job_id() -> str:
stamp = datetime.now().strftime("%Y%m%d-%H%M%S")
return f"{stamp}-{uuid.uuid4().hex[:6]}"
@dataclass
class JobPaths:
root: Path
original: Path
rembg: Path
intermediates: Path
variants: Path
manifest: Path
status: Path
def job_paths(job_id: str, original_suffix: str = ".jpg") -> JobPaths:
root = config.JOBS_DIR / job_id
return JobPaths(
root=root,
original=root / f"original{original_suffix}",
rembg=root / "rembg.png",
intermediates=root / "intermediates",
variants=root / "variants",
manifest=root / "manifest.json",
status=root / "status.json",
)
def read_status(job_id: str) -> dict[str, Any] | None:
paths = job_paths(job_id)
if not paths.status.exists():
return None
return json.loads(paths.status.read_text(encoding="utf-8"))
def write_status(paths: JobPaths, status: str, **extra: Any) -> None:
data = {}
if paths.status.exists():
try:
data = json.loads(paths.status.read_text(encoding="utf-8"))
except (OSError, ValueError):
data = {}
data["status"] = status
data["updated_at"] = now_iso()
data.setdefault("created_at", data["updated_at"])
data.update(extra)
paths.status.write_text(json.dumps(data, indent=2, ensure_ascii=False), encoding="utf-8")
def read_manifest(job_id: str) -> dict[str, Any] | None:
paths = job_paths(job_id)
if not paths.manifest.exists():
return None
return json.loads(paths.manifest.read_text(encoding="utf-8"))
def write_manifest(paths: JobPaths, manifest: dict[str, Any]) -> None:
manifest["updated_at"] = now_iso()
paths.manifest.write_text(json.dumps(manifest, indent=2, ensure_ascii=False), encoding="utf-8")
def next_variant_id(manifest: dict[str, Any]) -> str:
existing = {v["id"] for v in manifest.get("variants", [])}
i = len(manifest.get("variants", [])) + 1
while f"v{i}" in existing:
i += 1
return f"v{i}"
def compose_variant(
paths: JobPaths,
assets: FilterAssets,
*,
variant_id: str,
source: str,
bg_name: str | None = None,
bg_mode: str | None = None,
fg_name: str | None = None,
fg_mode: str | None = None,
opacity: str | None = None,
rng: random.Random | None = None,
) -> dict[str, Any]:
"""Compose one variant image from the job's cached original + rembg.
If bg_name/fg_name/bg_mode/fg_mode are given (remix path) they are used
directly. Otherwise a random working filter is picked with retries,
exactly like make_random.py's compose_one().
"""
rng = rng or random.Random()
tmp_dir = paths.intermediates
tmp_dir.mkdir(parents=True, exist_ok=True)
bg_mode = bg_mode or rng.choice(assets.blend_modes)
fg_mode = fg_mode or rng.choice(assets.blend_modes)
opacity = opacity or config.BLEND_OPACITY
if bg_name:
bg_command = assets.commands.get(bg_name)
if not bg_command:
raise PipelineError(f"Unbekannter Background-Filter: {bg_name}")
else:
bg_name, bg_command = pick_working_filter(assets.background_names, assets.commands, paths.original, tmp_dir, "bg", rng)
if fg_name:
fg_command = assets.commands.get(fg_name)
if not fg_command:
raise PipelineError(f"Unbekannter Foreground-Filter: {fg_name}")
else:
fg_name, fg_command = pick_working_filter(assets.foreground_names, assets.commands, paths.rembg, tmp_dir, "fg", rng)
p = {
"bg_filtered": tmp_dir / f"{variant_id}_bg_filtered.png",
"step1": tmp_dir / f"{variant_id}_bg_blend.png",
"step2": tmp_dir / f"{variant_id}_rembg_alpha.png",
"fg_filtered": tmp_dir / f"{variant_id}_fg_filtered.png",
"composed": tmp_dir / f"{variant_id}_composed.png",
"final": paths.variants / f"{variant_id}.png",
}
paths.variants.mkdir(parents=True, exist_ok=True)
ok, err = apply_filter(paths.original, bg_command, p["bg_filtered"])
if not ok:
raise PipelineError(f"Background-Filter fehlgeschlagen: {err}")
ok, err = blend_layers_opacity(paths.original, p["bg_filtered"], bg_mode, opacity, p["step1"])
if not ok:
raise PipelineError(f"Background-Blend fehlgeschlagen: {err}")
ok, err = alpha_composite(p["step1"], paths.rembg, p["step2"])
if not ok:
raise PipelineError(f"Rembg-Alpha fehlgeschlagen: {err}")
ok, err = apply_filter(paths.rembg, fg_command, p["fg_filtered"])
if not ok:
raise PipelineError(f"Foreground-Filter fehlgeschlagen: {err}")
ok, err = blend_layers_opacity(p["step2"], p["fg_filtered"], fg_mode, opacity, p["composed"])
if not ok:
raise PipelineError(f"Foreground-Blend fehlgeschlagen: {err}")
ok, err = apply_filter_chain(p["composed"], config.POST_FILTERS, p["final"], tmp_dir, variant_id)
if not ok:
raise PipelineError(f"Post-Processing fehlgeschlagen: {err}")
rel = lambda path: str(path.relative_to(paths.root))
return {
"id": variant_id,
"source": source,
"file": rel(p["final"]),
"background_filter": bg_name,
"background_command": bg_command,
"background_blend": bg_mode,
"foreground_filter": fg_name,
"foreground_command": fg_command,
"foreground_blend": fg_mode,
"blend_opacity": opacity,
"post_filters": list(config.POST_FILTERS),
"created_at": now_iso(),
"intermediates": {
"bg_filtered": rel(p["bg_filtered"]),
"bg_blend": rel(p["step1"]),
"rembg_alpha": rel(p["step2"]),
"fg_filtered": rel(p["fg_filtered"]),
"composed": rel(p["composed"]),
},
}
def process_job(job_id: str, source_path: Path, original_suffix: str) -> None:
"""Full pipeline for a freshly ingested inbox file: copy, rembg once,
generate OUTPUT_COUNT variants, write manifest + status.
`source_path` must already be a private copy (jobs/<id>/original.*) —
callers (worker.py) are responsible for copying out of inbox first, so
the inbox file itself is never touched here.
"""
paths = job_paths(job_id, original_suffix)
paths.root.mkdir(parents=True, exist_ok=True)
paths.intermediates.mkdir(parents=True, exist_ok=True)
paths.variants.mkdir(parents=True, exist_ok=True)
write_status(paths, "processing", source_file=str(source_path.name))
manifest: dict[str, Any] = {
"job_id": job_id,
"original_file": paths.original.name,
"rembg_file": paths.rembg.name,
"created_at": now_iso(),
"variants": [],
}
try:
ok, err = run_rembg(paths.original, paths.rembg)
if not ok:
raise PipelineError(f"rembg fehlgeschlagen: {err}")
assets = load_assets()
rng = random.Random()
variant_errors: list[str] = []
for i in range(1, config.OUTPUT_COUNT + 1):
variant_id = f"v{i}"
try:
entry = compose_variant(paths, assets, variant_id=variant_id, source="auto", rng=rng)
manifest["variants"].append(entry)
write_manifest(paths, manifest)
except PipelineError as exc:
logger.error("[%s] variant %s failed: %s", job_id, variant_id, exc)
variant_errors.append(f"{variant_id}: {exc}")
if not manifest["variants"]:
raise PipelineError("Keine Variante erfolgreich erzeugt: " + "; ".join(variant_errors))
write_status(paths, "done", variant_errors=variant_errors)
except PipelineError as exc:
logger.error("[%s] job failed: %s", job_id, exc)
write_status(paths, "error", error=str(exc))
except Exception as exc: # noqa: BLE001 - keep the worker loop alive
logger.exception("[%s] unexpected error", job_id)
write_status(paths, "error", error=f"Unerwarteter Fehler: {exc}")