diff --git a/recover_job.py b/recover_job.py new file mode 100755 index 0000000..77b7e45 --- /dev/null +++ b/recover_job.py @@ -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()