fix: lock meta RMW and harden concurrent remix paths
CI / lint-and-test (pull_request) Failing after 9s
CI / lint-and-test (pull_request) Failing after 9s
Prevent worker/web clobbering of meta variants via flock and merge-by-id, make filter timeouts thread-local, harden job-id/stem sanitization, migrate TemplateResponse API, and remove the compare-bg experiment. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+96
-28
@@ -15,6 +15,7 @@ import re
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
import uuid
|
||||
from collections.abc import Iterator
|
||||
@@ -25,14 +26,14 @@ from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from . import config
|
||||
from .io_utils import read_json, write_json_atomic
|
||||
from .io_utils import file_lock, read_json, write_json_atomic
|
||||
|
||||
logger = logging.getLogger("livef12.pipeline")
|
||||
|
||||
_ANSI_RE = re.compile(r"\x1b\[[0-9;]*m")
|
||||
|
||||
# Optional override for run_gmic timeout (used by the long-retry pass).
|
||||
_filter_timeout_override: int | None = None
|
||||
# Per-thread override for run_gmic timeout (used by the long-retry pass).
|
||||
_filter_timeout_state = threading.local()
|
||||
|
||||
|
||||
class PipelineError(RuntimeError):
|
||||
@@ -69,19 +70,19 @@ def _nice(cmd: list[str]) -> list[str]:
|
||||
|
||||
|
||||
def _active_filter_timeout() -> int:
|
||||
return _filter_timeout_override if _filter_timeout_override is not None else config.FILTER_TIMEOUT
|
||||
override = getattr(_filter_timeout_state, "override", None)
|
||||
return override if override is not None else config.FILTER_TIMEOUT
|
||||
|
||||
|
||||
@contextmanager
|
||||
def filter_timeout(seconds: int) -> Iterator[None]:
|
||||
"""Temporarily override FILTER_TIMEOUT for run_gmic (long-retry pass)."""
|
||||
global _filter_timeout_override
|
||||
previous = _filter_timeout_override
|
||||
_filter_timeout_override = seconds
|
||||
previous = getattr(_filter_timeout_state, "override", None)
|
||||
_filter_timeout_state.override = seconds
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
_filter_timeout_override = previous
|
||||
_filter_timeout_state.override = previous
|
||||
|
||||
|
||||
def _magick_bin() -> str:
|
||||
@@ -306,8 +307,10 @@ _UNSAFE_STEM_RE = re.compile(r"[^\w.\-]+", re.UNICODE)
|
||||
|
||||
def sanitize_stem(name: str) -> str:
|
||||
"""Safe basename stem for flat files (no path separators / junk)."""
|
||||
stem = Path(name).stem if name else ""
|
||||
stem = stem.replace("/", "_").replace("\\", "_").strip().strip(".")
|
||||
# Replace separators before Path.stem so "a/b" → "a_b", not "b".
|
||||
stem = (name or "").replace("/", "_").replace("\\", "_")
|
||||
stem = Path(stem).stem
|
||||
stem = stem.strip().strip(".")
|
||||
stem = _UNSAFE_STEM_RE.sub("_", stem).strip("._")
|
||||
return stem or "photo"
|
||||
|
||||
@@ -354,11 +357,39 @@ def read_meta(stem: str) -> dict[str, Any] | None:
|
||||
return data if isinstance(data, dict) else None
|
||||
|
||||
|
||||
def _write_meta(paths: JobPaths, data: dict[str, Any]) -> None:
|
||||
def _meta_lock_path(stem: str) -> Path:
|
||||
return config.META_DIR / f"{sanitize_stem(stem)}.lock"
|
||||
|
||||
|
||||
@contextmanager
|
||||
def meta_lock(stem: str) -> Iterator[None]:
|
||||
"""Exclusive lock for meta/{stem}.json read-modify-write."""
|
||||
with file_lock(_meta_lock_path(stem)):
|
||||
yield
|
||||
|
||||
|
||||
def _write_meta_unlocked(paths: JobPaths, data: dict[str, Any]) -> None:
|
||||
data["updated_at"] = now_iso()
|
||||
write_json_atomic(paths.meta, data)
|
||||
|
||||
|
||||
def _merge_variants(existing: list[Any], incoming: list[Any]) -> list[dict[str, Any]]:
|
||||
"""Union variants by id; incoming wins on conflict; preserve discovery order."""
|
||||
by_id: dict[str, dict[str, Any]] = {}
|
||||
order: list[str] = []
|
||||
for group in (existing, incoming):
|
||||
for item in group:
|
||||
if not isinstance(item, dict):
|
||||
continue
|
||||
vid = item.get("id")
|
||||
if not isinstance(vid, str) or not vid:
|
||||
continue
|
||||
if vid not in by_id:
|
||||
order.append(vid)
|
||||
by_id[vid] = item
|
||||
return [by_id[vid] for vid in order]
|
||||
|
||||
|
||||
def read_status(job_id: str) -> dict[str, Any] | None:
|
||||
data = read_meta(job_id)
|
||||
if not data:
|
||||
@@ -375,13 +406,14 @@ def read_status(job_id: str) -> dict[str, Any] | None:
|
||||
|
||||
|
||||
def write_status(paths: JobPaths, status: str, **extra: Any) -> None:
|
||||
data = read_meta(paths.stem) or {}
|
||||
data["job_id"] = paths.stem
|
||||
data["source_stem"] = paths.stem
|
||||
data["status"] = status
|
||||
data.setdefault("created_at", now_iso())
|
||||
data.update(extra)
|
||||
_write_meta(paths, data)
|
||||
with meta_lock(paths.stem):
|
||||
data = read_meta(paths.stem) or {}
|
||||
data["job_id"] = paths.stem
|
||||
data["source_stem"] = paths.stem
|
||||
data["status"] = status
|
||||
data.setdefault("created_at", now_iso())
|
||||
data.update(extra)
|
||||
_write_meta_unlocked(paths, data)
|
||||
|
||||
|
||||
def read_manifest(job_id: str) -> dict[str, Any] | None:
|
||||
@@ -400,16 +432,52 @@ def read_manifest(job_id: str) -> dict[str, Any] | None:
|
||||
|
||||
|
||||
def write_manifest(paths: JobPaths, manifest: dict[str, Any]) -> None:
|
||||
data = read_meta(paths.stem) or {}
|
||||
data["job_id"] = paths.stem
|
||||
data["source_stem"] = manifest.get("source_stem") or paths.stem
|
||||
data["original_file"] = manifest.get("original_file")
|
||||
data["rembg_file"] = manifest.get("rembg_file")
|
||||
data["variants"] = manifest.get("variants") or []
|
||||
if "created_at" in manifest:
|
||||
data.setdefault("created_at", manifest["created_at"])
|
||||
data.setdefault("status", data.get("status", "processing"))
|
||||
_write_meta(paths, data)
|
||||
"""Upsert manifest fields; merge variants by id so concurrent writers do not clobber."""
|
||||
with meta_lock(paths.stem):
|
||||
data = read_meta(paths.stem) or {}
|
||||
data["job_id"] = paths.stem
|
||||
if manifest.get("source_stem"):
|
||||
data["source_stem"] = manifest["source_stem"]
|
||||
else:
|
||||
data.setdefault("source_stem", paths.stem)
|
||||
if "original_file" in manifest and manifest["original_file"] is not None:
|
||||
data["original_file"] = manifest["original_file"]
|
||||
if "rembg_file" in manifest and manifest["rembg_file"] is not None:
|
||||
data["rembg_file"] = manifest["rembg_file"]
|
||||
data["variants"] = _merge_variants(data.get("variants") or [], manifest.get("variants") or [])
|
||||
if "created_at" in manifest:
|
||||
data.setdefault("created_at", manifest["created_at"])
|
||||
data.setdefault("status", data.get("status", "processing"))
|
||||
_write_meta_unlocked(paths, data)
|
||||
|
||||
|
||||
def append_manifest_variant(
|
||||
paths: JobPaths,
|
||||
entry: dict[str, Any],
|
||||
*,
|
||||
original_file: str | None = None,
|
||||
rembg_file: str | None = None,
|
||||
source_stem: str | None = None,
|
||||
created_at: str | None = None,
|
||||
) -> None:
|
||||
"""Append one variant under meta lock (re-reads disk so concurrent updates survive)."""
|
||||
with meta_lock(paths.stem):
|
||||
data = read_meta(paths.stem) or {}
|
||||
data["job_id"] = paths.stem
|
||||
if source_stem:
|
||||
data["source_stem"] = source_stem
|
||||
else:
|
||||
data.setdefault("source_stem", paths.stem)
|
||||
if original_file:
|
||||
data["original_file"] = original_file
|
||||
if rembg_file:
|
||||
data["rembg_file"] = rembg_file
|
||||
if created_at:
|
||||
data.setdefault("created_at", created_at)
|
||||
data.setdefault("created_at", now_iso())
|
||||
data.setdefault("status", data.get("status", "done"))
|
||||
data["variants"] = _merge_variants(data.get("variants") or [], [entry])
|
||||
_write_meta_unlocked(paths, data)
|
||||
|
||||
|
||||
def list_job_ids() -> list[str]:
|
||||
|
||||
Reference in New Issue
Block a user