Compare commits
No commits in common. "main" and "feat/media-handoff-orchestrator-20260905" have entirely different histories.
main
...
feat/media
4 changed files with 54 additions and 243 deletions
19
compose.yaml
19
compose.yaml
|
|
@ -1,21 +1,4 @@
|
||||||
services:
|
services:
|
||||||
# Synchronisiert das Git-versionierte Postprocessing vor jedem SAB-Start
|
|
||||||
# atomar auf den bestehenden NFS-Skriptpfad und setzt das Execute-Bit.
|
|
||||||
script-sync:
|
|
||||||
image: alpine:3.22
|
|
||||||
container_name: sabnzbd-script-sync
|
|
||||||
volumes:
|
|
||||||
- ./usenet-scripts:/source:ro
|
|
||||||
- usenet:/usenet
|
|
||||||
command:
|
|
||||||
- /bin/sh
|
|
||||||
- -ec
|
|
||||||
- |
|
|
||||||
mkdir -p /usenet/scripts
|
|
||||||
install -m 0755 /source/movetdarr.sh /usenet/scripts/movetdarr.sh.new
|
|
||||||
mv -f /usenet/scripts/movetdarr.sh.new /usenet/scripts/movetdarr.sh
|
|
||||||
restart: "no"
|
|
||||||
|
|
||||||
# WireGuard VPN-Exit -> Hetzner wg2 (dedizierter Tunnel). SABnzbd teilt diesen Netzstack,
|
# WireGuard VPN-Exit -> Hetzner wg2 (dedizierter Tunnel). SABnzbd teilt diesen Netzstack,
|
||||||
# damit der gesamte Usenet-Traffic ueber die Hetzner-IP rausgeht.
|
# damit der gesamte Usenet-Traffic ueber die Hetzner-IP rausgeht.
|
||||||
wireguard:
|
wireguard:
|
||||||
|
|
@ -50,8 +33,6 @@ services:
|
||||||
depends_on:
|
depends_on:
|
||||||
wireguard:
|
wireguard:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
script-sync:
|
|
||||||
condition: service_completed_successfully
|
|
||||||
environment:
|
environment:
|
||||||
- PUID=1000
|
- PUID=1000
|
||||||
- PGID=1000
|
- PGID=1000
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,7 @@
|
||||||
#!/bin/bash
|
#!/bin/bash
|
||||||
# movetdarr.sh v5 — SABnzbd -> FileFlows with durable n8n handoff lease
|
# movetdarr.sh v3 — SABnzbd -> FileFlows with durable n8n handoff lease
|
||||||
# Categories: serien4k->Sonarr UHD, serien->Sonarr FHD,
|
# Categories: serien4k->Sonarr UHD, serien->Sonarr FHD,
|
||||||
# video4k->Radarr UHD, video->Radarr 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.
|
|
||||||
|
|
||||||
DIR="$1"
|
DIR="$1"
|
||||||
NZB_NAME="$2"
|
NZB_NAME="$2"
|
||||||
|
|
@ -17,7 +15,6 @@ SOURCE_ROOT="${SOURCE_ROOT:-/usenet/complete}"
|
||||||
DEST_ROOT="${DEST_ROOT:-/tdarr/complete}"
|
DEST_ROOT="${DEST_ROOT:-/tdarr/complete}"
|
||||||
POLL_SECONDS="${POLL_SECONDS:-60}"
|
POLL_SECONDS="${POLL_SECONDS:-60}"
|
||||||
MAX_POLLS="${MAX_POLLS:-1440}" # 24 hours
|
MAX_POLLS="${MAX_POLLS:-1440}" # 24 hours
|
||||||
LEASE_LOCK="${LEASE_LOCK:-/config/.movetdarr.lease.lock}"
|
|
||||||
|
|
||||||
log() { echo "$(date '+%d.%m.%Y %H:%M:%S'): $1" >> "$LOGFILE"; }
|
log() { echo "$(date '+%d.%m.%Y %H:%M:%S'): $1" >> "$LOGFILE"; }
|
||||||
|
|
||||||
|
|
@ -42,7 +39,7 @@ notify_failure() {
|
||||||
handoff "$payload" >/dev/null 2>&1 || true
|
handoff "$payload" >/dev/null 2>&1 || true
|
||||||
}
|
}
|
||||||
|
|
||||||
log "Start v5 - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR"
|
log "Start v3 - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR"
|
||||||
|
|
||||||
if [ "$STATUS" != "0" ]; then
|
if [ "$STATUS" != "0" ]; then
|
||||||
log "Download fehlgeschlagen, ueberspringe"
|
log "Download fehlgeschlagen, ueberspringe"
|
||||||
|
|
@ -50,7 +47,7 @@ if [ "$STATUS" != "0" ]; then
|
||||||
fi
|
fi
|
||||||
|
|
||||||
case "$CATEGORY" in
|
case "$CATEGORY" in
|
||||||
serien4k|serien|serienen|video4k|video|videoen) ;;
|
serien4k|serien|video4k|video) ;;
|
||||||
*)
|
*)
|
||||||
log "FEHLER: nicht unterstuetzte Kategorie '$CATEGORY' - Quelle bleibt unangetastet"
|
log "FEHLER: nicht unterstuetzte Kategorie '$CATEGORY' - Quelle bleibt unangetastet"
|
||||||
exit 1
|
exit 1
|
||||||
|
|
@ -71,16 +68,8 @@ if [ ! -d "$DIR" ]; then
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
|
||||||
EXPECTED_FILES="$(find "$DIR" -type f \( -iname '*.mkv' -o -iname '*.mp4' -o -iname '*.avi' \) | wc -l | tr -d ' ')"
|
# Lease registrieren, solange der SAB-Job und sein Quellpfad noch sichtbar sind.
|
||||||
if [ -z "$EXPECTED_FILES" ] || [ "$EXPECTED_FILES" -lt 1 ] || [ "$EXPECTED_FILES" -gt 100 ]; then
|
START_PAYLOAD="$(json_payload action start category "$CATEGORY" directory "$DIR" release "$NZB_NAME" cleanName "$CLEAN_NAME")"
|
||||||
log "FEHLER: ungueltige Anzahl Videodateien: ${EXPECTED_FILES:-0} - kein Move"
|
|
||||||
exit 1
|
|
||||||
fi
|
|
||||||
|
|
||||||
# Lease registrieren und Move+Bestaetigung unter exclusivem Lock, damit parallele
|
|
||||||
# SAB-Auftraege sich nicht gegenseitig den n8n-staticData-Job zerstoeren (v4).
|
|
||||||
register_lease() {
|
|
||||||
START_PAYLOAD="$(json_payload action start category "$CATEGORY" directory "$DIR" release "$NZB_NAME" cleanName "$CLEAN_NAME" expectedFiles "$EXPECTED_FILES")"
|
|
||||||
START_RESPONSE="$(handoff "$START_PAYLOAD" 2>>"$LOGFILE")" || {
|
START_RESPONSE="$(handoff "$START_PAYLOAD" 2>>"$LOGFILE")" || {
|
||||||
log "FEHLER: n8n-Handoff konnte nicht registriert werden - kein Move"
|
log "FEHLER: n8n-Handoff konnte nicht registriert werden - kein Move"
|
||||||
exit 1
|
exit 1
|
||||||
|
|
@ -92,9 +81,8 @@ register_lease() {
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
log "Handoff registriert: Job=$JOB_ID"
|
log "Handoff registriert: Job=$JOB_ID"
|
||||||
export JOB_ID
|
|
||||||
|
|
||||||
DEST="${DEST_ROOT:?DEST_ROOT muss gesetzt sein}/${CATEGORY}"
|
DEST="${DEST_ROOT}/${CATEGORY}"
|
||||||
BASENAME="$(basename "$DIR")"
|
BASENAME="$(basename "$DIR")"
|
||||||
TARGET="$DEST/$BASENAME"
|
TARGET="$DEST/$BASENAME"
|
||||||
mkdir -p "$DEST" || {
|
mkdir -p "$DEST" || {
|
||||||
|
|
@ -142,38 +130,7 @@ register_lease() {
|
||||||
notify_failure "move_confirmation_failed"
|
notify_failure "move_confirmation_failed"
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
}
|
|
||||||
|
|
||||||
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
|
|
||||||
}
|
|
||||||
LEASE_ACQUIRED=false
|
|
||||||
for i in $(seq 1 900); do
|
|
||||||
if flock -n 9 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
|
|
||||||
fi
|
|
||||||
register_lease
|
|
||||||
flock -u 9
|
|
||||||
exec 9>&-
|
|
||||||
else
|
|
||||||
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).
|
|
||||||
STATUS_PAYLOAD="$(json_payload action status jobId "$JOB_ID")"
|
STATUS_PAYLOAD="$(json_payload action status jobId "$JOB_ID")"
|
||||||
for poll in $(seq 1 "$MAX_POLLS"); do
|
for poll in $(seq 1 "$MAX_POLLS"); do
|
||||||
if RESPONSE="$(handoff "$STATUS_PAYLOAD" 2>>"$LOGFILE")"; then
|
if RESPONSE="$(handoff "$STATUS_PAYLOAD" 2>>"$LOGFILE")"; then
|
||||||
|
|
|
||||||
|
|
@ -52,7 +52,6 @@ def run_case(root, category="serien4k", status="0", name="Show.S01E01", target_v
|
||||||
"DEST_ROOT": str(dst_root),
|
"DEST_ROOT": str(dst_root),
|
||||||
"LOGFILE": str(root / "postprocess.log"),
|
"LOGFILE": str(root / "postprocess.log"),
|
||||||
"HANDOFF_URL": f"http://127.0.0.1:{server.server_port}/media/handoff",
|
"HANDOFF_URL": f"http://127.0.0.1:{server.server_port}/media/handoff",
|
||||||
"LEASE_LOCK": str(root / "lease.lock"),
|
|
||||||
"POLL_SECONDS": "0",
|
"POLL_SECONDS": "0",
|
||||||
"MAX_POLLS": "3",
|
"MAX_POLLS": "3",
|
||||||
}
|
}
|
||||||
|
|
@ -77,7 +76,6 @@ try:
|
||||||
assert result.returncode == 0, (result.stdout, result.stderr)
|
assert result.returncode == 0, (result.stdout, result.stderr)
|
||||||
assert not source.exists() and (target / "episode.mkv").exists()
|
assert not source.exists() and (target / "episode.mkv").exists()
|
||||||
assert [a["action"] for a in actions] == ["start", "moved", "status", "status"]
|
assert [a["action"] for a in actions] == ["start", "moved", "status", "status"]
|
||||||
assert actions[0]["expectedFiles"] == "1"
|
|
||||||
|
|
||||||
with tempfile.TemporaryDirectory() as td:
|
with tempfile.TemporaryDirectory() as td:
|
||||||
actions.clear()
|
actions.clear()
|
||||||
|
|
@ -103,22 +101,7 @@ try:
|
||||||
assert result.returncode == 0 and not source.exists()
|
assert result.returncode == 0 and not source.exists()
|
||||||
assert not (target / "leftover.nfo").exists() and (target / "episode.mkv").exists()
|
assert not (target / "leftover.nfo").exists() and (target / "episode.mkv").exists()
|
||||||
|
|
||||||
# 16.09.2026: serienen/videoen (englische Arr-Instanzen sonarrEN:8991 /
|
print("PASS: success lease, category fail-closed, no-overwrite, stale-directory cleanup")
|
||||||
# radarrEN:7880 mit Root-Foldern /data/FHD/serienen bzw. /data/FHD/videoen)
|
|
||||||
# fehlten in der Kategorie-Whitelist. Folge: Mutiny.2026 x2 blieben am
|
|
||||||
# 08./09.09.2026 mit "nicht unterstuetzte Kategorie 'videoen'" in
|
|
||||||
# /usenet/complete liegen und wurden nie importiert.
|
|
||||||
for extra_category in ("serienen", "videoen"):
|
|
||||||
with tempfile.TemporaryDirectory() as td:
|
|
||||||
actions.clear()
|
|
||||||
status_polls = 0
|
|
||||||
result, source, target = run_case(Path(td), category=extra_category)
|
|
||||||
assert result.returncode == 0, (extra_category, result.stdout, result.stderr)
|
|
||||||
assert not source.exists(), extra_category
|
|
||||||
assert (target / "episode.mkv").exists(), extra_category
|
|
||||||
assert [a["action"] for a in actions] == ["start", "moved", "status", "status"], extra_category
|
|
||||||
|
|
||||||
print("PASS: success lease, category fail-closed, no-overwrite, stale-directory cleanup, EN-Kategorien")
|
|
||||||
finally:
|
finally:
|
||||||
server.shutdown()
|
server.shutdown()
|
||||||
server.server_close()
|
server.server_close()
|
||||||
|
|
|
||||||
|
|
@ -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")
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue