From 68ed37eefb79dc1966dd46f74fc10882ce6f356e Mon Sep 17 00:00:00 2001 From: David Kifer Date: Thu, 25 Jun 2026 20:48:53 -0400 Subject: [PATCH] Fix subprocess pipe deadlock in upscaler and update static assets --- app/upscaler.py | 95 +++++++++++++++++++++++++++++-------------------- 1 file changed, 57 insertions(+), 38 deletions(-) diff --git a/app/upscaler.py b/app/upscaler.py index 1ca1746..8d8c4aa 100644 --- a/app/upscaler.py +++ b/app/upscaler.py @@ -136,15 +136,15 @@ class UpscaleJob: pass self._processes.clear() - def run_command(self, cmd: list, shell=False) -> subprocess.Popen: + 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=subprocess.PIPE, - stderr=subprocess.PIPE, + stdout=stdout, + stderr=stderr, text=True, shell=shell ) @@ -379,43 +379,59 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic upscale_cmd.append("-x") upscale_start_time = time.time() - p_upscale = job.run_command(upscale_cmd) + upscale_stdout_path = os.path.join(job_temp_dir, "upscale_stdout.log") + upscale_stderr_path = os.path.join(job_temp_dir, "upscale_stderr.log") - # 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% + 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) - # 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 + # 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% - # Format ETA - if eta_sec > 60: - eta_str = f"{int(eta_sec // 60)}m {int(eta_sec % 60)}s" + # 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 = 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) + 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 = "" - stdout, stderr = p_upscale.communicate() + 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) if job._is_cancelled: @@ -636,10 +652,13 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic job.update_status("failed", error=str(e)) on_progress_update(job.job_id, {"status": "failed", "error": str(e)}) finally: - # Clean up temp frames to save space + # Clean up temp frames to save space only if completed or cancelled try: - if os.path.exists(job_temp_dir): - shutil.rmtree(job_temp_dir) + 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}")