From 7c89ba3491f80d872724c9d22a220781c0cd4ae8 Mon Sep 17 00:00:00 2001 From: Thales <> Date: Fri, 17 Jul 2026 01:33:20 +0100 Subject: [PATCH] fix(registry): persist race on Windows; recover metadata-less done jobs persist() is called concurrently from the pipeline thread, API threads, and the sweep loop, all sharing one temp path. Two writers could collide, and on Windows os.replace over a file another writer holds open raises an uncaught PermissionError. The write+replace now happens under the existing lock with a unique temp name per call (the _ensure_cached_mp3 pattern), best-effort like the settings store. _recover_done_job required metadata.json, which is written after status flips to done -- a crash in that window left a complete stems dir permanently unrecoverable. Such dirs now recover with a placeholder title, and a minimal metadata.json is written immediately so the next restart takes the normal path (self-healing, not a lasting special case). The stems-present requirement is unchanged. Closes #281 Closes #284 --- app/core/registry.py | 45 ++++++++++++++++------ tests/test_registry_persistence.py | 61 ++++++++++++++++++++++++++++-- 2 files changed, 92 insertions(+), 14 deletions(-) diff --git a/app/core/registry.py b/app/core/registry.py index 06df6822..39e981e4 100644 --- a/app/core/registry.py +++ b/app/core/registry.py @@ -4,6 +4,7 @@ import logging import subprocess import threading +import uuid from pathlib import Path from app.core.config import JOB_ID_RE, STEM_NAMES @@ -70,12 +71,20 @@ def _migrate(data: dict) -> dict: def persist(jobs_dir: Path) -> None: - """Persist terminal jobs so completed library entries survive restarts.""" + """Persist terminal jobs so completed library entries survive restarts. + + Callers run on the pipeline thread, API threads, and the sweep loop + concurrently, so the write+replace happens under the lock with a unique + temp name per call (#281) -- a shared temp path let two writers collide, + and on Windows os.replace over a file another writer holds open raises + PermissionError. Best-effort like the settings store: a failed persist + logs and returns rather than killing the caller.""" try: jobs_dir.mkdir(parents=True, exist_ok=True) except OSError: logger.warning("cannot create jobs dir %s; skipping persist", jobs_dir, exc_info=True) return + path = jobs_dir / _REGISTRY_FILE with _lock: records = [ job.to_record() @@ -83,10 +92,14 @@ def persist(jobs_dir: Path) -> None: if job.status in _TERMINAL ] payload = json.dumps({"version": REGISTRY_VERSION, "jobs": records}, indent=2) + "\n" - path = jobs_dir / _REGISTRY_FILE - tmp = path.with_suffix(".json.tmp") - tmp.write_text(payload, encoding="utf-8") - tmp.replace(path) + tmp = jobs_dir / f".registry.{uuid.uuid4().hex}.tmp" + try: + tmp.write_text(payload, encoding="utf-8") + tmp.replace(path) + except OSError: + logger.warning("could not persist registry to %s", path, exc_info=True) + finally: + tmp.unlink(missing_ok=True) def restore(jobs_dir: Path) -> None: @@ -137,13 +150,23 @@ def _recover_done_job(job_dir: Path) -> Job | None: mix_url = f"/api/jobs/{job_dir.name}/stems/mix.wav" selected = [stem["name"] for stem in stems if stem["name"] in STEM_NAMES] or list(STEM_NAMES) meta_path = job_dir / "metadata.json" - if not meta_path.is_file(): - return None meta: dict = {} - try: - meta = json.loads(meta_path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError): - pass + if meta_path.is_file(): + try: + meta = json.loads(meta_path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + pass + else: + # Crash window (#284): the process died between status=done and the + # metadata write, leaving a complete stems dir that used to be + # unrecoverable. Recover with a placeholder title and write a minimal + # metadata.json immediately, so the NEXT restart takes the normal + # path -- self-healing, not a permanent special case. + meta = {"title": f"Recovered track {job_dir.name[:6]}"} + try: + meta_path.write_text(json.dumps(meta, indent=2) + "\n", encoding="utf-8") + except OSError: + logger.warning("could not write recovery metadata for %s", job_dir.name, exc_info=True) return Job( id=job_dir.name, status="done", diff --git a/tests/test_registry_persistence.py b/tests/test_registry_persistence.py index 36840753..a1279eab 100644 --- a/tests/test_registry_persistence.py +++ b/tests/test_registry_persistence.py @@ -59,14 +59,69 @@ def test_restore_recovers_orphan_done_job_from_stems(tmp_path: Path): assert {stem["name"] for stem in restored.stems} == {"vocals", "drums"} -def test_restore_skips_orphan_without_metadata(tmp_path: Path): - stems_dir = tmp_path / "abcdefabcde0" / "stems" +def test_restore_recovers_orphan_without_metadata(tmp_path: Path): + """#284: a crash between status=done and the metadata write used to leave + a complete stems dir permanently unrecoverable. Now it comes back with a + placeholder title, and a minimal metadata.json is written so the next + restart takes the normal recovery path (self-healing).""" + job_dir = tmp_path / "abcdefabcde0" + stems_dir = job_dir / "stems" stems_dir.mkdir(parents=True) (stems_dir / "vocals.wav").write_bytes(b"RIFF") restore_registry(tmp_path) - assert "abcdefabcde0" not in _jobs + restored = _jobs["abcdefabcde0"] + assert restored.status == "done" + assert restored.title == "Recovered track abcdef" + assert {stem["name"] for stem in restored.stems} == {"vocals"} + # Self-healed: metadata.json now exists with the placeholder title. + meta = json.loads((job_dir / "metadata.json").read_text(encoding="utf-8")) + assert meta["title"] == "Recovered track abcdef" + + +def test_restore_still_ignores_dir_without_stems(tmp_path: Path): + """The stems requirement stays: an empty/partial job dir is not a track.""" + (tmp_path / "abcdefabcde1" / "stems").mkdir(parents=True) # no WAVs + (tmp_path / "abcdefabcde2").mkdir(parents=True) # no stems dir at all + + restore_registry(tmp_path) + + assert "abcdefabcde1" not in _jobs + assert "abcdefabcde2" not in _jobs + + +def test_persist_concurrent_writers_no_corruption(tmp_path: Path): + """#281: pipeline thread, API threads, and the sweep all call persist() + concurrently. A shared temp path let writers collide (PermissionError on + Windows os.replace). Hammer it from threads: no exception, valid JSON, + no stray temp files.""" + import threading + + for i in range(5): + job = Job(id=f"abcdefabcd{i:02x}", status="done", title=f"t{i}") + _jobs[job.id] = job + + errors: list[Exception] = [] + + def hammer(): + try: + for _ in range(30): + persist_registry(tmp_path) + except Exception as e: # pragma: no cover - the failure being tested + errors.append(e) + + threads = [threading.Thread(target=hammer) for _ in range(8)] + for t in threads: + t.start() + for t in threads: + t.join() + + assert errors == [] + data = json.loads((tmp_path / "registry.json").read_text(encoding="utf-8")) + assert len(data["jobs"]) == 5 + assert not list(tmp_path.glob("*.tmp")), "no temp files may be left behind" + assert not list(tmp_path.glob(".registry.*")), "no temp files may be left behind" def test_restored_job_serves_stems(tmp_path: Path, monkeypatch):