diff --git a/usenet-scripts/movetdarr.sh b/usenet-scripts/movetdarr.sh index f1a19f2..4e8024f 100644 --- a/usenet-scripts/movetdarr.sh +++ b/usenet-scripts/movetdarr.sh @@ -1,9 +1,9 @@ #!/bin/bash -# movetdarr.sh v5 — SABnzbd -> FileFlows with durable n8n handoff lease +# movetdarr.sh v4 — SABnzbd -> FileFlows with durable n8n handoff lease # Categories: serien4k->Sonarr UHD, serien->Sonarr FHD, # video4k->Radarr UHD, video->Radarr FHD. -# v5: BusyBox-kompatibler FD-Lock auf lokalem /config-Dateisystem. -# Der alte Lock unter /usenet lag auf NFS und konnte dauerhaft fehlschlagen. +# v4: flock-Serialisierung des Lease-Blocks gegen n8n-staticData-Race +# bei parallelen SAB-Aufträgen (Last-Writer-Wins verschluckt Jobs). DIR="$1" NZB_NAME="$2" @@ -17,7 +17,7 @@ SOURCE_ROOT="${SOURCE_ROOT:-/usenet/complete}" DEST_ROOT="${DEST_ROOT:-/tdarr/complete}" POLL_SECONDS="${POLL_SECONDS:-60}" MAX_POLLS="${MAX_POLLS:-1440}" # 24 hours -LEASE_LOCK="${LEASE_LOCK:-/config/.movetdarr.lease.lock}" +LEASE_LOCK="${LEASE_LOCK:-/usenet/scripts/.movetdarr.lease.lock}" log() { echo "$(date '+%d.%m.%Y %H:%M:%S'): $1" >> "$LOGFILE"; } @@ -42,7 +42,7 @@ notify_failure() { handoff "$payload" >/dev/null 2>&1 || true } -log "Start v5 - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR" +log "Start v4 - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR" if [ "$STATUS" != "0" ]; then log "Download fehlgeschlagen, ueberspringe" @@ -145,32 +145,21 @@ register_lease() { } if command -v flock >/dev/null 2>&1; then - # BusyBox flock hat kein -w: FD öffnen und -n mit Retry-Loop verwenden. - # Der FD bleibt bis nach start -> mv -> moved offen und wird vor dem - # lang laufenden Status-Polling ausdrücklich entsperrt. - exec 9>"$LEASE_LOCK" || { - log "FEHLER: Lease-Lock-Datei konnte nicht geoeffnet werden - kein Move" - exit 1 - } + # BusyBox flock hat kein -w (timeout) → -n (non-blocking) mit Retry-Loop LEASE_ACQUIRED=false for i in $(seq 1 900); do - if flock -n 9 2>/dev/null; then + if flock -n "$LEASE_LOCK" -c true 2>/dev/null; then LEASE_ACQUIRED=true break fi sleep 1 done if [ "$LEASE_ACQUIRED" != "true" ]; then - log "FEHLER: Lease-Lock nach 15 Minuten nicht erhalten - kein Move" - exec 9>&- - exit 1 + log "WARNUNG: Lease-Lock nach 15 Minuten nicht erhalten; fahre ohne Lock fort" fi register_lease - flock -u 9 - exec 9>&- else - log "FEHLER: flock fehlt - kein sicherer Handoff und kein Move" - exit 1 + register_lease fi # Ab hier: Status-Polling ohne Lock (lang laufend, stoert keine andere Registrierung). diff --git a/usenet-scripts/test_movetdarr.py b/usenet-scripts/test_movetdarr.py index cef101f..61a7488 100644 --- a/usenet-scripts/test_movetdarr.py +++ b/usenet-scripts/test_movetdarr.py @@ -52,7 +52,6 @@ def run_case(root, category="serien4k", status="0", name="Show.S01E01", target_v "DEST_ROOT": str(dst_root), "LOGFILE": str(root / "postprocess.log"), "HANDOFF_URL": f"http://127.0.0.1:{server.server_port}/media/handoff", - "LEASE_LOCK": str(root / "lease.lock"), "POLL_SECONDS": "0", "MAX_POLLS": "3", } diff --git a/usenet-scripts/test_movetdarr_parallel.py b/usenet-scripts/test_movetdarr_parallel.py deleted file mode 100644 index f3cad42..0000000 --- a/usenet-scripts/test_movetdarr_parallel.py +++ /dev/null @@ -1,110 +0,0 @@ -import json -import os -import subprocess -import tempfile -import threading -import time -from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer -from pathlib import Path - - -SCRIPT = Path(__file__).with_name("movetdarr.sh") -state_lock = threading.Lock() -next_job = 0 -open_leases = 0 -max_open_leases = 0 -job_ids = {} - - -class Handler(BaseHTTPRequestHandler): - def do_POST(self): - global next_job, open_leases, max_open_leases - data = json.loads(self.rfile.read(int(self.headers.get("Content-Length", "0")))) - action = data.get("action") - if action == "start": - with state_lock: - next_job += 1 - job_id = f"job-{next_job}" - job_ids[data["directory"]] = job_id - open_leases += 1 - max_open_leases = max(max_open_leases, open_leases) - time.sleep(0.3) - result = {"ok": True, "jobId": job_id, "state": "registered"} - elif action == "moved": - with state_lock: - open_leases -= 1 - result = {"ok": True, "jobId": data["jobId"], "state": "processing"} - elif action == "status": - result = {"ok": True, "jobId": data["jobId"], "state": "success"} - else: - result = {"ok": False, "state": "failed"} - body = json.dumps(result).encode() - self.send_response(200) - self.send_header("Content-Type", "application/json") - self.send_header("Content-Length", str(len(body))) - self.end_headers() - self.wfile.write(body) - - def log_message(self, *_args): - pass - - -def start_script(root, name, server_port): - source_root = root / "source" - source = source_root / "serien" / name - source.mkdir(parents=True) - (source / f"{name}.mkv").write_bytes(b"video") - env = os.environ | { - "SOURCE_ROOT": str(source_root), - "DEST_ROOT": str(root / "dest"), - "LOGFILE": str(root / "postprocess.log"), - "LEASE_LOCK": str(root / "lease.lock"), - "HANDOFF_URL": f"http://127.0.0.1:{server_port}/media/handoff", - "POLL_SECONDS": "0", - "MAX_POLLS": "2", - } - return subprocess.Popen( - ["bash", str(SCRIPT), str(source), name + ".nzb", name, "", "serien", "", "0"], - env=env, - text=True, - stdout=subprocess.PIPE, - stderr=subprocess.PIPE, - ) - - -def test_parallel_lease_blocks_are_serialized(): - global next_job, open_leases, max_open_leases, job_ids - next_job = 0 - open_leases = 0 - max_open_leases = 0 - job_ids = {} - server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) - thread = threading.Thread(target=server.serve_forever, daemon=True) - thread.start() - try: - with tempfile.TemporaryDirectory() as td: - root = Path(td) - first = start_script(root, "Show.S01E01", server.server_port) - second = start_script(root, "Show.S01E02", server.server_port) - first_result = first.communicate(timeout=20) - second_result = second.communicate(timeout=20) - assert first.returncode == 0, first_result - assert second.returncode == 0, second_result - assert len(set(job_ids.values())) == 2 - assert max_open_leases == 1, "start→move→moved blocks overlapped" - finally: - server.shutdown() - server.server_close() - - -def test_lock_is_held_on_a_file_descriptor_until_register_returns(): - source = SCRIPT.read_text() - assert 'exec 9>"$LEASE_LOCK"' in source - assert 'flock -n 9' in source - assert 'flock -u 9' in source - - -if __name__ == "__main__": - test_lock_is_held_on_a_file_descriptor_until_register_returns() - test_parallel_lease_blocks_are_serialized() - print("PASS: parallel start→move→moved blocks are serialized")