import subprocess import sys import os import queue import time import threading class ProcessManager: _instance = None _lock = threading.Lock() @classmethod def instance(cls): with cls._lock: if cls._instance is None: cls._instance = cls() return cls._instance def __init__(self): self.active_processes = {} # pid -> (name, Popen object) def build_subprocess_env(self): env = os.environ.copy() env["PROTOCOL_BUFFERS_PYTHON_IMPLEMENTATION"] = "python" env["PYTHONUTF8"] = "1" base_dir = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..")) nvidia_root = os.path.join(base_dir, "env", "Lib", "site-packages", "nvidia") extra_paths = [] if os.path.isdir(nvidia_root): for name in os.listdir(nvidia_root): for sub in ("bin", "lib"): p = os.path.join(nvidia_root, name, sub) if os.path.isdir(p): extra_paths.append(p) if extra_paths: env["PATH"] = os.pathsep.join(extra_paths + [env.get("PATH", "")]) return env def register(self, name, process): with self._lock: self.active_processes[process.pid] = (name, process) def unregister(self, pid): with self._lock: if pid in self.active_processes: del self.active_processes[pid] def kill_process_tree(self, pid): """Cleanly kill a process tree on Windows using taskkill""" if sys.platform == 'win32': cmd = ["taskkill", "/F", "/T", "/PID", str(pid)] try: startupinfo = subprocess.STARTUPINFO() startupinfo.dwFlags |= subprocess.STARTF_USESHOWWINDOW subprocess.run(cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, startupinfo=startupinfo) except Exception as e: print(f"Failed to run taskkill for PID {pid}: {e}") else: # Unix fallback try: import signal os.killpg(os.getpgid(pid), signal.SIGKILL) except Exception: try: os.kill(pid, signal.SIGKILL) except Exception: pass def cancel_all(self): """Kill all active processes immediately""" with self._lock: pids = list(self.active_processes.keys()) for pid in pids: name, proc = self.active_processes.get(pid, ("Unknown", None)) print(f"Force-terminating active subprocess: {name} (PID: {pid})") if proc: try: self.kill_process_tree(pid) except Exception as e: print(f"Error killing PID {pid}: {e}") self.unregister(pid) def run_subprocess_sync(self, cmd, timeout=300, cwd=None, log_fn=None, startupinfo=None, env=None): """ Runs a subprocess and monitors its output. Automatically registers to active processes, respects timeouts, and cleans up on errors/cancellations. """ if startupinfo is None and sys.platform == 'win32': startupinfo = subprocess.STARTUPINFO() startupinfo.dwFlags |= subprocess.STARTF_USESHOWWINDOW # Extract name from cmd # FIX: trước đây lấy cmd[1] -> lệnh ffmpeg hiện tên "-y". Với [python, script] # thì tên đúng là script; còn lại là binary ở cmd[0]. if (len(cmd) > 1 and os.path.basename(str(cmd[0])).lower().startswith("python") and not str(cmd[1]).startswith("-")): name = os.path.basename(str(cmd[1])) else: name = os.path.basename(str(cmd[0])) try: process = subprocess.Popen( cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, encoding="utf-8", errors="ignore", cwd=cwd, startupinfo=startupinfo, env=env or self.build_subprocess_env() ) except Exception as e: raise Exception(f"Failed to start subprocess {name}: {e}") self.register(name, process) start_time = time.time() output_lines = [] try: # FIX: đọc stdout bằng reader thread + queue. Trước đây readline() block # trong lúc model load im lặng hàng phút -> watchdog timeout không bao # giờ được kiểm tra, job treo thay vì bị kill đúng hẹn. line_queue: queue.Queue = queue.Queue() def _reader(): try: for _line in process.stdout: line_queue.put(_line) except Exception: pass finally: line_queue.put(None) reader_thread = threading.Thread(target=_reader, daemon=True) reader_thread.start() # Non-blocking check loop with timeout stdout_closed = False while True: # Check timeout if time.time() - start_time > timeout: self.kill_process_tree(process.pid) raise TimeoutError(f"Subprocess {name} timed out after {timeout} seconds.") # Check if process finished retcode = process.poll() if retcode is not None and stdout_closed and line_queue.empty(): break # Read line (poll queue, never block the watchdog) try: line = line_queue.get(timeout=0.2) except queue.Empty: continue if line is None: stdout_closed = True continue line_str = line.strip() if not line_str: continue if log_fn: log_fn(line_str) output_lines.append(line_str) if process.returncode != 0: last_logs = "\n".join(output_lines[-25:]) if output_lines else "(Không có log output)" raise Exception( f"Subprocess '{name}' thất bại với exit code {process.returncode}.\n" f"--- CHI TIẾT LOG SUBPROCESS ---\n{last_logs}" ) return "\n".join(output_lines) finally: self.unregister(process.pid)