from __future__ import annotations import os import sys import subprocess import pathlib import shutil import re import uuid import json import glob import random import time import base64 import asyncio from functools import lru_cache from typing import Any import requests as http_requests # ============================================================ # 1. HUGGINGFACE_HUB SELF-HEALING REPAIR # ============================================================ def _hf_hub_version() -> str: try: from importlib.metadata import version as _pkg_version return _pkg_version("huggingface_hub") except Exception: return "" def _hf_hub_is_broken() -> bool: import importlib import importlib.util for module_name in ("huggingface_hub._snapshot_download", "huggingface_hub._tree_cache"): try: if importlib.util.find_spec(module_name) is None: continue except Exception: return True try: importlib.import_module(module_name) except ImportError: return True except Exception: continue return False def _hf_hub_reinstall(upgrade: bool) -> None: cmd = [ sys.executable, "-m", "pip", "install", "--no-cache-dir", "--force-reinstall", "--no-deps", ] if upgrade: cmd += ["--upgrade", "huggingface_hub"] else: pinned = _hf_hub_version() cmd.append(f"huggingface_hub=={pinned}" if pinned else "huggingface_hub") print(f"[hf-repair] {' '.join(cmd)}", flush=True) subprocess.run(cmd, check=False) def _repair_huggingface_hub_and_restart() -> None: stage = int(os.environ.get("_HF_HUB_REPAIR_STAGE", "0") or "0") if stage >= 2 or not _hf_hub_is_broken(): return _hf_hub_reinstall(upgrade=stage == 1) os.environ["_HF_HUB_REPAIR_STAGE"] = str(stage + 1) os.execv(sys.executable, [sys.executable, *sys.argv]) _repair_huggingface_hub_and_restart() # ============================================================ # 2. IMPORTS UTAMA (SPACES WAJIB PERTAMA SEBELUM TORCH) # ============================================================ import spaces # WAJIB PERTAMA sebelum torch! import torch # ============================================================ # ZEROGPU COMPATIBILITY PATCH FOR PYTORCH CUDA MOCK PROPERTIES # ============================================================ if hasattr(torch, "cuda") and hasattr(torch.cuda, "get_device_properties"): _orig_cuda_get_device_properties = torch.cuda.get_device_properties def _safe_cuda_get_device_properties(device=None): props = _orig_cuda_get_device_properties(device) if not hasattr(props, "is_integrated"): try: setattr(props, "is_integrated", False) except Exception: class _PropsProxy: def __init__(self, p): self._p = p self.is_integrated = False def __getattr__(self, name): return getattr(self._p, name) return _PropsProxy(props) return props torch.cuda.get_device_properties = _safe_cuda_get_device_properties from fastapi import Request, Response, HTTPException import gradio as gr import gradio_client.utils from gradio_client import Client, handle_file from huggingface_hub import hf_hub_download # ============================================================ # 2.1 MONKEY-PATCH GRADIO_CLIENT OPENAPI SCHEMA BUG # ============================================================ _orig_get_type = gradio_client.utils.get_type def _safe_get_type(schema): if isinstance(schema, bool): return "boolean" if not isinstance(schema, dict): return "str" return _orig_get_type(schema) gradio_client.utils.get_type = _safe_get_type _orig_json_schema = gradio_client.utils._json_schema_to_python_type def _safe_json_schema(schema, defs=None): if isinstance(schema, bool): return "bool" if not isinstance(schema, dict): return "str" return _orig_json_schema(schema, defs) gradio_client.utils._json_schema_to_python_type = _safe_json_schema # ============================================================ # 3. KONFIGURASI PATH & DIREKTORI # ============================================================ ROOT = pathlib.Path(__file__).resolve().parent COMFY = ROOT / "ComfyUI" MODELS = COMFY / "models" INPUT = COMFY / "input" OUTPUT = COMFY / "output" LOCAL_CUSTOM_NODES = ROOT / "custom_nodes" WORKFLOW_FILE = ROOT / "workflow_generator.json" NODE_OUTPUT_ID = "92" CONDITIONER_SPACE = os.environ.get("H3_CONDITIONER_SPACE", "alibaybay/dualspace-h3-clip") # ============================================================ # 4. DAFTAR MODEL GENERATOR MINIMAX-H3 (INT8 + TAOMATE 3-STEP) # ============================================================ DOWNLOADS = [ { "repo": "Comfy-Org/MiniMax-H3", "file": "diffusion_models/minimax_h3_fl2va_pruned_int8_convrot.safetensors", "dest": MODELS / "diffusion_models" / "minimax_h3_fl2va_pruned_int8_convrot.safetensors", "alt_dest": MODELS / "unet" / "minimax_h3_fl2va_pruned_int8_convrot.safetensors", "label": "Diffusion Model (FL2VA Pruned INT8 ~21.0GB)", }, { "repo": "Kijai/MiniMax-H3_comfy", "file": "loras/minimax_h3_taomate_3step_lora_avg_rank_19_bf16.safetensors", "dest": MODELS / "loras" / "minimax_h3_taomate_3step_lora_avg_rank_19_bf16.safetensors", "alt_dest": None, "label": "TaoMate 3-Step LoRA BF16 (~1.96GB)", }, { "repo": "Comfy-Org/MiniMax-H3", "file": "vae/minimax_h3_video_vae_int8_convrot.safetensors", "dest": MODELS / "vae" / "minimax_h3_video_vae_int8_convrot.safetensors", "alt_dest": None, "label": "Video VAE INT8 ConvRot (~2.6GB)", }, { "repo": "Comfy-Org/MiniMax-H3", "file": "vae/minimax_h3_audio_vae_fp32.safetensors", "dest": MODELS / "vae" / "minimax_h3_audio_vae_fp32.safetensors", "alt_dest": None, "label": "Audio VAE FP32 (~605MB)", }, ] CUSTOM_NODES: list[tuple[str, str]] = [] _comfy_ready = False _nodes_ready = False server_instance = None gpu_lock = asyncio.Lock() # Job Tracking Asynchronous REST API JOBS: dict[str, dict[str, Any]] = {} # ============================================================ # 5. HELPER & MODEL DOWNLOADER # ============================================================ def _run_cmd(cmd: list[str], cwd: pathlib.Path = ROOT, check: bool = True) -> None: print(f"[*] Menjalankan: {' '.join(cmd)} di {cwd}", flush=True) subprocess.run(cmd, cwd=cwd, check=check) def _link_or_copy(src: pathlib.Path, dest: pathlib.Path) -> None: dest.parent.mkdir(parents=True, exist_ok=True) if dest.is_symlink(): dest.unlink() if dest.exists() and dest.stat().st_size > 1000: return try: os.link(src, dest) return except OSError: pass shutil.copy2(src, dest) def _download_to_dest(repo: str, file_path: str, dest: pathlib.Path, token: str | None) -> None: dest.parent.mkdir(parents=True, exist_ok=True) if dest.is_symlink(): dest.unlink() if dest.exists() and dest.stat().st_size > 1000: return p = pathlib.Path(file_path) filename = p.name subfolder = str(p.parent) if str(p.parent) != "." else None print(f"[*] Mengunduh {filename} dari {repo} ke {dest.parent}...", flush=True) downloaded_str = hf_hub_download( repo_id=repo, filename=filename, subfolder=subfolder, local_dir=str(dest.parent), token=token, ) downloaded = pathlib.Path(downloaded_str) if downloaded.resolve() == dest.resolve(): return if dest.exists() or dest.is_symlink(): dest.unlink() dest.parent.mkdir(parents=True, exist_ok=True) try: os.replace(downloaded, dest) except OSError: shutil.copy2(downloaded, dest) if downloaded.exists(): downloaded.unlink() sub_dir = dest.parent / "split_files" if sub_dir.exists(): shutil.rmtree(sub_dir, ignore_errors=True) def _install_filtered_requirements(req_path: pathlib.Path, cwd: pathlib.Path) -> None: if not req_path.exists(): return blocked = {"torch", "torchvision", "torchaudio", "transformers", "huggingface-hub", "accelerate", "xformers"} safe: list[str] = [] for line in req_path.read_text(encoding="utf-8", errors="ignore").splitlines(): item = line.strip() if not item or item.startswith("#"): continue low = item.lower().replace("_", "-") package = re.split(r"[<>=!~;\[\s]", low, maxsplit=1)[0] if package in blocked: continue safe.append(item) if safe: filtered_file = cwd / "requirements_filtered.txt" filtered_file.write_text("\n".join(safe) + "\n", encoding="utf-8") _run_cmd([sys.executable, "-m", "pip", "install", "-r", "requirements_filtered.txt", "--no-cache-dir"], cwd=cwd, check=False) def _apply_comfy_utils_namespace_fix() -> None: utils_path = COMFY / "utils" utilities_path = COMFY / "utilities" if utils_path.exists() and not utilities_path.exists(): try: utils_path.rename(utilities_path) except OSError: pass replacements = [ (re.compile(r"(^|\n)(\s*)from utils(\s|\.)"), r"\1\2from utilities\3"), (re.compile(r"(^|\n)(\s*)import utils(\s|\.|$)"), r"\1\2import utilities\3"), ] for path in COMFY.rglob("*.py"): if "__pycache__" in path.parts: continue try: text = path.read_text(encoding="utf-8") except UnicodeDecodeError: continue updated = text for pattern, repl in replacements: updated = pattern.sub(repl, updated) updated = updated.replace("from utils import", "from utilities import") if updated != text: path.write_text(updated, encoding="utf-8") def _ensure_comfy() -> None: global _comfy_ready if _comfy_ready: return print("[1/3] Menyiapkan ComfyUI Runtime untuk H3 Generator...", flush=True) if not COMFY.exists(): _run_cmd(["git", "clone", "--depth", "1", "--branch", "v0.38.2", "https://github.com/comfyanonymous/ComfyUI.git", str(COMFY)]) _install_filtered_requirements(COMFY / "requirements.txt", COMFY) custom_root = COMFY / "custom_nodes" custom_root.mkdir(parents=True, exist_ok=True) # Pasang local thin wire node: external_h3_conditioning if LOCAL_CUSTOM_NODES.exists(): for src_node in LOCAL_CUSTOM_NODES.iterdir(): if src_node.is_dir() and not src_node.name.startswith("."): target_node = custom_root / src_node.name if target_node.exists(): shutil.rmtree(target_node, ignore_errors=True) shutil.copytree(src_node, target_node) print(f"[*] Terpasang local custom node: {src_node.name}", flush=True) _apply_comfy_utils_namespace_fix() for folder in ("diffusion_models", "unet", "loras", "vae"): (MODELS / folder).mkdir(parents=True, exist_ok=True) INPUT.mkdir(parents=True, exist_ok=True) OUTPUT.mkdir(parents=True, exist_ok=True) _comfy_ready = True print("[1/3] ComfyUI Runtime H3 Generator Siap.", flush=True) def _ensure_models(progress=None) -> None: print("[2/3] Memeriksa & Mengunduh Model Generator MiniMax-H3...", flush=True) token = os.environ.get("HF_TOKEN") or os.environ.get("HUGGINGFACE_HUB_TOKEN") for row in DOWNLOADS: dest = pathlib.Path(row["dest"]) dest.parent.mkdir(parents=True, exist_ok=True) if dest.is_symlink(): dest.unlink() if not (dest.exists() and dest.stat().st_size > 1000): print(f"[*] Mengunduh {row['label']}...", flush=True) _download_to_dest(row["repo"], row["file"], dest, token) alt = row.get("alt_dest") if alt is not None: alt_path = pathlib.Path(alt) if dest.exists() and dest.stat().st_size > 1000: _link_or_copy(dest, alt_path) print("[2/3] Semua model generator MiniMax-H3 telah siap.", flush=True) def _init_comfy_nodes() -> None: global _nodes_ready, server_instance if _nodes_ready: return print("[3/3] Menginisialisasi Engine ComfyUI Generator (Full Standby)...", flush=True) comfy_path = str(COMFY) sys.path = [p for p in sys.path if p != comfy_path] sys.path.insert(0, comfy_path) for module_name in list(sys.modules): if module_name in ("utils", "app") or module_name.startswith(("utils.", "app.")): del sys.modules[module_name] os.chdir(COMFY) import execution import nodes import server loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) import inspect sig = inspect.signature(server.PromptServer.__init__) if "asset_manager" in sig.parameters: try: from app.assets.manager import default_asset_manager asset_mgr = default_asset_manager() except Exception: class DummyAssetManager: enabled = False def startup(self): pass def shutdown(self): pass def register_routes(self, app, user_manager=None): pass def ensure_scan_started(self): pass def pause_background_scan(self): pass def queue_output_scan(self): pass def resume_background_scan(self): pass def register_upload(self, *args, **kwargs): return None def register_executed_output(self, *args, **kwargs): return None def register_cached_output(self, *args, **kwargs): return None def set_event_sink(self, sink): pass asset_mgr = DummyAssetManager() server_instance = server.PromptServer(loop, asset_mgr) else: server_instance = server.PromptServer(loop) try: execution.PromptQueue(server_instance) except Exception: pass res = nodes.init_extra_nodes() if asyncio.iscoroutine(res): loop.run_until_complete(res) _nodes_ready = True print("[3/3] Engine ComfyUI Generator Siap & Berada dalam Mode Hot-Standby.", flush=True) executor_instance = None def _get_or_create_executor(): global executor_instance if executor_instance is None: import execution executor_instance = execution.PromptExecutor( server_instance, cache_type=execution.CacheType.RAM_PRESSURE, cache_args={"lru": 32, "ram": 60.0, "ram_inactive": 60.0}, ) return executor_instance def _preload_models_to_ram(): """Me-load UNet INT8, LoRA TaoMate 3-Step, Video VAE INT8, dan Audio VAE ke RAM saat boot.""" print("[*] Pre-loading bobot model (UNet INT8 ~21GB, LoRA, VAE INT8) ke RAM...", flush=True) try: import nodes unet_loader = nodes.UNETLoader() vae_loader = nodes.VAELoader() lora_loader = nodes.LoraLoaderModelOnly() # 1. Load UNet INT8 print("[*] Pre-loading UNet INT8...", flush=True) unet_res = unet_loader.load_unet("minimax_h3_fl2va_pruned_int8_convrot.safetensors", weight_dtype="default") unet_model = unet_res[0] # 2. Patch LoRA TaoMate 3-Step print("[*] Pre-patching LoRA TaoMate 3-Step...", flush=True) lora_loader.load_lora_model_only(unet_model, "minimax_h3_taomate_3step_lora_avg_rank_19_bf16.safetensors", 1.0) # 3. Load Video VAE INT8 & Audio VAE print("[*] Pre-loading Video VAE INT8 & Audio VAE...", flush=True) vae_loader.load_vae("minimax_h3_video_vae_int8_convrot.safetensors") vae_loader.load_vae("minimax_h3_audio_vae_fp32.safetensors") print("[*] Pre-load model ke RAM berhasil! Semua model telah siap di memory (Zero Disk Reload).", flush=True) except Exception as e: print(f"[!] Warning saat pre-load model: {e}", flush=True) # ============================================================ # 6. ROOT STARTUP PRE-WARMING (HOT-STANDBY OPTIMIZATION) # ============================================================ def _startup_prewarm(): print("=" * 60, flush=True) print("[startup] Memulai Pre-Warming Engine H3 Generator Hub...", flush=True) _ensure_comfy() _ensure_models() _init_comfy_nodes() _get_or_create_executor() _preload_models_to_ram() print("[startup] Pre-Warming Selesai. Siap Melayani Permintaan Instan.", flush=True) print("=" * 60, flush=True) _startup_prewarm() # ============================================================ # 7. PERSISTENT LRU CLIENT POOL KE SPACE KONDISIONER # ============================================================ @lru_cache(maxsize=32) def _get_conditioner_client(space_id: str, ip_token: str | None) -> Client: """Membuka dan mempertahankan sesi koneksi HTTP / WebSocket ke Space 1.""" headers = {"x-ip-token": ip_token} if ip_token else {} print(f"[*] [LRU Pool] Membuka persistent connection ke Conditioner: {space_id}", flush=True) return Client(space_id, headers=headers) def fetch_remote_conditioning( prompt: str, first_frame_path: str, last_frame_path: str | None, duration: float, ip_token: str | None, ) -> str: """Mengirim parameter ke Space 1 via Thin Wire API dan menerima file .safetensors.""" client = _get_conditioner_client(CONDITIONER_SPACE, ip_token) print(f"[*] [Space 2] Meminta conditioning dari Space 1 ({CONDITIONER_SPACE})...", flush=True) res = client.predict( prompt=prompt, first_frame_path=handle_file(first_frame_path), last_frame_path=handle_file(last_frame_path) if last_frame_path else None, duration=f"{duration:.0f}s", api_name="/encode", ) if isinstance(res, (tuple, list)): remote_file = res[0] elif isinstance(res, dict) and "path" in res: remote_file = res["path"] else: remote_file = str(res) size_kb = os.path.getsize(remote_file) / 1024.0 if os.path.exists(remote_file) else 0 print(f"[*] [Space 2] Berhasil menerima Thin Wire safetensors: {remote_file} ({size_kb:.2f} KB)", flush=True) return remote_file # ============================================================ # 8. LOGIKA RUNNER ZERO-OVERHEAD @SPACES.GPU (DYNAMIC DURATION) # ============================================================ DURATION_GPU_MAP = { 5.0: 55, # Waktu riil ~35s -> Minta 55s (ZeroGPU reserve: 82.5s) 10.0: 75, # Waktu riil ~55s -> Minta 75s (ZeroGPU reserve: 112.5s) 15.0: 100, # Waktu riil ~74.6s -> Minta 100s (ZeroGPU reserve: 150.0s) } def get_generator_duration(safetensors_path: str, video_duration: float | str = 5.0, seed: int = 0) -> int: try: dur = float(str(video_duration).replace("s", "").strip()) except Exception: dur = 5.0 return DURATION_GPU_MAP.get(dur, int(min(max(50, 30 + dur * 4.5), 110))) def _load_base_workflow() -> dict[str, Any]: with open(WORKFLOW_FILE, "r", encoding="utf-8") as f: return json.load(f) def _run_comfy_generator_workflow(workflow: dict[str, Any]) -> str: """Eksekusi workflow Denoising di Space 2.""" import execution executor = _get_or_create_executor() prompt_id = str(uuid.uuid4()) node_durations: list[tuple[str, str, float]] = [] orig_get_output_data = execution.get_output_data def _profiling_get_output_data(obj, input_data_all, *args, **kwargs): if isinstance(obj, str): node_id = obj node_info = workflow.get(node_id, {}) class_type = node_info.get("class_type", "UnknownNode") node_title = node_info.get("_meta", {}).get("title", class_type) label = f"[Node {node_id}: {node_title}]" else: node_id = "?" class_type = obj.__class__.__name__ label = f"[{class_type}]" t0 = time.time() print(f"🚀 {label} Mulai dieksekusi...", flush=True) try: res = orig_get_output_data(obj, input_data_all, *args, **kwargs) dur = time.time() - t0 node_durations.append((node_id, label, dur)) print(f"⏱️ {label} Selesai dalam: {dur:.2f}s", flush=True) return res except Exception as e: dur = time.time() - t0 print(f"❌ {label} Gagal setelah: {dur:.2f}s ({e})", flush=True) raise execution.get_output_data = _profiling_get_output_data t_workflow_start = time.time() try: executor.execute( workflow, prompt_id, extra_data={}, execute_outputs=[NODE_OUTPUT_ID], ) finally: execution.get_output_data = orig_get_output_data t_workflow_total = time.time() - t_workflow_start print("\n" + "=" * 70, flush=True) print("📊 REKAPITULASI PROFILING WAKTU SPACE 2 (PER NODE):", flush=True) print("=" * 70, flush=True) sorted_nodes = sorted(node_durations, key=lambda x: x[2], reverse=True) for nid, label, dur in sorted_nodes: pct = (dur / t_workflow_total * 100) if t_workflow_total > 0 else 0 bar = "█" * int(pct // 5) print(f" {label:<45} : {dur:>6.2f}s ({pct:>5.1f}%) {bar}", flush=True) print("-" * 70, flush=True) print(f" ⏱️ TOTAL DURASI SAMPLING & DECODE : {t_workflow_total:.2f} detik", flush=True) print("=" * 70 + "\n", flush=True) if not executor.success: err = ( executor.status_messages[-1] if hasattr(executor, "status_messages") and executor.status_messages else "ComfyUI execution gagal" ) raise RuntimeError(str(err)) files = [ pathlib.Path(p) for p in glob.glob(str(OUTPUT / "**" / "*.mp4"), recursive=True) ] if not files: files = [ pathlib.Path(p) for p in glob.glob(str(OUTPUT / "**" / "*.*"), recursive=True) if p.endswith((".mp4", ".webm", ".mkv", ".mov")) ] if not files: raise RuntimeError("Generation selesai tetapi file video output (.mp4) tidak ditemukan di folder output.") latest_video = sorted(files, key=lambda p: p.stat().st_mtime, reverse=True)[0] return str(latest_video) @spaces.GPU(duration=get_generator_duration) def run_generator_gpu( safetensors_path: str, video_duration: float | str = 5.0, seed: int = 0, ) -> str: """Eksekusi murni GPU forward: Denoising TaoMate 3-Step LoRA + Video & Audio VAE Decode.""" wf = _load_base_workflow() try: dur_val = float(str(video_duration).replace("s", "").strip()) except Exception: dur_val = 5.0 prefix = f"H3_vid_{uuid.uuid4().hex[:8]}" wf["ext_h3_cond"]["inputs"]["safetensors_path"] = safetensors_path wf["105_15"]["inputs"]["noise_seed"] = int(seed) wf["105_9"]["inputs"]["steps"] = 3 # Kunci standar TaoMate 3-Step LoRA wf["92"]["inputs"]["filename_prefix"] = f"video/{prefix}" print(f"[*] [GPU Space 2] Menjalankan MiniMax-H3 Denoising (TaoMate 3-Step, Durasi={dur_val}s, Seed={seed})...", flush=True) target_video = _run_comfy_generator_workflow(wf) print(f"[*] [GPU Selesai] Video MP4 berhasil dibuat: {target_video}", flush=True) return target_video # ============================================================ # 9. PIPELINE END-TO-END DUAL-SPACE # ============================================================ def generate_video_pipeline( first_frame: str | None, last_frame: str | None, prompt: str, duration_choice: str, seed: float, randomize_seed: bool, request: gr.Request = None, progress=gr.Progress(track_tqdm=True), ) -> tuple[str | None, str, int]: if not first_frame or not os.path.exists(first_frame): raise gr.Error("First Frame (gambar keyframe awal) wajib diunggah untuk mode Image-to-Video!") seed_int = int(seed) if randomize_seed or seed_int == 0: seed_int = random.randint(1, 1000000000000000) try: dur_val = float(str(duration_choice).replace("s", "").strip()) except Exception: dur_val = 5.0 ip_token = None if request and hasattr(request, "headers"): ip_token = request.headers.get("x-ip-token") t_total_start = time.time() # Step 1: Conditioning via Space 1 (Thin Wire) progress(0.1, desc="⚡ [Space 1] Menghubungi Qwen3-VL Conditioner & Keyframes Encoder...") t_cond_start = time.time() try: safetensors_file = fetch_remote_conditioning( prompt=prompt or "", first_frame_path=first_frame, last_frame_path=last_frame, duration=dur_val, ip_token=ip_token, ) except Exception as e: raise gr.Error(f"Gagal memproses conditioning di Space 1 ({CONDITIONER_SPACE}): {e}") t_cond_elapsed = time.time() - t_cond_start # Step 2: Denoising & Video Decode via Space 2 GPU progress(0.4, desc=f"🎬 [Space 2] Denoising TaoMate 3-Step & Decoding Video ({dur_val:.0f}s)...") t_gen_start = time.time() try: video_path = run_generator_gpu( safetensors_path=safetensors_file, video_duration=dur_val, seed=seed_int, ) except Exception as e: raise gr.Error(f"Gagal saat proses denoising/video generation: {e}") t_gen_elapsed = time.time() - t_gen_start t_total_elapsed = time.time() - t_total_start report = ( f"✅ Video Berhasil Dibuat!\n" f"⏱️ Total Waktu: {t_total_elapsed:.2f}s | " f"🧠 Space 1 (Conditioner): {t_cond_elapsed:.2f}s | " f"🎬 Space 2 (Denoise & Decode): {t_gen_elapsed:.2f}s\n" f"⚙️ Resolusi: Otomatis (0.4 MP) | Durasi: {dur_val:.0f}s | Steps: 3 (TaoMate LoRA) | Seed: {seed_int}" ) return video_path, report, seed_int # ============================================================ # 10. ASYNCHRONOUS REST TWIN-API (STANDAR INDUSTRI REPLICATE/FAL) # ============================================================ def _cleanup_expired_jobs(ttl_seconds: int = 900) -> None: now = time.time() expired = [] for jid, info in list(JOBS.items()): # 1. Grace Period Burn (60 detik setelah unduhan pertama) burn_at = info.get("burn_at") if burn_at and now > burn_at: expired.append(jid) continue # 2. TTL Max Lifetime (15 menit jika tidak pernah diunduh) if now - info.get("created_at", 0) > ttl_seconds: expired.append(jid) for jid in expired: path = JOBS[jid].get("output_path") if path and os.path.exists(path): try: os.remove(path) print(f"[*] [Burned] File {path} telah dihapus permanen dari disk.", flush=True) except Exception: pass JOBS.pop(jid, None) def process_media_value(val: str, prefix: str = "media") -> str: """Download image dari URL atau decode dari base64, return absolute filepath.""" if not isinstance(val, str) or not val.strip(): return "" val = val.strip() if os.path.exists(val): return val if val.startswith("http://") or val.startswith("https://"): ext = "png" target_path = INPUT / f"{prefix}_{uuid.uuid4().hex[:8]}.{ext}" r = http_requests.get(val, stream=True, timeout=60) r.raise_for_status() with open(target_path, "wb") as f: for chunk in r.iter_content(chunk_size=8192): f.write(chunk) return str(target_path) if val.startswith("data:"): header, encoded = val.split(",", 1) ext = "png" if "jpeg" in header or "jpg" in header: ext = "jpg" elif "webp" in header: ext = "webp" target_path = INPUT / f"{prefix}_{uuid.uuid4().hex[:8]}.{ext}" target_path.write_bytes(base64.b64decode(encoded)) return str(target_path) return val async def _job_worker(job_id: str, payload: dict[str, Any]) -> None: JOBS[job_id]["status"] = "PROCESSING" JOBS[job_id]["start_time"] = time.time() try: raw_first = payload.get("image") or payload.get("first_frame") raw_last = payload.get("last_image") or payload.get("last_frame") prompt = payload.get("prompt", "") raw_dur = payload.get("duration", 5.0) raw_seed = payload.get("seed", 0) if not raw_first: raise ValueError("Field 'image' (first frame) wajib disertakan.") first_frame = await asyncio.to_thread(process_media_value, raw_first, "first_frame") last_frame = await asyncio.to_thread(process_media_value, raw_last, "last_frame") if raw_last else None try: dur_val = float(str(raw_dur).replace("s", "").strip()) except Exception: dur_val = 5.0 seed_val = int(raw_seed) if seed_val <= 0: seed_val = random.randint(1, 1000000000000000) # 1. Hubungi Space 1 untuk conditioning safetensors_file = await asyncio.to_thread( fetch_remote_conditioning, prompt=prompt, first_frame_path=first_frame, last_frame_path=last_frame, duration=dur_val, ip_token=None, ) # 2. Eksekusi forward GPU Space 2 async with gpu_lock: video_path = await asyncio.to_thread( run_generator_gpu, safetensors_path=safetensors_file, video_duration=dur_val, seed=seed_val, ) if not video_path or not os.path.exists(video_path): raise RuntimeError("Eksekusi selesai tetapi video output tidak ditemukan.") JOBS[job_id]["status"] = "COMPLETED" JOBS[job_id]["end_time"] = time.time() JOBS[job_id]["output_path"] = str(video_path) JOBS[job_id]["media_type"] = "video/mp4" JOBS[job_id]["filename"] = pathlib.Path(video_path).name except Exception as err: JOBS[job_id]["status"] = "FAILED" JOBS[job_id]["error"] = str(err) JOBS[job_id]["end_time"] = time.time() async def api_submit_job(request: Request): try: item = await request.json() except Exception: raise HTTPException(status_code=400, detail="Payload harus berupa JSON yang valid.") _cleanup_expired_jobs() job_id = f"job_{uuid.uuid4().hex}" # Full 128-bit Cryptographic UUIDv4 (32 hex characters) now = time.time() JOBS[job_id] = { "job_id": job_id, "status": "QUEUED", "created_at": now, "start_time": None, "end_time": None, "output_path": None, "media_type": None, "burn_at": None, "error": None, } asyncio.create_task(_job_worker(job_id, item)) return Response( content=json.dumps({ "status": "QUEUED", "job_id": job_id, "created_at": round(now, 3), "check_url": f"/api/status/{job_id}", }), status_code=202, media_type="application/json", ) async def api_get_status(job_id: str): _cleanup_expired_jobs() job = JOBS.get(job_id) if not job: raise HTTPException(status_code=404, detail="Job ID tidak ditemukan, telah diunduh dan dibakar (Burned After Read), atau telah kedaluwarsa.") status = job["status"] now = time.time() if status == "QUEUED": return { "job_id": job_id, "status": "QUEUED", "elapsed_seconds": round(now - job["created_at"], 2), } elif status == "PROCESSING": start = job["start_time"] or job["created_at"] return { "job_id": job_id, "status": "PROCESSING", "elapsed_seconds": round(now - start, 2), } elif status == "COMPLETED": exec_time = round((job["end_time"] or now) - (job["start_time"] or job["created_at"]), 2) burn_msg = f"Hangus dalam {max(0, round(job['burn_at'] - now, 1))} detik" if job.get("burn_at") else "Hangus 60 detik setelah unduhan pertama" return { "job_id": job_id, "status": "COMPLETED", "execution_time_seconds": exec_time, "media_type": job["media_type"], "result_url": f"/api/result/{job_id}", "privacy_policy": "Burn After Read (60s Grace Period)", "burn_status": burn_msg, "error": None, } else: # FAILED return { "job_id": job_id, "status": "FAILED", "error": job.get("error"), } async def api_get_result(job_id: str): _cleanup_expired_jobs() job = JOBS.get(job_id) if not job: raise HTTPException( status_code=404, detail="Job ID tidak ditemukan, telah diunduh dan dibakar (Burned After Read), atau telah kedaluwarsa.", ) if job["status"] != "COMPLETED": raise HTTPException(status_code=400, detail=f"Job belum selesai. Status saat ini: {job['status']}") now = time.time() # Jika ini unduhan pertama kali, aktifkan grace period countdown (60 detik) if not job.get("burn_at"): job["burn_at"] = now + 60.0 print(f"[*] [Burn-After-Read] Job {job_id} pertama kali diunduh. Dijadwalkan hangus permanen dalam 60 detik.", flush=True) output_path = job.get("output_path") if not output_path or not os.path.exists(output_path): raise HTTPException( status_code=404, detail="File output telah hangus dari server (Burned After Read).", ) media_bytes = pathlib.Path(output_path).read_bytes() exec_time = round((job["end_time"] or now) - (job["start_time"] or job["created_at"]), 2) remaining_burn = max(0, round(job["burn_at"] - now, 1)) headers = { "Content-Disposition": f'inline; filename="{job.get("filename", "output.mp4")}"', "X-Status": "success", "X-Execution-Time": f"{exec_time:.2f}s", "X-Burn-In-Seconds": str(remaining_burn), } return Response(content=media_bytes, media_type="video/mp4", headers=headers) # ============================================================ # 11. ANTARMUKA GRADIO STUDIO MINIMALIS & ELEGAN # ============================================================ custom_css = """ .container { max-width: 1200px; margin: auto; } .generate-btn { font-size: 1.15rem !important; padding: 12px !important; font-weight: bold !important; } """ with gr.Blocks(title="MiniMax-H3 Studio — Dual-Space ComfyUI", css=custom_css, theme=gr.themes.Soft()) as demo: gr.Markdown( """ # ⚡ MiniMax-H3 FL2VA — Dual-Space Video Studio & REST API Aplikasi Video Generative canggih multimodal berbasis **MiniMax-H3 (FL2VA)** dengan arsitektur **Dual-Space ZeroGPU**. - **Space 1 (`Conditioner`)**: Memproses Qwen3-VL 32B + Video VAE Keyframes via Thin Wire Protocol. - **Space 2 (`Generator Hub`)**: Melakukan Denoising INT8 FL2VA + TaoMate 3-Step LoRA + Decode Video & Audio. - **REST API Aktif**: Mendukung endpoint asinkron `/api/generate`, `/api/status/{job_id}`, dan `/api/result/{job_id}`. """ ) with gr.Row(): with gr.Column(scale=5): with gr.Row(): first_frame_input = gr.Image(type="filepath", label="First Frame (Keyframe Awal - Wajib)") last_frame_input = gr.Image(type="filepath", label="Last Frame (Keyframe Akhir - Opsional)") prompt_input = gr.Textbox( label="Prompt Teks (Opsional)", placeholder="Deskripsikan aksi atau suasana video sinematik...", lines=2, value="cinematic camera movement, natural motion, hyper-detailed", ) with gr.Group(): duration_input = gr.Radio( choices=["5s", "10s", "15s"], value="5s", label="Durasi Video", ) with gr.Row(): seed_input = gr.Number(value=0, label="Seed (0 = Random)", precision=0) randomize_seed_input = gr.Checkbox(label="🎲 Randomize Seed Setiap Generate", value=True) btn_generate = gr.Button("🚀 Generate Video (Dual-Space I2V)", variant="primary", elem_classes=["generate-btn"]) with gr.Column(scale=5): video_output = gr.Video(label="Hasil Video MiniMax-H3 (dengan Audio)", autoplay=True) report_output = gr.Markdown(label="Status & Benchmark Log") btn_generate.click( fn=generate_video_pipeline, inputs=[ first_frame_input, last_frame_input, prompt_input, duration_input, seed_input, randomize_seed_input, ], outputs=[video_output, report_output, seed_input], ) if __name__ == "__main__": # 1. Jalankan Gradio dengan prevent_thread_lock=True fastapi_app, _, _ = demo.queue(max_size=20).launch(show_error=True, prevent_thread_lock=True, ssr_mode=False) # 2. Daftarkan 100% Asynchronous REST API fastapi_app.post("/api/generate")(api_submit_job) fastapi_app.post("/")(api_submit_job) fastapi_app.get("/api/status/{job_id}")(api_get_status) fastapi_app.get("/api/result/{job_id}")(api_get_result) # 3. Kunci main thread agar melayani request demo.block_thread()