Implement batch processing, pause/resume, custom temp directory, and detailed pipeline logs

This commit is contained in:
2026-06-29 08:04:35 -04:00
parent d1ea029afb
commit d35a48bf64
5 changed files with 507 additions and 71 deletions
+41 -5
View File
@@ -26,8 +26,9 @@ class UpscaleJob:
denoise: bool = False, sharpen: bool = False, interpolation: bool = False,
webhook_url: str = None, transcode_format: str = "mp4", is_preview: bool = False,
ai_face_restoration: bool = False, ai_rife_interpolation: bool = False,
ai_audio_denoise: bool = False):
ai_audio_denoise: bool = False, temp_dir: str = None):
self.job_id = job_id
self.temp_dir = temp_dir
self.video_path = video_path
self.model = model
self.scale = scale
@@ -79,9 +80,10 @@ class UpscaleJob:
self.start_time = None
self.output_file = None
# Track processes to allow cancellation
# Track processes to allow cancellation/pause
self._processes = []
self._is_cancelled = False
self._is_paused = False
self._lock = threading.Lock()
def to_dict(self) -> dict:
@@ -106,6 +108,7 @@ class UpscaleJob:
setattr(job, k, v)
job._processes = []
job._is_cancelled = False
job._is_paused = data.get('_is_paused', False) or (data.get('status') == 'paused')
job._lock = threading.Lock()
return job
@@ -137,10 +140,32 @@ class UpscaleJob:
pass
self._processes.clear()
def pause(self):
with self._lock:
if self.status in ["queued", "pending"]:
self.status = "paused"
self.eta = "Paused"
elif self.status in ["analyzing", "extracting", "upscaling", "restoring_faces", "interpolating", "assembling"]:
self._is_paused = True
self.status = "paused"
self.eta = "Paused"
for p in self._processes:
try:
p.terminate()
p.wait(timeout=2)
except Exception:
try:
p.kill()
except Exception:
pass
self._processes.clear()
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")
if getattr(self, "_is_paused", False):
raise InterruptedError("Job was paused")
p = subprocess.Popen(
cmd,
@@ -240,6 +265,7 @@ def upscale_image_file(input_path: str, output_path: str, model: str, scale: int
return False
def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dict[str, Any]], None]):
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] starting pipeline. Video path: {job.video_path}, Model: {job.model}, Scale: {job.scale}")
job.start_time = time.time()
job.update_status("analyzing", progress=5)
@@ -254,7 +280,8 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
fps = info["fps"]
# Create job temp directories
job_temp_dir = os.path.join(TEMP_DIR, job.job_id)
base_temp = job.temp_dir if getattr(job, "temp_dir", None) else TEMP_DIR
job_temp_dir = os.path.join(base_temp, job.job_id)
input_frames_dir = os.path.join(job_temp_dir, "input_frames")
output_frames_dir = os.path.join(job_temp_dir, "output_frames")
@@ -299,6 +326,7 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
if not skip_extraction:
job.update_status("extracting", progress=10)
on_progress_update(job.job_id, {"status": "extracting", "progress": 10})
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] extracting frames. Command: {' '.join(extract_cmd)}")
# High quality JPG frames to balance disk usage and speed
extract_cmd = ["ffmpeg", "-y"]
@@ -327,6 +355,7 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
# Count actual frames extracted
extracted_files = sorted([f for f in os.listdir(input_frames_dir) if f.startswith("frame_")])
actual_total = len(extracted_files)
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] extraction completed. Extracted {actual_total} frames.")
if actual_total == 0:
raise RuntimeError("No frames extracted from video")
@@ -384,6 +413,7 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
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:
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] starting realesrgan upscaling. Model: {job.model}, Tile size: {current_tile_size}. Command: {' '.join(upscale_cmd)}")
p_upscale = job.run_command(upscale_cmd, stdout=f_out, stderr=f_err)
# Monitor thread for output files
@@ -479,7 +509,7 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
gfpgan_installed = importlib.util.find_spec("gfpgan") is not None
if gfpgan_installed:
print(f"Job {job.job_id}: GFPGAN detected. Running Face Restoration...")
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] starting GFPGAN face restoration.")
restored_dir = os.path.join(job_temp_dir, "restored_frames")
os.makedirs(restored_dir, exist_ok=True)
@@ -521,6 +551,7 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
if getattr(job, "ai_rife_interpolation", False):
rife_bin = os.path.join(BASE_DIR, "rife-bin", "rife-ncnn-vulkan")
if os.path.isfile(rife_bin):
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] starting RIFE frame interpolation.")
job.update_status("interpolating", progress=83.0)
on_progress_update(job.job_id, {"status": "interpolating", "progress": 83.0})
os.makedirs(rife_frames_dir, exist_ok=True)
@@ -630,6 +661,7 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
out_filepath
])
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] starting video assembly. Command: {' '.join(assemble_cmd)}")
p_assemble = job.run_command(assemble_cmd)
stdout, stderr = p_assemble.communicate()
job.cleanup_process(p_assemble)
@@ -645,13 +677,17 @@ def run_upscale_pipeline(job: UpscaleJob, on_progress_update: Callable[[str, Dic
"eta": "Done",
"output_file": out_filename
})
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] pipeline completed successfully. Output file: {job.output_file}")
except Exception as e:
import traceback
traceback.print_exc()
if not job._is_cancelled:
if not job._is_cancelled and not getattr(job, "_is_paused", False) and job.status != "paused":
job.update_status("failed", error=str(e))
on_progress_update(job.job_id, {"status": "failed", "error": str(e)})
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] pipeline failed. Error: {e}")
else:
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] [Job {job.job_id}] pipeline halted. Status: {job.status}")
finally:
# Clean up temp frames to save space only if completed or cancelled
try: