Compare commits
4 commits
bb748741cd
...
64892fadb1
| Author | SHA1 | Date | |
|---|---|---|---|
| 64892fadb1 | |||
| 35efa23cde | |||
| e0d59f213c | |||
| 32d68201fd |
3 changed files with 131 additions and 9 deletions
|
|
@ -1,9 +1,9 @@
|
|||
#!/bin/bash
|
||||
# movetdarr.sh v4 — SABnzbd -> FileFlows with durable n8n handoff lease
|
||||
# movetdarr.sh v5 — SABnzbd -> FileFlows with durable n8n handoff lease
|
||||
# Categories: serien4k->Sonarr UHD, serien->Sonarr FHD,
|
||||
# video4k->Radarr UHD, video->Radarr FHD.
|
||||
# v4: flock-Serialisierung des Lease-Blocks gegen n8n-staticData-Race
|
||||
# bei parallelen SAB-Aufträgen (Last-Writer-Wins verschluckt Jobs).
|
||||
# v5: BusyBox-kompatibler FD-Lock auf lokalem /config-Dateisystem.
|
||||
# Der alte Lock unter /usenet lag auf NFS und konnte dauerhaft fehlschlagen.
|
||||
|
||||
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:-/usenet/scripts/.movetdarr.lease.lock}"
|
||||
LEASE_LOCK="${LEASE_LOCK:-/config/.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 v4 - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR"
|
||||
log "Start v5 - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR"
|
||||
|
||||
if [ "$STATUS" != "0" ]; then
|
||||
log "Download fehlgeschlagen, ueberspringe"
|
||||
|
|
@ -145,21 +145,32 @@ register_lease() {
|
|||
}
|
||||
|
||||
if command -v flock >/dev/null 2>&1; then
|
||||
# BusyBox flock hat kein -w (timeout) → -n (non-blocking) mit Retry-Loop
|
||||
# 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
|
||||
}
|
||||
LEASE_ACQUIRED=false
|
||||
for i in $(seq 1 900); do
|
||||
if flock -n "$LEASE_LOCK" -c true 2>/dev/null; then
|
||||
if flock -n 9 2>/dev/null; then
|
||||
LEASE_ACQUIRED=true
|
||||
break
|
||||
fi
|
||||
sleep 1
|
||||
done
|
||||
if [ "$LEASE_ACQUIRED" != "true" ]; then
|
||||
log "WARNUNG: Lease-Lock nach 15 Minuten nicht erhalten; fahre ohne Lock fort"
|
||||
log "FEHLER: Lease-Lock nach 15 Minuten nicht erhalten - kein Move"
|
||||
exec 9>&-
|
||||
exit 1
|
||||
fi
|
||||
register_lease
|
||||
flock -u 9
|
||||
exec 9>&-
|
||||
else
|
||||
register_lease
|
||||
log "FEHLER: flock fehlt - kein sicherer Handoff und kein Move"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# Ab hier: Status-Polling ohne Lock (lang laufend, stoert keine andere Registrierung).
|
||||
|
|
|
|||
|
|
@ -52,6 +52,7 @@ 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",
|
||||
}
|
||||
|
|
|
|||
110
usenet-scripts/test_movetdarr_parallel.py
Normal file
110
usenet-scripts/test_movetdarr_parallel.py
Normal file
|
|
@ -0,0 +1,110 @@
|
|||
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")
|
||||
Loading…
Add table
Add a link
Reference in a new issue