Fix subprocess pipe deadlock in upscaler and update static assets

This commit is contained in:
2026-06-25 20:48:56 -04:00
parent cbd071c872
commit 68ed37eefb
+57 -38
View File
@@ -136,15 +136,15 @@ class UpscaleJob:
pass pass
self._processes.clear() 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: with self._lock:
if self._is_cancelled: if self._is_cancelled:
raise InterruptedError("Job was cancelled") raise InterruptedError("Job was cancelled")
p = subprocess.Popen( p = subprocess.Popen(
cmd, cmd,
stdout=subprocess.PIPE, stdout=stdout,
stderr=subprocess.PIPE, stderr=stderr,
text=True, text=True,
shell=shell shell=shell
) )
@@ -379,43 +379,59 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
upscale_cmd.append("-x") upscale_cmd.append("-x")
upscale_start_time = time.time() 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 with open(upscale_stdout_path, "w") as f_out, open(upscale_stderr_path, "w") as f_err:
while p_upscale.poll() is None: p_upscale = job.run_command(upscale_cmd, stdout=f_out, stderr=f_err)
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 # Monitor thread for output files
elapsed = time.time() - upscale_start_time while p_upscale.poll() is None:
this_run_processed = processed_files - (actual_total - remaining_inputs) if job._is_cancelled:
if this_run_processed > 0: return
sec_per_frame = elapsed / this_run_processed
rem_frames = actual_total - processed_files processed_files = len(os.listdir(output_frames_dir))
eta_sec = rem_frames * sec_per_frame progress_pct = 20.0 + (float(processed_files) / actual_total) * 60.0 # upscaling is 20% to 80%
# Format ETA # Estimate ETA
if eta_sec > 60: elapsed = time.time() - upscale_start_time
eta_str = f"{int(eta_sec // 60)}m {int(eta_sec % 60)}s" 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: else:
eta_str = f"{int(eta_sec)}s" eta_str = "Calculating..."
else:
eta_str = "Calculating..." job.update_status("upscaling", progress=progress_pct, current_frame=processed_files, eta=eta_str)
on_progress_update(job.job_id, {
job.update_status("upscaling", progress=progress_pct, current_frame=processed_files, eta=eta_str) "status": "upscaling",
on_progress_update(job.job_id, { "progress": progress_pct,
"status": "upscaling", "current_frame": processed_files,
"progress": progress_pct, "total_frames": actual_total,
"current_frame": processed_files, "eta": eta_str
"total_frames": actual_total, })
"eta": eta_str time.sleep(0.5)
})
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) job.cleanup_process(p_upscale)
if job._is_cancelled: 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)) job.update_status("failed", error=str(e))
on_progress_update(job.job_id, {"status": "failed", "error": str(e)}) on_progress_update(job.job_id, {"status": "failed", "error": str(e)})
finally: finally:
# Clean up temp frames to save space # Clean up temp frames to save space only if completed or cancelled
try: try:
if os.path.exists(job_temp_dir): if job.status in ["completed", "cancelled"]:
shutil.rmtree(job_temp_dir) 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: except Exception as cleanup_err:
print(f"Error during temp cleanup: {cleanup_err}") print(f"Error during temp cleanup: {cleanup_err}")