Add automated recover_job.py script for quick recovery on other systems
This commit is contained in:
Executable
+324
@@ -0,0 +1,324 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
recover_job.py — Automates the recovery of a deadlocked AI Video Upscaler job on any system.
|
||||
This script patches the upscaler code, kills the stuck server and upscaler processes
|
||||
without trigger cleanups (preserving temp frames), restarts the server, and resumes the job.
|
||||
|
||||
Usage:
|
||||
python3 recover_job.py
|
||||
"""
|
||||
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
import subprocess
|
||||
import signal
|
||||
import json
|
||||
import urllib.request
|
||||
|
||||
BASE_DIR = os.path.dirname(os.path.abspath(__file__))
|
||||
|
||||
def patch_upscaler():
|
||||
upscaler_path = os.path.join(BASE_DIR, "app", "upscaler.py")
|
||||
if not os.path.exists(upscaler_path):
|
||||
print(f"❌ Error: {upscaler_path} not found.")
|
||||
return False
|
||||
|
||||
with open(upscaler_path, "r") as f:
|
||||
content = f.read()
|
||||
|
||||
# 1. Update run_command signature and Popen call
|
||||
old_run_command = """ def run_command(self, cmd: list, shell=False) -> subprocess.Popen:
|
||||
with self._lock:
|
||||
if self._is_cancelled:
|
||||
raise InterruptedError("Job was cancelled")
|
||||
|
||||
p = subprocess.Popen(
|
||||
cmd,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
shell=shell
|
||||
)"""
|
||||
|
||||
new_run_command = """ def run_command(self, cmd: list, stdout=subprocess.PIPE, stderr=subprocess.PIPE, shell=False) -> subprocess.Popen:
|
||||
with self._lock:
|
||||
if self._is_cancelled:
|
||||
raise InterruptedError("Job was cancelled")
|
||||
|
||||
p = subprocess.Popen(
|
||||
cmd,
|
||||
stdout=stdout,
|
||||
stderr=stderr,
|
||||
text=True,
|
||||
shell=shell
|
||||
)"""
|
||||
|
||||
# 2. Update upscale command redirection
|
||||
old_upscale_loop = """ upscale_start_time = time.time()
|
||||
p_upscale = job.run_command(upscale_cmd)
|
||||
|
||||
# Monitor thread for output files
|
||||
while p_upscale.poll() is None:
|
||||
if job._is_cancelled:
|
||||
return
|
||||
|
||||
processed_files = len(os.listdir(output_frames_dir))
|
||||
progress_pct = 20.0 + (float(processed_files) / actual_total) * 60.0 # upscaling is 20% to 80%
|
||||
|
||||
# Estimate ETA
|
||||
elapsed = time.time() - upscale_start_time
|
||||
this_run_processed = processed_files - (actual_total - remaining_inputs)
|
||||
if this_run_processed > 0:
|
||||
sec_per_frame = elapsed / this_run_processed
|
||||
rem_frames = actual_total - processed_files
|
||||
eta_sec = rem_frames * sec_per_frame
|
||||
|
||||
# Format ETA
|
||||
if eta_sec > 60:
|
||||
eta_str = f"{int(eta_sec // 60)}m {int(eta_sec % 60)}s"
|
||||
else:
|
||||
eta_str = f"{int(eta_sec)}s"
|
||||
else:
|
||||
eta_str = "Calculating..."
|
||||
|
||||
job.update_status("upscaling", progress=progress_pct, current_frame=processed_files, eta=eta_str)
|
||||
on_progress_update(job.job_id, {
|
||||
"status": "upscaling",
|
||||
"progress": progress_pct,
|
||||
"current_frame": processed_files,
|
||||
"total_frames": actual_total,
|
||||
"eta": eta_str
|
||||
})
|
||||
time.sleep(0.5)
|
||||
|
||||
stdout, stderr = p_upscale.communicate()
|
||||
job.cleanup_process(p_upscale)"""
|
||||
|
||||
new_upscale_loop = """ upscale_start_time = time.time()
|
||||
upscale_stdout_path = os.path.join(job_temp_dir, "upscale_stdout.log")
|
||||
upscale_stderr_path = os.path.join(job_temp_dir, "upscale_stderr.log")
|
||||
|
||||
with open(upscale_stdout_path, "w") as f_out, open(upscale_stderr_path, "w") as f_err:
|
||||
p_upscale = job.run_command(upscale_cmd, stdout=f_out, stderr=f_err)
|
||||
|
||||
# Monitor thread for output files
|
||||
while p_upscale.poll() is None:
|
||||
if job._is_cancelled:
|
||||
return
|
||||
|
||||
processed_files = len(os.listdir(output_frames_dir))
|
||||
progress_pct = 20.0 + (float(processed_files) / actual_total) * 60.0 # upscaling is 20% to 80%
|
||||
|
||||
# Estimate ETA
|
||||
elapsed = time.time() - upscale_start_time
|
||||
this_run_processed = processed_files - (actual_total - remaining_inputs)
|
||||
if this_run_processed > 0:
|
||||
sec_per_frame = elapsed / this_run_processed
|
||||
rem_frames = actual_total - processed_files
|
||||
eta_sec = rem_frames * sec_per_frame
|
||||
|
||||
# Format ETA
|
||||
if eta_sec > 60:
|
||||
eta_str = f"{int(eta_sec // 60)}m {int(eta_sec % 60)}s"
|
||||
else:
|
||||
eta_str = f"{int(eta_sec)}s"
|
||||
else:
|
||||
eta_str = "Calculating..."
|
||||
|
||||
job.update_status("upscaling", progress=progress_pct, current_frame=processed_files, eta=eta_str)
|
||||
on_progress_update(job.job_id, {
|
||||
"status": "upscaling",
|
||||
"progress": progress_pct,
|
||||
"current_frame": processed_files,
|
||||
"total_frames": actual_total,
|
||||
"eta": eta_str
|
||||
})
|
||||
time.sleep(0.5)
|
||||
|
||||
# Read stdout/stderr from files
|
||||
if os.path.exists(upscale_stdout_path):
|
||||
with open(upscale_stdout_path, "r") as f_out:
|
||||
stdout = f_out.read()
|
||||
else:
|
||||
stdout = ""
|
||||
|
||||
if os.path.exists(upscale_stderr_path):
|
||||
with open(upscale_stderr_path, "r") as f_err:
|
||||
stderr = f_err.read()
|
||||
else:
|
||||
stderr = ""
|
||||
|
||||
job.cleanup_process(p_upscale)"""
|
||||
|
||||
# 3. Update finally block
|
||||
old_finally = """ finally:
|
||||
# Clean up temp frames to save space
|
||||
try:
|
||||
if os.path.exists(job_temp_dir):
|
||||
shutil.rmtree(job_temp_dir)
|
||||
except Exception as cleanup_err:
|
||||
print(f"Error during temp cleanup: {cleanup_err}")"""
|
||||
|
||||
new_finally = """ finally:
|
||||
# Clean up temp frames to save space only if completed or cancelled
|
||||
try:
|
||||
if job.status in ["completed", "cancelled"]:
|
||||
if os.path.exists(job_temp_dir):
|
||||
shutil.rmtree(job_temp_dir)
|
||||
else:
|
||||
print(f"Job {job.job_id} finished with status {job.status}. Preserving temp directory {job_temp_dir} for potential resume.")
|
||||
except Exception as cleanup_err:
|
||||
print(f"Error during temp cleanup: {cleanup_err}")"""
|
||||
|
||||
patched = False
|
||||
if old_run_command in content:
|
||||
content = content.replace(old_run_command, new_run_command)
|
||||
patched = True
|
||||
elif "def run_command(self, cmd: list, stdout=" in content:
|
||||
print("⚙️ run_command already patched.")
|
||||
else:
|
||||
print("⚠ Warning: Could not find old run_command signature to patch.")
|
||||
|
||||
if old_upscale_loop in content:
|
||||
content = content.replace(old_upscale_loop, new_upscale_loop)
|
||||
patched = True
|
||||
elif "upscale_stdout_path" in content:
|
||||
print("⚙️ upscale_loop already patched.")
|
||||
else:
|
||||
print("⚠ Warning: Could not find old upscale_loop to patch.")
|
||||
|
||||
if old_finally in content:
|
||||
content = content.replace(old_finally, new_finally)
|
||||
patched = True
|
||||
elif "job.status in [\"completed\", \"cancelled\"]" in content:
|
||||
print("⚙️ finally block already patched.")
|
||||
else:
|
||||
print("⚠ Warning: Could not find old finally block to patch.")
|
||||
|
||||
if patched:
|
||||
with open(upscaler_path, "w") as f:
|
||||
f.write(content)
|
||||
print("✅ app/upscaler.py successfully patched.")
|
||||
return True
|
||||
return False
|
||||
|
||||
def find_pids():
|
||||
server_pids = []
|
||||
realesrgan_pids = []
|
||||
|
||||
# Try reading .uvicorn.pid
|
||||
pid_file = os.path.join(BASE_DIR, ".uvicorn.pid")
|
||||
if os.path.exists(pid_file):
|
||||
try:
|
||||
with open(pid_file, "r") as f:
|
||||
server_pids.append(int(f.read().strip()))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Scan /proc for running processes matching names
|
||||
if os.path.exists("/proc"):
|
||||
for name in os.listdir("/proc"):
|
||||
if name.isdigit():
|
||||
pid = int(name)
|
||||
try:
|
||||
cmdline_path = f"/proc/{pid}/cmdline"
|
||||
if os.path.exists(cmdline_path):
|
||||
with open(cmdline_path, "r") as f:
|
||||
cmdline = f.read()
|
||||
if "realesrgan" in cmdline:
|
||||
realesrgan_pids.append(pid)
|
||||
elif "uvicorn" in cmdline or "app.main" in cmdline:
|
||||
if pid not in server_pids:
|
||||
server_pids.append(pid)
|
||||
except Exception:
|
||||
pass
|
||||
return list(set(server_pids)), list(set(realesrgan_pids))
|
||||
|
||||
def kill_processes(pids, name):
|
||||
for pid in pids:
|
||||
try:
|
||||
print(f"💥 Force killing {name} process (PID {pid}) using SIGKILL...")
|
||||
os.kill(pid, signal.SIGKILL)
|
||||
except OSError as e:
|
||||
print(f"⚠ Failed to kill PID {pid}: {e}")
|
||||
|
||||
def get_stuck_jobs():
|
||||
jobs_file = os.path.join(BASE_DIR, "jobs.json")
|
||||
if not os.path.exists(jobs_file):
|
||||
return []
|
||||
try:
|
||||
with open(jobs_file, "r") as f:
|
||||
data = json.load(f)
|
||||
stuck = []
|
||||
for jid, job in data.items():
|
||||
if job.get("status") in ["upscaling", "failed", "interrupted"]:
|
||||
stuck.append(jid)
|
||||
return stuck
|
||||
except Exception as e:
|
||||
print(f"⚠ Error reading jobs.json: {e}")
|
||||
return []
|
||||
|
||||
def main():
|
||||
print("🚀 Starting AI Video Upscaler Recovery Process...")
|
||||
|
||||
# 1. Patch the codebase
|
||||
patch_upscaler()
|
||||
|
||||
# 2. Identify running deadlocked processes
|
||||
server_pids, realesrgan_pids = find_pids()
|
||||
|
||||
# 3. Read jobs.json before killing processes to see which job is currently active
|
||||
stuck_jobs = get_stuck_jobs()
|
||||
print(f"📋 Detected active/stuck jobs to resume: {stuck_jobs}")
|
||||
|
||||
# 4. Terminate processes
|
||||
# CRITICAL: We use SIGKILL (kill -9) on the python process so that it exits immediately
|
||||
# and has NO opportunity to run its `finally` block which deletes the extracted/upscaled frames!
|
||||
if server_pids:
|
||||
kill_processes(server_pids, "server")
|
||||
else:
|
||||
print("ℹ No running server processes found.")
|
||||
|
||||
if realesrgan_pids:
|
||||
kill_processes(realesrgan_pids, "realesrgan")
|
||||
else:
|
||||
print("ℹ No running realesrgan processes found.")
|
||||
|
||||
# Give it a second to clear
|
||||
time.sleep(2)
|
||||
|
||||
# 5. Start the server again
|
||||
print("🚀 Launching the server...")
|
||||
try:
|
||||
# Check port in start.py parameters or default to 8000
|
||||
subprocess.Popen(
|
||||
[sys.executable, "start.py", "--no-browser"],
|
||||
cwd=BASE_DIR,
|
||||
start_new_session=True
|
||||
)
|
||||
print("✅ Server process launched in background.")
|
||||
except Exception as e:
|
||||
print(f"❌ Failed to launch server: {e}")
|
||||
sys.exit(1)
|
||||
|
||||
# Wait for server startup
|
||||
print("⏳ Waiting 5 seconds for server to initialize and bind...")
|
||||
time.sleep(5)
|
||||
|
||||
# 6. Call the resume API for the stuck jobs
|
||||
for jid in stuck_jobs:
|
||||
print(f"🔄 Requesting resume for job {jid}...")
|
||||
url = f"http://127.0.0.1:8000/api/upscale/resume/{jid}"
|
||||
try:
|
||||
req = urllib.request.Request(url, method="POST")
|
||||
with urllib.request.urlopen(req, timeout=10) as response:
|
||||
res_data = json.loads(response.read().decode('utf-8'))
|
||||
print(f"✅ Success response: {res_data}")
|
||||
except Exception as e:
|
||||
print(f"❌ Failed to resume job {jid}: {e}")
|
||||
|
||||
print("\n🎉 Recovery complete! The server is running and jobs have been resumed.")
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user