feat: keep SAB jobs leased until Arr import #1
2 changed files with 241 additions and 26 deletions
|
|
@ -1,55 +1,163 @@
|
||||||
#!/bin/bash
|
#!/bin/bash
|
||||||
# movetdarr.sh v2 — SABnzbd PostProcessing: Download -> FileFlows-Eingang (/tdarr/complete)
|
# movetdarr.sh v3 — SABnzbd -> FileFlows with durable n8n handoff lease
|
||||||
# Fixes 05.09.2026:
|
# Categories: serien4k->Sonarr UHD, serien->Sonarr FHD,
|
||||||
# - Vor dem mv: bestehendes Zielverzeichnis ohne Video-Datei (FileFlows-Leichen) wird entfernt
|
# video4k->Radarr UHD, video->Radarr FHD.
|
||||||
# - Fehler beim mv werden GEMELDET (exit 1) statt still zu scheitern
|
|
||||||
# - Ziel mit existierender Video-Datei: Download behalten + WARNUNG (kein Datenverlust)
|
|
||||||
DIR="$1"
|
DIR="$1"
|
||||||
|
NZB_NAME="$2"
|
||||||
|
CLEAN_NAME="$3"
|
||||||
CATEGORY="$5"
|
CATEGORY="$5"
|
||||||
STATUS="$7"
|
STATUS="$7"
|
||||||
|
|
||||||
LOGFILE="/usenet/scripts/postprocess.log"
|
LOGFILE="${LOGFILE:-/usenet/scripts/postprocess.log}"
|
||||||
|
HANDOFF_URL="${HANDOFF_URL:-http://10.5.85.2:8888/media/handoff}"
|
||||||
|
SOURCE_ROOT="${SOURCE_ROOT:-/usenet/complete}"
|
||||||
|
DEST_ROOT="${DEST_ROOT:-/tdarr/complete}"
|
||||||
|
POLL_SECONDS="${POLL_SECONDS:-60}"
|
||||||
|
MAX_POLLS="${MAX_POLLS:-1440}" # 24 hours
|
||||||
|
|
||||||
log() { echo "$(date '+%d.%m.%Y %H:%M:%S'): $1" >> "$LOGFILE"; }
|
log() { echo "$(date '+%d.%m.%Y %H:%M:%S'): $1" >> "$LOGFILE"; }
|
||||||
|
|
||||||
log "Start - Kategorie=$CATEGORY Status=$STATUS Dir=$DIR"
|
json_payload() {
|
||||||
|
python3 -c 'import json,sys; print(json.dumps(dict(zip(sys.argv[1::2],sys.argv[2::2]))))' "$@"
|
||||||
|
}
|
||||||
|
|
||||||
|
handoff() {
|
||||||
|
curl --silent --show-error --fail-with-body \
|
||||||
|
--connect-timeout 5 --max-time 20 \
|
||||||
|
-H 'Content-Type: application/json' \
|
||||||
|
--data "$1" "$HANDOFF_URL"
|
||||||
|
}
|
||||||
|
|
||||||
|
json_field() {
|
||||||
|
python3 -c 'import json,sys; print(json.load(sys.stdin).get(sys.argv[1], ""))' "$1"
|
||||||
|
}
|
||||||
|
|
||||||
|
notify_failure() {
|
||||||
|
[ -n "${JOB_ID:-}" ] || return 0
|
||||||
|
payload="$(json_payload action fail jobId "$JOB_ID" reason "$1")"
|
||||||
|
handoff "$payload" >/dev/null 2>&1 || true
|
||||||
|
}
|
||||||
|
|
||||||
|
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"
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
|
||||||
|
case "$CATEGORY" in
|
||||||
|
serien4k|serien|video4k|video) ;;
|
||||||
|
*)
|
||||||
|
log "FEHLER: nicht unterstuetzte Kategorie '$CATEGORY' - Quelle bleibt unangetastet"
|
||||||
|
exit 1
|
||||||
|
;;
|
||||||
|
esac
|
||||||
|
|
||||||
|
EXPECTED_PREFIX="${SOURCE_ROOT}/${CATEGORY}/"
|
||||||
|
case "$DIR/" in
|
||||||
|
"$EXPECTED_PREFIX"*) ;;
|
||||||
|
*)
|
||||||
|
log "FEHLER: Quellpfad passt nicht zur Kategorie - Quelle bleibt unangetastet"
|
||||||
|
exit 1
|
||||||
|
;;
|
||||||
|
esac
|
||||||
|
|
||||||
if [ ! -d "$DIR" ]; then
|
if [ ! -d "$DIR" ]; then
|
||||||
log "Quellverzeichnis existiert nicht: $DIR"
|
log "Quellverzeichnis existiert nicht: $DIR"
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
|
||||||
DEST="/tdarr/complete/${CATEGORY}"
|
# Lease registrieren, solange der SAB-Job und sein Quellpfad noch sichtbar sind.
|
||||||
BASENAME="$(basename "$DIR")"
|
START_PAYLOAD="$(json_payload action start category "$CATEGORY" directory "$DIR" release "$NZB_NAME" cleanName "$CLEAN_NAME")"
|
||||||
TARGET="$DEST/$BASENAME"
|
START_RESPONSE="$(handoff "$START_PAYLOAD" 2>>"$LOGFILE")" || {
|
||||||
|
log "FEHLER: n8n-Handoff konnte nicht registriert werden - kein Move"
|
||||||
mkdir -p "$DEST"
|
exit 1
|
||||||
|
}
|
||||||
# Leichen-Handling: existiert das Ziel schon?
|
JOB_ID="$(printf '%s' "$START_RESPONSE" | json_field jobId 2>>"$LOGFILE")"
|
||||||
if [ -d "$TARGET" ]; then
|
START_OK="$(printf '%s' "$START_RESPONSE" | json_field ok 2>>"$LOGFILE")"
|
||||||
# Enthaelt es Video-Dateien?
|
if [ "$START_OK" != "True" ] && [ "$START_OK" != "true" ] || [ -z "$JOB_ID" ]; then
|
||||||
if ls "$TARGET"/*.mkv "$TARGET"/*.mp4 "$TARGET"/*.avi >/dev/null 2>&1; then
|
log "FEHLER: ungueltige Handoff-Antwort - kein Move"
|
||||||
log "WARNUNG: Ziel $TARGET enthaelt bereits Video-Dateien - NICHT ueberschrieben, Download bleibt in $DIR"
|
|
||||||
exit 2
|
|
||||||
fi
|
|
||||||
log "Entferne FileFlows-Restverzeichnis (keine Videos enthalten): $TARGET"
|
|
||||||
rm -rf "$TARGET"
|
|
||||||
if [ $? -ne 0 ]; then
|
|
||||||
log "FEHLER: konnte Restverzeichnis $TARGET nicht entfernen"
|
|
||||||
exit 1
|
|
||||||
fi
|
|
||||||
fi
|
|
||||||
|
|
||||||
mv "$DIR" "$DEST/"
|
|
||||||
|
|
||||||
if [ $? -eq 0 ]; then
|
|
||||||
log "Erfolgreich verschoben nach $TARGET"
|
|
||||||
exit 0
|
|
||||||
else
|
|
||||||
log "FEHLER beim Verschieben nach $TARGET (Quelle bleibt: $DIR)"
|
|
||||||
exit 1
|
exit 1
|
||||||
fi
|
fi
|
||||||
|
log "Handoff registriert: Job=$JOB_ID"
|
||||||
|
|
||||||
|
DEST="${DEST_ROOT}/${CATEGORY}"
|
||||||
|
BASENAME="$(basename "$DIR")"
|
||||||
|
TARGET="$DEST/$BASENAME"
|
||||||
|
mkdir -p "$DEST" || {
|
||||||
|
log "FEHLER: Zielbasis konnte nicht erstellt werden: $DEST"
|
||||||
|
notify_failure "mkdir_failed"
|
||||||
|
exit 1
|
||||||
|
}
|
||||||
|
|
||||||
|
if [ -d "$TARGET" ]; then
|
||||||
|
if find "$TARGET" -maxdepth 1 -type f \( -iname '*.mkv' -o -iname '*.mp4' -o -iname '*.avi' \) -print -quit | grep -q .; then
|
||||||
|
log "WARNUNG: Ziel $TARGET enthaelt bereits Video-Dateien - kein Ueberschreiben"
|
||||||
|
notify_failure "target_contains_video"
|
||||||
|
exit 2
|
||||||
|
fi
|
||||||
|
log "Entferne FileFlows-Restverzeichnis ohne Video: $TARGET"
|
||||||
|
rm -rf -- "$TARGET" || {
|
||||||
|
log "FEHLER: Restverzeichnis konnte nicht entfernt werden"
|
||||||
|
notify_failure "stale_target_remove_failed"
|
||||||
|
exit 1
|
||||||
|
}
|
||||||
|
fi
|
||||||
|
|
||||||
|
if ! mv -- "$DIR" "$DEST/"; then
|
||||||
|
log "FEHLER beim Verschieben nach $TARGET (Quelle bleibt: $DIR)"
|
||||||
|
notify_failure "move_failed"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
log "Erfolgreich verschoben nach $TARGET; SAB-Job bleibt bis Arr-Import aktiv"
|
||||||
|
|
||||||
|
MOVED_PAYLOAD="$(json_payload action moved jobId "$JOB_ID")"
|
||||||
|
MOVED_OK=false
|
||||||
|
for attempt in 1 2 3 4 5; do
|
||||||
|
if MOVED_RESPONSE="$(handoff "$MOVED_PAYLOAD" 2>>"$LOGFILE")"; then
|
||||||
|
MOVED_STATE="$(printf '%s' "$MOVED_RESPONSE" | json_field state 2>>"$LOGFILE")"
|
||||||
|
if [ "$MOVED_STATE" = "processing" ] || [ "$MOVED_STATE" = "success" ]; then
|
||||||
|
MOVED_OK=true
|
||||||
|
break
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
log "Handoff-Move-Bestaetigung Versuch $attempt/5 fehlgeschlagen"
|
||||||
|
sleep 15
|
||||||
|
done
|
||||||
|
if [ "$MOVED_OK" != "true" ]; then
|
||||||
|
log "FEHLER: Move konnte n8n nicht bestaetigt werden; Datei bleibt fuer Recovery in $TARGET"
|
||||||
|
notify_failure "move_confirmation_failed"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
STATUS_PAYLOAD="$(json_payload action status jobId "$JOB_ID")"
|
||||||
|
for poll in $(seq 1 "$MAX_POLLS"); do
|
||||||
|
if RESPONSE="$(handoff "$STATUS_PAYLOAD" 2>>"$LOGFILE")"; then
|
||||||
|
STATE="$(printf '%s' "$RESPONSE" | json_field state 2>>"$LOGFILE")"
|
||||||
|
case "$STATE" in
|
||||||
|
success)
|
||||||
|
log "Arr-Import verifiziert: Job=$JOB_ID - SAB darf jetzt abschliessen/aufraeumen"
|
||||||
|
exit 0
|
||||||
|
;;
|
||||||
|
failed|timeout)
|
||||||
|
REASON="$(printf '%s' "$RESPONSE" | json_field reason 2>>"$LOGFILE")"
|
||||||
|
log "FEHLER: Handoff Job=$JOB_ID State=$STATE Reason=$REASON"
|
||||||
|
exit 1
|
||||||
|
;;
|
||||||
|
registered|processing)
|
||||||
|
if [ $((poll % 10)) -eq 0 ]; then
|
||||||
|
log "Warte auf FileFlows/Arr-Import: Job=$JOB_ID State=$STATE Poll=$poll/$MAX_POLLS"
|
||||||
|
fi
|
||||||
|
;;
|
||||||
|
*) log "WARNUNG: unbekannter Handoff-State '$STATE' fuer Job=$JOB_ID" ;;
|
||||||
|
esac
|
||||||
|
else
|
||||||
|
log "WARNUNG: Handoff-Status temporaer nicht erreichbar; Job=$JOB_ID Poll=$poll/$MAX_POLLS"
|
||||||
|
fi
|
||||||
|
sleep "$POLL_SECONDS"
|
||||||
|
done
|
||||||
|
|
||||||
|
log "FEHLER: lokaler 24h-Timeout fuer Job=$JOB_ID"
|
||||||
|
notify_failure "local_poll_timeout"
|
||||||
|
exit 1
|
||||||
|
|
|
||||||
107
usenet-scripts/test_movetdarr.py
Normal file
107
usenet-scripts/test_movetdarr.py
Normal file
|
|
@ -0,0 +1,107 @@
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
import subprocess
|
||||||
|
import tempfile
|
||||||
|
import threading
|
||||||
|
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
SCRIPT = Path(__file__).with_name("movetdarr.sh")
|
||||||
|
actions = []
|
||||||
|
status_polls = 0
|
||||||
|
|
||||||
|
|
||||||
|
class Handler(BaseHTTPRequestHandler):
|
||||||
|
def do_POST(self):
|
||||||
|
global status_polls
|
||||||
|
data = json.loads(self.rfile.read(int(self.headers.get("Content-Length", "0"))))
|
||||||
|
actions.append(data)
|
||||||
|
action = data.get("action")
|
||||||
|
if action == "start":
|
||||||
|
result = {"ok": True, "jobId": "test-job-1234", "state": "registered"}
|
||||||
|
elif action == "moved":
|
||||||
|
result = {"ok": True, "jobId": "test-job-1234", "state": "processing"}
|
||||||
|
elif action == "status":
|
||||||
|
status_polls += 1
|
||||||
|
result = {"ok": True, "jobId": "test-job-1234", "state": "success" if status_polls >= 2 else "processing"}
|
||||||
|
else:
|
||||||
|
result = {"ok": True, "jobId": "test-job-1234", "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 run_case(root, category="serien4k", status="0", name="Show.S01E01", target_video=False):
|
||||||
|
src_root = root / "source"
|
||||||
|
dst_root = root / "dest"
|
||||||
|
source = src_root / category / name
|
||||||
|
source.mkdir(parents=True)
|
||||||
|
(source / "episode.mkv").write_bytes(b"video")
|
||||||
|
target = dst_root / category / name
|
||||||
|
if target_video:
|
||||||
|
target.mkdir(parents=True)
|
||||||
|
(target / "existing.mkv").write_bytes(b"existing")
|
||||||
|
env = os.environ | {
|
||||||
|
"SOURCE_ROOT": str(src_root),
|
||||||
|
"DEST_ROOT": str(dst_root),
|
||||||
|
"LOGFILE": str(root / "postprocess.log"),
|
||||||
|
"HANDOFF_URL": f"http://127.0.0.1:{server.server_port}/media/handoff",
|
||||||
|
"POLL_SECONDS": "0",
|
||||||
|
"MAX_POLLS": "3",
|
||||||
|
}
|
||||||
|
result = subprocess.run(
|
||||||
|
["bash", str(SCRIPT), str(source), name + ".nzb", name, "", category, "", status],
|
||||||
|
env=env,
|
||||||
|
text=True,
|
||||||
|
capture_output=True,
|
||||||
|
timeout=20,
|
||||||
|
)
|
||||||
|
return result, source, target
|
||||||
|
|
||||||
|
|
||||||
|
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:
|
||||||
|
actions.clear()
|
||||||
|
status_polls = 0
|
||||||
|
result, source, target = run_case(Path(td))
|
||||||
|
assert result.returncode == 0, (result.stdout, result.stderr)
|
||||||
|
assert not source.exists() and (target / "episode.mkv").exists()
|
||||||
|
assert [a["action"] for a in actions] == ["start", "moved", "status", "status"]
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as td:
|
||||||
|
actions.clear()
|
||||||
|
status_polls = 0
|
||||||
|
result, source, _ = run_case(Path(td), category="other")
|
||||||
|
assert result.returncode == 1 and source.exists() and actions == []
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as td:
|
||||||
|
actions.clear()
|
||||||
|
status_polls = 0
|
||||||
|
result, source, target = run_case(Path(td), target_video=True)
|
||||||
|
assert result.returncode == 2 and source.exists() and (target / "existing.mkv").exists()
|
||||||
|
assert [a["action"] for a in actions] == ["start", "fail"]
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory() as td:
|
||||||
|
actions.clear()
|
||||||
|
status_polls = 0
|
||||||
|
root = Path(td)
|
||||||
|
stale = root / "dest" / "video" / "Movie.2026"
|
||||||
|
stale.mkdir(parents=True)
|
||||||
|
(stale / "leftover.nfo").write_text("stale")
|
||||||
|
result, source, target = run_case(root, category="video", name="Movie.2026")
|
||||||
|
assert result.returncode == 0 and not source.exists()
|
||||||
|
assert not (target / "leftover.nfo").exists() and (target / "episode.mkv").exists()
|
||||||
|
|
||||||
|
print("PASS: success lease, category fail-closed, no-overwrite, stale-directory cleanup")
|
||||||
|
finally:
|
||||||
|
server.shutdown()
|
||||||
|
server.server_close()
|
||||||
Loading…
Add table
Add a link
Reference in a new issue