Refine FFMPEGVideoWriter to prevent pipe blocking
Building on @thehhmdb's fix in PR #418 which identified the stderr pipe blocking issue. This refinement simplifies the solution by redirecting stderr to DEVNULL and adding stdin.flush() to prevent buffering deadlocks. Improvements: - Simpler implementation without threading complexity - Zero memory overhead - Better error messages for debugging - Maintains fix for the original hanging issue Co-authored-by: thehhmdb <thehhmdb@users.noreply.github.com> Fixes numz/ComfyUI-SeedVR2_VideoUpscaler#418
This commit is contained in:
+17
-46
@@ -109,7 +109,6 @@ import torch
|
||||
import cv2
|
||||
import numpy as np
|
||||
import subprocess
|
||||
import threading
|
||||
import shutil
|
||||
|
||||
# Project imports
|
||||
@@ -172,25 +171,22 @@ class FFMPEGVideoWriter:
|
||||
['ffmpeg', '-y', '-f', 'rawvideo', '-pix_fmt', 'rgb24',
|
||||
'-s', f'{width}x{height}', '-r', str(fps), '-i', '-',
|
||||
'-c:v', codec, '-pix_fmt', pix_fmt, '-preset', 'medium', '-crf', '12', path],
|
||||
stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=subprocess.PIPE
|
||||
stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL
|
||||
)
|
||||
# Background reader to continuously consume stderr so ffmpeg cannot block
|
||||
self._stderr_buffer = bytearray()
|
||||
self._stderr_thread = threading.Thread(target=self._read_stderr, daemon=True)
|
||||
self._stderr_thread.start()
|
||||
|
||||
def write(self, frame_bgr: np.ndarray):
|
||||
frame_rgb = cv2.cvtColor(frame_bgr, cv2.COLOR_BGR2RGB)
|
||||
if not self.isOpened():
|
||||
raise RuntimeError("ffmpeg process is not running")
|
||||
|
||||
raise RuntimeError("FFMPEGVideoWriter: ffmpeg process is not running")
|
||||
|
||||
frame_rgb = cv2.cvtColor(frame_bgr, cv2.COLOR_BGR2RGB)
|
||||
try:
|
||||
self.proc.stdin.write(frame_rgb.astype(np.uint8).tobytes())
|
||||
# ensure data is flushed to the subprocess pipe
|
||||
self.proc.stdin.flush()
|
||||
self.proc.stdin.flush() # Critical: prevent buffering issues
|
||||
except BrokenPipeError:
|
||||
stderr = bytes(self._stderr_buffer).decode(errors='replace')
|
||||
raise RuntimeError(f"ffmpeg process closed (BrokenPipe). Stderr:\n{stderr}")
|
||||
raise RuntimeError(
|
||||
"FFMPEGVideoWriter: ffmpeg process terminated unexpectedly. "
|
||||
"Check video path, codec support, and disk space."
|
||||
)
|
||||
|
||||
def isOpened(self) -> bool:
|
||||
return self.proc is not None and self.proc.poll() is None
|
||||
@@ -200,43 +196,18 @@ class FFMPEGVideoWriter:
|
||||
try:
|
||||
self.proc.stdin.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Wait for process to exit and for stderr thread to finish
|
||||
pass # Ignore errors on close
|
||||
|
||||
self.proc.wait()
|
||||
try:
|
||||
self._stderr_thread.join(timeout=1.0)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Merge any remaining stderr
|
||||
remaining = b''
|
||||
try:
|
||||
if self.proc.stderr:
|
||||
remaining = self.proc.stderr.read() or b''
|
||||
except Exception:
|
||||
remaining = b''
|
||||
|
||||
stderr = bytes(self._stderr_buffer) + remaining
|
||||
|
||||
if self.proc.returncode != 0:
|
||||
debug.log(f"ffmpeg error: {stderr.decode(errors='replace')}", level="WARNING", category="file")
|
||||
|
||||
debug.log(
|
||||
f"ffmpeg exited with code {self.proc.returncode}. "
|
||||
"Check output file for corruption.",
|
||||
level="WARNING", force=True, category="file"
|
||||
)
|
||||
self.proc = None
|
||||
|
||||
def _read_stderr(self):
|
||||
# Continuously read stderr in small chunks to avoid filling the pipe buffer
|
||||
try:
|
||||
if not self.proc or not self.proc.stderr:
|
||||
return
|
||||
while True:
|
||||
chunk = self.proc.stderr.read(1024)
|
||||
if not chunk:
|
||||
break
|
||||
self._stderr_buffer.extend(chunk)
|
||||
except Exception:
|
||||
# Ignore read errors - best effort only
|
||||
return
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Device Management Helpers
|
||||
|
||||
Reference in New Issue
Block a user