diff --git a/__init__.py b/__init__.py index 9d5288a..ffce45f 100644 --- a/__init__.py +++ b/__init__.py @@ -20,16 +20,14 @@ from .model_management_mgpu import ( force_full_system_cleanup, ) - MGPU_MM_LOG = True +DEBUG_LOG = False -# Set to "E" for Engineering (DEBUG) or "P" for Production (INFO) -LOG_LEVEL = "P" logger = logging.getLogger("MultiGPU") logger.propagate = False if not logger.handlers: - log_level = logging.DEBUG if LOG_LEVEL == "E" else logging.INFO + log_level = logging.DEBUG if DEBUG_LOG else logging.INFO handler = logging.StreamHandler() formatter = logging.Formatter('%(message)s') handler.setFormatter(formatter) @@ -41,7 +39,15 @@ def mgpu_mm_log_method(self, msg): self.info(f"[MultiGPU Model Management] {msg}") logger.mgpu_mm_log = mgpu_mm_log_method.__get__(logger, type(logger)) -# Global device state management +def check_module_exists(module_path): + full_path = os.path.join(folder_paths.get_folder_paths("custom_nodes")[0], module_path) + logger.debug(f"[MultiGPU] Checking for module at {full_path}") + if not os.path.exists(full_path): + logger.debug(f"[MultiGPU] Module {module_path} not found - skipping") + return False + logger.debug(f"[MultiGPU] Found {module_path}, creating compatible MultiGPU nodes") + return True + current_device = mm.get_torch_device() current_text_encoder_device = mm.text_encoder_device() @@ -55,81 +61,6 @@ def set_current_text_encoder_device(device): current_text_encoder_device = device logger.debug(f"[MultiGPU Initialization] current_text_encoder_device set to: {device}") -def override_class(cls): - class NodeOverride(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) - return inputs - - CATEGORY = "multigpu" - FUNCTION = "override" - - def override(self, *args, device=None, **kwargs): - - if device is not None: - set_current_device(device) - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - return out - - return NodeOverride - -def override_class_clip(cls): - class NodeOverride(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) - return inputs - - CATEGORY = "multigpu" - FUNCTION = "override" - - def override(self, *args, device=None, **kwargs): - if device is not None: - set_current_text_encoder_device(device) - kwargs['device'] = 'default' - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - return out - - return NodeOverride - -def override_class_clip_no_device(cls): - class NodeOverride(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) - return inputs - - CATEGORY = "multigpu" - FUNCTION = "override" - - def override(self, *args, device=None, **kwargs): - if device is not None: - set_current_text_encoder_device(device) - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - return out - - return NodeOverride - - def get_torch_device_patched(): device = None if (not is_accelerator_available() or mm.cpu_state == mm.CPUState.CPU or "cpu" in str(current_device).lower()): @@ -150,23 +81,12 @@ def text_encoder_device_patched(): logger.debug(f"[MultiGPU Core Patching] text_encoder_device_patched returning device: {device} (current_text_encoder_device={current_text_encoder_device})") return device - logger.info(f"[MultiGPU Core Patching] Patching mm.get_torch_device and mm.text_encoder_device") logger.debug(f"[MultiGPU DEBUG] Initial current_device: {current_device}") logger.debug(f"[MultiGPU DEBUG] Initial current_text_encoder_device: {current_text_encoder_device}") mm.get_torch_device = get_torch_device_patched mm.text_encoder_device = text_encoder_device_patched -def check_module_exists(module_path): - full_path = os.path.join(folder_paths.get_folder_paths("custom_nodes")[0], module_path) - logger.debug(f"[MultiGPU] Checking for module at {full_path}") - if not os.path.exists(full_path): - logger.debug(f"[MultiGPU] Module {module_path} not found - skipping") - return False - logger.debug(f"[MultiGPU] Found {module_path}, creating compatible MultiGPU nodes") - return True - -# Import from nodes.py from .nodes import ( DeviceSelectorMultiGPU, HunyuanVideoEmbeddingsAdapter, @@ -194,7 +114,6 @@ from .nodes import ( FullCleanupMultiGPU, ) -# Import from wanvideo.py from .wanvideo import ( WanVideoModelLoader, WanVideoModelLoader_2, @@ -206,101 +125,32 @@ from .wanvideo import ( WanVideoSampler ) -# Import from distorch.py -from .distorch import ( - model_allocation_store, - create_model_hash, - register_patched_ggufmodelpatcher, - analyze_ggml_loading, - calculate_vvram_allocation_string, +from .wrappers import ( + override_class, + override_class_clip, + override_class_clip_no_device, override_class_with_distorch_gguf, override_class_with_distorch_gguf_v2, override_class_with_distorch_clip, override_class_with_distorch_clip_no_device, - override_class_with_distorch + override_class_with_distorch, + override_class_with_distorch_safetensor_v2, + override_class_with_distorch_safetensor_v2_clip, + override_class_with_distorch_safetensor_v2_clip_no_device, ) - -# Import from distorch_2.py for DisTorch v2 SafeTensor support from .distorch_2 import ( safetensor_allocation_store, create_safetensor_model_hash, register_patched_safetensor_modelpatcher, analyze_safetensor_loading, calculate_safetensor_vvram_allocation, - override_class_with_distorch_safetensor_v2, - override_class_with_distorch_safetensor_v2_clip, - override_class_with_distorch_safetensor_v2_clip_no_device ) -logger.info("[MultiGPU Core Patching] Patching mm.soft_empty_cache for Comprehensive Memory Management (VRAM + CPU + Store Pruning)") - -original_soft_empty_cache = mm.soft_empty_cache - -def soft_empty_cache_distorch2_patched(force=False): - """ - Patched mm.soft_empty_cache. - - Prunes DisTorch store bookkeeping to avoid stale references - - Manages VRAM: if DisTorch2 models are active, clear allocator caches on all devices; - otherwise delegate to original mm.soft_empty_cache. - - Manages CPU RAM: adaptive threshold-based PromptExecutor cache reset; - and force-triggered reset when explicitly requested (mirrors ComfyUI 'Free memory' button). - """ - multigpu_memory_log("patched_soft_empty", f"start:force={force}") - is_distorch_active = False - - # Detect DisTorch2-managed models - logger.mgpu_mm_log(f"[DETECT_DEBUG] Checking DisTorch2 active status - loaded models: {len(mm.current_loaded_models)}, store entries: {len(safetensor_allocation_store)}") - - for i, lm in enumerate(mm.current_loaded_models): - mp = lm.model # weakref call to ModelPatcher - if mp is not None: - try: - model_hash = create_safetensor_model_hash(mp, "cache_patch_check") - in_store = model_hash in safetensor_allocation_store - alloc_value = safetensor_allocation_store.get(model_hash, "") - model_name = type(getattr(mp, 'model', mp)).__name__ - unload_distorch_model = getattr(getattr(mp, 'model', None), '_mgpu_unload_distorch_model', False) - - logger.mgpu_mm_log(f"[DETECT_DEBUG] Model {i}: {model_name}, hash={model_hash[:8]}, in_store={in_store}, alloc_value='{alloc_value}', unload_distorch_model={unload_distorch_model}") - - if in_store and alloc_value: - is_distorch_active = True - logger.mgpu_mm_log(f"[DETECT_DEBUG] DisTorch2 ACTIVE detected on model: {model_name}") - break - except Exception as e: - logger.mgpu_mm_log(f"[DETECT_DEBUG] Model {i}: Error during detection - {e}") - - logger.mgpu_mm_log(f"[DETECT_DEBUG] Final DisTorch2 active status: {is_distorch_active}") - - # Phase 2: adaptive CPU memory management - check_cpu_memory_threshold() - - # VRAM allocator management - if is_distorch_active: - logger.mgpu_mm_log("DisTorch2 active: clearing allocator caches on all devices (VRAM)") - soft_empty_cache_multigpu() - else: - logger.mgpu_mm_log("DisTorch2 not active: delegating allocator cache clear (VRAM) to original mm.soft_empty_cache") - original_soft_empty_cache(force) - # Optional: return CPU heap to OS (not part of Comfy Core) - - # Phase 1/3: forced executor reset mirrors ComfyUI 'Free memory' semantics - if force: - logger.mgpu_mm_log("Force flag active: triggering executor cache reset (CPU)") - trigger_executor_cache_reset(reason="forced_soft_empty", force=True) - multigpu_memory_log("patched_soft_empty", "end") - -mm.soft_empty_cache = soft_empty_cache_distorch2_patched - -LARGE_MODEL_THRESHOLD = 2 * (1024**3) # 2 GB threshold for "large" models - -# Import advanced checkpoint loaders from .checkpoint_multigpu import ( CheckpointLoaderAdvancedMultiGPU, CheckpointLoaderAdvancedDisTorch2MultiGPU ) -# Initialize NODE_CLASS_MAPPINGS NODE_CLASS_MAPPINGS = { "DeviceSelectorMultiGPU": DeviceSelectorMultiGPU, "HunyuanVideoEmbeddingsAdapter": HunyuanVideoEmbeddingsAdapter, @@ -309,41 +159,29 @@ NODE_CLASS_MAPPINGS = { "UNetLoaderLP": UNetLoaderLP, } -# Standard MultiGPU nodes NODE_CLASS_MAPPINGS["UNETLoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["UNETLoader"]) NODE_CLASS_MAPPINGS["VAELoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["VAELoader"]) NODE_CLASS_MAPPINGS["CLIPLoaderMultiGPU"] = override_class_clip(GLOBAL_NODE_CLASS_MAPPINGS["CLIPLoader"]) NODE_CLASS_MAPPINGS["DualCLIPLoaderMultiGPU"] = override_class_clip(GLOBAL_NODE_CLASS_MAPPINGS["DualCLIPLoader"]) -if "TripleCLIPLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["TripleCLIPLoaderMultiGPU"] = override_class_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["TripleCLIPLoader"]) -if "QuadrupleCLIPLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["QuadrupleCLIPLoaderMultiGPU"] = override_class_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["QuadrupleCLIPLoader"]) +NODE_CLASS_MAPPINGS["TripleCLIPLoaderMultiGPU"] = override_class_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["TripleCLIPLoader"]) +NODE_CLASS_MAPPINGS["QuadrupleCLIPLoaderMultiGPU"] = override_class_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["QuadrupleCLIPLoader"]) NODE_CLASS_MAPPINGS["CLIPVisionLoaderMultiGPU"] = override_class_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["CLIPVisionLoader"]) NODE_CLASS_MAPPINGS["CheckpointLoaderSimpleMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["CheckpointLoaderSimple"]) NODE_CLASS_MAPPINGS["ControlNetLoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["ControlNetLoader"]) -if "DiffusersLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["DiffusersLoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["DiffusersLoader"]) -if "DiffControlNetLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["DiffControlNetLoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["DiffControlNetLoader"]) - -# DisTorch 2 SafeTensor nodes for FLUX and other safetensor models +NODE_CLASS_MAPPINGS["DiffusersLoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["DiffusersLoader"]) +NODE_CLASS_MAPPINGS["DiffControlNetLoaderMultiGPU"] = override_class(GLOBAL_NODE_CLASS_MAPPINGS["DiffControlNetLoader"]) NODE_CLASS_MAPPINGS["UNETLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["UNETLoader"]) NODE_CLASS_MAPPINGS["VAELoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["VAELoader"]) NODE_CLASS_MAPPINGS["CLIPLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip(GLOBAL_NODE_CLASS_MAPPINGS["CLIPLoader"]) NODE_CLASS_MAPPINGS["DualCLIPLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip(GLOBAL_NODE_CLASS_MAPPINGS["DualCLIPLoader"]) -if "TripleCLIPLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["TripleCLIPLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["TripleCLIPLoader"]) -if "QuadrupleCLIPLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["QuadrupleCLIPLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["QuadrupleCLIPLoader"]) +NODE_CLASS_MAPPINGS["TripleCLIPLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["TripleCLIPLoader"]) +NODE_CLASS_MAPPINGS["QuadrupleCLIPLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["QuadrupleCLIPLoader"]) NODE_CLASS_MAPPINGS["CLIPVisionLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2_clip_no_device(GLOBAL_NODE_CLASS_MAPPINGS["CLIPVisionLoader"]) NODE_CLASS_MAPPINGS["CheckpointLoaderSimpleDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["CheckpointLoaderSimple"]) NODE_CLASS_MAPPINGS["ControlNetLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["ControlNetLoader"]) -if "DiffusersLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["DiffusersLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["DiffusersLoader"]) -if "DiffControlNetLoader" in GLOBAL_NODE_CLASS_MAPPINGS: - NODE_CLASS_MAPPINGS["DiffControlNetLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["DiffControlNetLoader"]) +NODE_CLASS_MAPPINGS["DiffusersLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["DiffusersLoader"]) +NODE_CLASS_MAPPINGS["DiffControlNetLoaderDisTorch2MultiGPU"] = override_class_with_distorch_safetensor_v2(GLOBAL_NODE_CLASS_MAPPINGS["DiffControlNetLoader"]) -# --- Registration Table --- logger.info("[MultiGPU] Initiating custom_node Registration. . .") dash_line = "-" * 47 fmt_reg = "{:<30}{:>5}{:>10}" @@ -370,26 +208,21 @@ def register_and_count(module_names, node_map): registration_data.append({"name": module_names[0], "found": "Y" if found else "N", "count": count}) return found -# ComfyUI-LTXVideo ltx_nodes = {"LTXVLoaderMultiGPU": override_class(LTXVLoader)} register_and_count(["ComfyUI-LTXVideo", "comfyui-ltxvideo"], ltx_nodes) -# ComfyUI-Florence2 florence_nodes = { "Florence2ModelLoaderMultiGPU": override_class(Florence2ModelLoader), "DownloadAndLoadFlorence2ModelMultiGPU": override_class(DownloadAndLoadFlorence2Model) } register_and_count(["ComfyUI-Florence2", "comfyui-florence2"], florence_nodes) -# ComfyUI_bitsandbytes_NF4 nf4_nodes = {"CheckpointLoaderNF4MultiGPU": override_class(CheckpointLoaderNF4)} register_and_count(["ComfyUI_bitsandbytes_NF4", "comfyui_bitsandbytes_nf4"], nf4_nodes) -# x-flux-comfyui flux_controlnet_nodes = {"LoadFluxControlNetMultiGPU": override_class(LoadFluxControlNet)} register_and_count(["x-flux-comfyui"], flux_controlnet_nodes) -# ComfyUI-MMAudio mmaudio_nodes = { "MMAudioModelLoaderMultiGPU": override_class(MMAudioModelLoader), "MMAudioFeatureUtilsLoaderMultiGPU": override_class(MMAudioFeatureUtilsLoader), @@ -397,7 +230,6 @@ mmaudio_nodes = { } register_and_count(["ComfyUI-MMAudio", "comfyui-mmaudio"], mmaudio_nodes) -# ComfyUI-GGUF gguf_nodes = { "UnetLoaderGGUFDisTorchMultiGPU": override_class_with_distorch_gguf(UnetLoaderGGUF), "UnetLoaderGGUFAdvancedDisTorchMultiGPU": override_class_with_distorch_gguf(UnetLoaderGGUFAdvanced), @@ -420,7 +252,6 @@ gguf_nodes = { } register_and_count(["ComfyUI-GGUF", "comfyui-gguf"], gguf_nodes) -# PuLID_ComfyUI pulid_nodes = { "PulidModelLoaderMultiGPU": override_class(PulidModelLoader), "PulidInsightFaceLoaderMultiGPU": override_class(PulidInsightFaceLoader), @@ -428,7 +259,6 @@ pulid_nodes = { } register_and_count(["PuLID_ComfyUI", "pulid_comfyui"], pulid_nodes) -# ComfyUI-HunyuanVideoWrapper hunyuan_nodes = { "HyVideoModelLoaderMultiGPU": override_class(HyVideoModelLoader), "HyVideoVAELoaderMultiGPU": override_class(HyVideoVAELoader), @@ -436,7 +266,6 @@ hunyuan_nodes = { } register_and_count(["ComfyUI-HunyuanVideoWrapper", "comfyui-hunyuanvideowrapper"], hunyuan_nodes) -# ComfyUI-WanVideoWrapper wanvideo_nodes = { "WanVideoModelLoaderMultiGPU": WanVideoModelLoader, "WanVideoModelLoaderMultiGPU_2": WanVideoModelLoader_2, @@ -449,13 +278,8 @@ wanvideo_nodes = { } register_and_count(["ComfyUI-WanVideoWrapper", "comfyui-wanvideowrapper"], wanvideo_nodes) -# Print the registration table for item in registration_data: logger.info(fmt_reg.format(item['name'], item['found'], str(item['count']))) logger.info(dash_line) - -# Register maintenance node -NODE_CLASS_MAPPINGS["FullCleanupMultiGPU"] = FullCleanupMultiGPU - logger.info(f"[MultiGPU] Registration complete. Final mappings: {', '.join(NODE_CLASS_MAPPINGS.keys())}") diff --git a/device_utils.py b/device_utils.py index 96ce073..d48962a 100644 --- a/device_utils.py +++ b/device_utils.py @@ -291,6 +291,74 @@ def soft_empty_cache_multigpu(): multigpu_memory_log("general", "post-soft-empty") + +# ========================================================================================== +# Comprehensive Memory Management (VRAM + CPU + Store Pruning) +# ========================================================================================== + +logger.info("[MultiGPU Core Patching] Patching mm.soft_empty_cache for Comprehensive Memory Management (VRAM + CPU + Store Pruning)") + +original_soft_empty_cache = mm.soft_empty_cache + +def soft_empty_cache_distorch2_patched(force=False): + """ + Patched mm.soft_empty_cache. + - Prunes DisTorch store bookkeeping to avoid stale references + - Manages VRAM: if DisTorch2 models are active, clear allocator caches on all devices; + otherwise delegate to original mm.soft_empty_cache. + - Manages CPU RAM: adaptive threshold-based PromptExecutor cache reset; + and force-triggered reset when explicitly requested (mirrors ComfyUI 'Free memory' button). + """ + from .model_management_mgpu import multigpu_memory_log, check_cpu_memory_threshold, trigger_executor_cache_reset + from .distorch_2 import safetensor_allocation_store, create_safetensor_model_hash + + multigpu_memory_log("patched_soft_empty", f"start:force={force}") + is_distorch_active = False + + # Detect DisTorch2-managed models + logger.mgpu_mm_log(f"[DETECT_DEBUG] Checking DisTorch2 active status - loaded models: {len(mm.current_loaded_models)}, store entries: {len(safetensor_allocation_store)}") + + for i, lm in enumerate(mm.current_loaded_models): + mp = lm.model # weakref call to ModelPatcher + if mp is not None: + try: + model_hash = create_safetensor_model_hash(mp, "cache_patch_check") + in_store = model_hash in safetensor_allocation_store + alloc_value = safetensor_allocation_store.get(model_hash, "") + model_name = type(getattr(mp, 'model', mp)).__name__ + unload_distorch_model = getattr(getattr(mp, 'model', None), '_mgpu_unload_distorch_model', False) + + logger.mgpu_mm_log(f"[DETECT_DEBUG] Model {i}: {model_name}, hash={model_hash[:8]}, in_store={in_store}, alloc_value='{alloc_value}', unload_distorch_model={unload_distorch_model}") + + if in_store and alloc_value: + is_distorch_active = True + logger.mgpu_mm_log(f"[DETECT_DEBUG] DisTorch2 ACTIVE detected on model: {model_name}") + break + except Exception as e: + logger.mgpu_mm_log(f"[DETECT_DEBUG] Model {i}: Error during detection - {e}") + + logger.mgpu_mm_log(f"[DETECT_DEBUG] Final DisTorch2 active status: {is_distorch_active}") + + # Phase 2: adaptive CPU memory management + check_cpu_memory_threshold() + + # VRAM allocator management + if is_distorch_active: + logger.mgpu_mm_log("DisTorch2 active: clearing allocator caches on all devices (VRAM)") + soft_empty_cache_multigpu() + else: + logger.mgpu_mm_log("DisTorch2 not active: delegating allocator cache clear (VRAM) to original mm.soft_empty_cache") + original_soft_empty_cache(force) + # Optional: return CPU heap to OS (not part of Comfy Core) + + # Phase 1/3: forced executor reset mirrors ComfyUI 'Free memory' semantics + if force: + logger.mgpu_mm_log("Force flag active: triggering executor cache reset (CPU)") + trigger_executor_cache_reset(reason="forced_soft_empty", force=True) + multigpu_memory_log("patched_soft_empty", "end") + +mm.soft_empty_cache = soft_empty_cache_distorch2_patched + # ========================================================================================== # Memory Inspection Utilities # ========================================================================================== diff --git a/distorch.py b/distorch.py deleted file mode 100644 index 113aba9..0000000 --- a/distorch.py +++ /dev/null @@ -1,529 +0,0 @@ -""" -DisTorch GGUF/GGML Memory Management Module -Contains all GGUF/GGML related code for distributed memory management -""" - -import sys -import torch -import logging -import hashlib - -logger = logging.getLogger("MultiGPU") -import copy -from collections import defaultdict -import comfy.model_management as mm -from .device_utils import get_device_list, soft_empty_cache_multigpu -from .model_management_mgpu import multigpu_memory_log - -# Global store for model allocations -model_allocation_store = {} - - -def create_model_hash(model, caller): - """Create a unique hash for a model to track allocations""" - model_type = type(model.model).__name__ - model_size = model.model_size() - first_layers = str(list(model.model_state_dict().keys())[:3]) - identifier = f"{model_type}_{model_size}_{first_layers}" - final_hash = hashlib.sha256(identifier.encode()).hexdigest() - logger.debug(f"[MultiGPU_DisTorch_HASH] Created hash for {caller}: {final_hash[:8]}...") - return final_hash - - -def register_patched_ggufmodelpatcher(): - """Register and patch the GGUFModelPatcher for distributed loading""" - from nodes import NODE_CLASS_MAPPINGS - original_loader = NODE_CLASS_MAPPINGS["UnetLoaderGGUF"] - module = sys.modules[original_loader.__module__] - - if not hasattr(module.GGUFModelPatcher, '_patched'): - original_load = module.GGUFModelPatcher.load - - def new_load(self, *args, force_patch_weights=False, **kwargs): - global model_allocation_store - - debug_hash = create_model_hash(self, "patcher") - multigpu_memory_log(f"gguf:{debug_hash[:8]}", "pre-load") - super(module.GGUFModelPatcher, self).load(*args, force_patch_weights=True, **kwargs) - multigpu_memory_log(f"gguf:{debug_hash[:8]}", "post-load") - linked = [] - module_count = 0 - for n, m in self.model.named_modules(): - module_count += 1 - if hasattr(m, "weight"): - device = getattr(m.weight, "device", None) - if device is not None: - linked.append((n, m)) - continue - if hasattr(m, "bias"): - device = getattr(m.bias, "device", None) - if device is not None: - linked.append((n, m)) - continue - if linked: - if hasattr(self, 'model'): - debug_hash = create_model_hash(self, "patcher") - debug_allocations = model_allocation_store.get(debug_hash) - if debug_allocations: - logger.info("[MultiGPU DisTorch GGUF] Invoking soft_empty_cache_multigpu before GGUF device assignment") - soft_empty_cache_multigpu() - device_assignments = analyze_ggml_loading(self.model, debug_allocations)['device_assignments'] - for device, layers in device_assignments.items(): - target_device = torch.device(device) - for n, m, _ in layers: - m.to(self.load_device).to(target_device) - - self.mmap_released = True - - module.GGUFModelPatcher.load = new_load - module.GGUFModelPatcher._patched = True - - -def analyze_ggml_loading(model, allocations_str): - """Analyze and distribute GGML model layers across devices""" - DEVICE_RATIOS_DISTORCH = {} - device_table = {} - distorch_alloc = allocations_str - virtual_vram_gb = 0.0 - - if '#' in allocations_str: - distorch_alloc, virtual_vram_str = allocations_str.split('#') - if not distorch_alloc: - distorch_alloc = calculate_vvram_allocation_string(model, virtual_vram_str) - - eq_line = "=" * 47 - dash_line = "-" * 47 - fmt_assign = "{:<12}{:>10}{:>14}{:>10}" - - for allocation in distorch_alloc.split(';'): - dev_name, fraction = allocation.split(',') - fraction = float(fraction) - total_mem_bytes = mm.get_total_memory(torch.device(dev_name)) - alloc_gb = (total_mem_bytes * fraction) / (1024**3) - DEVICE_RATIOS_DISTORCH[dev_name] = alloc_gb - device_table[dev_name] = { - "fraction": fraction, - "total_gb": total_mem_bytes / (1024**3), - "alloc_gb": alloc_gb - } - - logger.info(eq_line) - logger.info(" DisTorch Model Device Allocations") - logger.info(eq_line) - logger.info(fmt_assign.format("Device", "Alloc %", "Total (GB)", " Alloc (GB)")) - logger.info(dash_line) - - sorted_devices = sorted(device_table.keys(), key=lambda d: (d == "cpu", d)) - - for dev in sorted_devices: - frac = device_table[dev]["fraction"] - tot_gb = device_table[dev]["total_gb"] - alloc_gb = device_table[dev]["alloc_gb"] - logger.info(fmt_assign.format(dev,f"{int(frac * 100)}%",f"{tot_gb:.2f}",f"{alloc_gb:.2f}")) - - logger.info(dash_line) - - layer_summary = {} - layer_list = [] - memory_by_type = defaultdict(int) - total_memory = 0 - - for name, module in model.named_modules(): - if hasattr(module, "weight"): - layer_type = type(module).__name__ - layer_summary[layer_type] = layer_summary.get(layer_type, 0) + 1 - layer_list.append((name, module, layer_type)) - layer_memory = 0 - if module.weight is not None: - layer_memory += module.weight.numel() * module.weight.element_size() - if hasattr(module, "bias") and module.bias is not None: - layer_memory += module.bias.numel() * module.bias.element_size() - memory_by_type[layer_type] += layer_memory - total_memory += layer_memory - - logger.info(" DisTorch Model Layer Distribution") - logger.info(dash_line) - fmt_layer = "{:<12}{:>10}{:>14}{:>10}" - logger.info(fmt_layer.format("Layer Type", "Layers", "Memory (MB)", "% Total")) - logger.info(dash_line) - for layer_type, count in layer_summary.items(): - mem_mb = memory_by_type[layer_type] / (1024 * 1024) - mem_percent = (memory_by_type[layer_type] / total_memory) * 100 if total_memory > 0 else 0 - logger.info(fmt_layer.format(layer_type,str(count),f"{mem_mb:.2f}",f"{mem_percent:.1f}%")) - logger.info(dash_line) - - nonzero_devices = [d for d, r in DEVICE_RATIOS_DISTORCH.items() if r > 0] - nonzero_total_ratio = sum(DEVICE_RATIOS_DISTORCH[d] for d in nonzero_devices) - device_assignments = {device: [] for device in DEVICE_RATIOS_DISTORCH.keys()} - total_layers = len(layer_list) - current_layer = 0 - - for idx, device in enumerate(nonzero_devices): - ratio = DEVICE_RATIOS_DISTORCH[device] - if idx == len(nonzero_devices) - 1: - device_layer_count = total_layers - current_layer - else: - device_layer_count = int((ratio / nonzero_total_ratio) * total_layers) - start_idx = current_layer - end_idx = current_layer + device_layer_count - device_assignments[device] = layer_list[start_idx:end_idx] - current_layer += device_layer_count - - logger.info("DisTorch Model Final Device/Layer Assignments") - logger.info(dash_line) - fmt_assign = "{:<12}{:>10}{:>14}{:>10}" - logger.info(fmt_assign.format("Device", "Layers", "Memory (MB)", "% Total")) - logger.info(dash_line) - total_assigned_memory = 0 - device_memories = {} - for device, layers in device_assignments.items(): - device_memory = 0 - for layer_type in layer_summary: - type_layers = sum(1 for _, _, lt in layers if lt == layer_type) - if layer_summary[layer_type] > 0: - mem_per_layer = memory_by_type[layer_type] / layer_summary[layer_type] - device_memory += mem_per_layer * type_layers - device_memories[device] = device_memory - total_assigned_memory += device_memory - - sorted_assignments = sorted(device_assignments.keys(), key=lambda d: (d == "cpu", d)) - - for dev in sorted_assignments: - layers = device_assignments[dev] - mem_mb = device_memories[dev] / (1024 * 1024) - mem_percent = (device_memories[dev] / total_memory) * 100 if total_memory > 0 else 0 - logger.info(fmt_assign.format(dev,str(len(layers)),f"{mem_mb:.2f}",f"{mem_percent:.1f}%")) - logger.info(dash_line) - - return {"device_assignments": device_assignments} - - -def calculate_vvram_allocation_string(model, virtual_vram_str): - """Calculate virtual VRAM allocation string for distributed loading""" - recipient_device, vram_amount, donors = virtual_vram_str.split(';') - virtual_vram_gb = float(vram_amount) - - eq_line = "=" * 47 - dash_line = "-" * 47 - fmt_assign = "{:<8} {:<6} {:>11} {:>9} {:>9}" - - logger.info(eq_line) - logger.info(" DisTorch Model Virtual VRAM Analysis") - logger.info(eq_line) - logger.info(fmt_assign.format("Object", "Role", "Original(GB)", "Total(GB)", "Virt(GB)")) - logger.info(dash_line) - - recipient_vram = mm.get_total_memory(torch.device(recipient_device)) / (1024**3) - recipient_virtual = recipient_vram + virtual_vram_gb - - logger.info(fmt_assign.format(recipient_device, 'recip', f"{recipient_vram:.2f}GB",f"{recipient_virtual:.2f}GB", f"+{virtual_vram_gb:.2f}GB")) - - ram_donors = [d for d in donors.split(',') if d != 'cpu'] - remaining_vram_needed = virtual_vram_gb - - donor_device_info = {} - donor_allocations = {} - - for donor in ram_donors: - donor_vram = mm.get_total_memory(torch.device(donor)) / (1024**3) - max_donor_capacity = donor_vram * 0.9 - - donation = min(remaining_vram_needed, max_donor_capacity) - donor_virtual = donor_vram - donation - remaining_vram_needed -= donation - donor_allocations[donor] = donation - - donor_device_info[donor] = (donor_vram, donor_virtual) - logger.info(fmt_assign.format(donor, 'donor', f"{donor_vram:.2f}GB", f"{donor_virtual:.2f}GB", f"-{donation:.2f}GB")) - - system_dram_gb = mm.get_total_memory(torch.device('cpu')) / (1024**3) - cpu_donation = remaining_vram_needed - cpu_virtual = system_dram_gb - cpu_donation - donor_allocations['cpu'] = cpu_donation - logger.info(fmt_assign.format('cpu', 'donor', f"{system_dram_gb:.2f}GB", f"{cpu_virtual:.2f}GB", f"-{cpu_donation:.2f}GB")) - - logger.info(dash_line) - - layer_summary = {} - layer_list = [] - memory_by_type = defaultdict(int) - total_memory = 0 - - for name, module in model.named_modules(): - if hasattr(module, "weight"): - layer_type = type(module).__name__ - layer_summary[layer_type] = layer_summary.get(layer_type, 0) + 1 - layer_list.append((name, module, layer_type)) - layer_memory = 0 - if module.weight is not None: - layer_memory += module.weight.numel() * module.weight.element_size() - if hasattr(module, "bias") and module.bias is not None: - layer_memory += module.bias.numel() * module.bias.element_size() - memory_by_type[layer_type] += layer_memory - total_memory += layer_memory - - model_size_gb = total_memory / (1024**3) - new_model_size_gb = max(0, model_size_gb - virtual_vram_gb) - - logger.info(fmt_assign.format('model', 'model', f"{model_size_gb:.2f}GB",f"{new_model_size_gb:.2f}GB", f"-{virtual_vram_gb:.2f}GB")) - - if model_size_gb > (recipient_vram * 0.9): - on_recipient = recipient_vram * 0.9 - on_virtuals = model_size_gb - on_recipient - logger.info(f"\nWarning: Model size is greater than 90% of recipient VRAM. {on_virtuals:.2f} GB of GGML Layers Offloaded Automatically to Virtual VRAM.\n") - else: - on_recipient = model_size_gb - on_virtuals = 0 - - new_on_recipient = max(0, on_recipient - virtual_vram_gb) - - allocation_parts = [] - recipient_percent = new_on_recipient / recipient_vram - allocation_parts.append(f"{recipient_device},{recipient_percent:.4f}") - - for donor in ram_donors: - donor_vram = donor_device_info[donor][0] - donor_percent = donor_allocations[donor] / donor_vram - allocation_parts.append(f"{donor},{donor_percent:.4f}") - - cpu_percent = donor_allocations['cpu'] / system_dram_gb - allocation_parts.append(f"cpu,{cpu_percent:.4f}") - - allocation_string = ";".join(allocation_parts) - fmt_mem = "{:<20}{:>20}" - logger.info(fmt_mem.format("\n v1 Expert String", allocation_string)) - - return allocation_string - - -def override_class_with_distorch_gguf(cls): - """Legacy DisTorch wrapper for GGUF models for backward compatibility.""" - from . import current_device - - class NodeOverrideDisTorchGGUFLegacy(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 24.0, "step": 0.1}) - inputs["optional"]["use_other_vram"] = ("BOOLEAN", {"default": False}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", { - "multiline": False, - "default": "", - }) - return inputs - - CATEGORY = "multigpu/legacy" - FUNCTION = "override" - if hasattr(cls, 'TITLE'): - TITLE = f"{cls.TITLE} (Legacy)" - else: - TITLE = "Legacy DisTorch Node" - - def override(self, *args, device=None, expert_mode_allocations=None, use_other_vram=None, virtual_vram_gb=0.0, **kwargs): - from . import set_current_device - if device is not None: - set_current_device(device) - - register_patched_ggufmodelpatcher() - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - vram_string = "" - if virtual_vram_gb > 0: - if use_other_vram: - available_devices = [d for d in get_device_list() if d != "cpu"] - other_devices = [d for d in available_devices if d != device] - other_devices.sort(key=lambda x: int(x.split(':')[1] if ':' in x else x[-1]), reverse=False) - device_string = ','.join(other_devices + ['cpu']) - vram_string = f"{device};{virtual_vram_gb};{device_string}" - else: - vram_string = f"{device};{virtual_vram_gb};cpu" - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - if hasattr(out[0], 'model'): - model_hash = create_model_hash(out[0], "override") - model_allocation_store[model_hash] = full_allocation - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_hash = create_model_hash(out[0].patcher, "override") - model_allocation_store[model_hash] = full_allocation - - return out - - return NodeOverrideDisTorchGGUFLegacy - - -def override_class_with_distorch_gguf_v2(cls): - """DisTorch 2.0 wrapper for GGUF models.""" - from . import current_device - - class NodeOverrideDisTorchGGUFv2(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - compute_device = devices[1] if len(devices) > 1 else devices[0] - - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["compute_device"] = (devices, {"default": compute_device}) - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 128.0, "step": 0.1}) - inputs["optional"]["donor_device"] = (devices, {"default": "cpu"}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) - return inputs - - CATEGORY = "multigpu/distorch_2" - FUNCTION = "override" - - def override(self, *args, compute_device=None, virtual_vram_gb=4.0, - donor_device="cpu", expert_mode_allocations="", **kwargs): - from . import set_current_device - if compute_device is not None: - set_current_device(compute_device) - - register_patched_ggufmodelpatcher() - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - vram_string = "" - if virtual_vram_gb > 0: - vram_string = f"{compute_device};{virtual_vram_gb};{donor_device}" - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - logger.info(f"[MultiGPU_DisTorch] Full allocation string: {full_allocation}") - - if hasattr(out[0], 'model'): - model_hash = create_model_hash(out[0], "override") - model_allocation_store[model_hash] = full_allocation - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_hash = create_model_hash(out[0].patcher, "override") - model_allocation_store[model_hash] = full_allocation - - return out - - return NodeOverrideDisTorchGGUFv2 - - -def override_class_with_distorch_clip(cls): - """DisTorch wrapper for CLIP models with GGUF support""" - from . import current_text_encoder_device - - class NodeOverrideDisTorch(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 24.0, "step": 0.1}) - inputs["optional"]["use_other_vram"] = ("BOOLEAN", {"default": False}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", { - "multiline": False, - "default": "", - "tooltip": "Expert use only: Manual VRAM allocation string. Incorrect values can cause crashes. Do not modify unless you fully understand DisTorch memory management." - }) - return inputs - - CATEGORY = "multigpu" - FUNCTION = "override" - - def override(self, *args, device=None, expert_mode_allocations=None, use_other_vram=None, virtual_vram_gb=0.0, **kwargs): - from . import set_current_text_encoder_device - if device is not None: - set_current_text_encoder_device(device) - - register_patched_ggufmodelpatcher() - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - vram_string = "" - if virtual_vram_gb > 0: - if use_other_vram: - available_devices = [d for d in get_device_list() if d != "cpu"] - other_devices = [d for d in available_devices if d != device] - other_devices.sort(key=lambda x: int(x.split(':')[1] if ':' in x else x[-1]), reverse=False) - device_string = ','.join(other_devices + ['cpu']) - vram_string = f"{device};{virtual_vram_gb};{device_string}" - else: - vram_string = f"{device};{virtual_vram_gb};cpu" - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - logging.info(f"[MultiGPU_DisTorch] Full allocation string: {full_allocation}") - - if hasattr(out[0], 'model'): - model_hash = create_model_hash(out[0], "override") - model_allocation_store[model_hash] = full_allocation - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_hash = create_model_hash(out[0].patcher, "override") - model_allocation_store[model_hash] = full_allocation - - return out - - return NodeOverrideDisTorch -def override_class_with_distorch_clip_no_device(cls): - """DisTorch wrapper for CLIP models with GGUF support""" - from . import current_text_encoder_device - - class NodeOverrideDisTorchClipNoDevice(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 24.0, "step": 0.1}) - inputs["optional"]["use_other_vram"] = ("BOOLEAN", {"default": False}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", { - "multiline": False, - "default": "", - "tooltip": "Expert use only: Manual VRAM allocation string. Incorrect values can cause crashes. Do not modify unless you fully understand DisTorch memory management." - }) - return inputs - - CATEGORY = "multigpu" - FUNCTION = "override" - - def override(self, *args, device=None, expert_mode_allocations=None, use_other_vram=None, virtual_vram_gb=0.0, **kwargs): - from . import set_current_text_encoder_device - if device is not None: - set_current_text_encoder_device(device) - - register_patched_ggufmodelpatcher() - fn = getattr(super(), cls.FUNCTION) - out = fn(*args, **kwargs) - - vram_string = "" - if virtual_vram_gb > 0: - if use_other_vram: - available_devices = [d for d in get_device_list() if d != "cpu"] - other_devices = [d for d in available_devices if d != device] - other_devices.sort(key=lambda x: int(x.split(':')[1] if ':' in x else x[-1]), reverse=False) - device_string = ','.join(other_devices + ['cpu']) - vram_string = f"{device};{virtual_vram_gb};{device_string}" - else: - vram_string = f"{device};{virtual_vram_gb};cpu" - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - logging.info(f"[MultiGPU_DisTorch] Full allocation string: {full_allocation}") - - if hasattr(out[0], 'model'): - model_hash = create_model_hash(out[0], "override") - model_allocation_store[model_hash] = full_allocation - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_hash = create_model_hash(out[0].patcher, "override") - model_allocation_store[model_hash] = full_allocation - - return out - - return NodeOverrideDisTorchClipNoDevice - -# Alias for backward compatibility -override_class_with_distorch = override_class_with_distorch_gguf diff --git a/distorch_2.py b/distorch_2.py index f595397..4769c45 100644 --- a/distorch_2.py +++ b/distorch_2.py @@ -843,406 +843,9 @@ def calculate_safetensor_vvram_allocation(model_patcher, virtual_vram_str): allocations_string = ";".join(allocation_parts) return allocations_string -def override_class_with_distorch_safetensor_v2(cls): - """DisTorch 2.0 wrapper for safetensor models""" - - class NodeOverrideDisTorchSafetensorV2(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - compute_device = devices[1] if len(devices) > 1 else devices[0] - - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["compute_device"] = (devices, {"default": compute_device}) - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 128.0, "step": 0.1}) - inputs["optional"]["donor_device"] = (devices, {"default": "cpu"}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) - inputs["optional"]["keep_loaded"] = ("BOOLEAN", {"default": True}) - return inputs - - CATEGORY = "multigpu/distorch_2" - FUNCTION = "override" - TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch2)" - - @classmethod - def IS_CHANGED(s, *args, compute_device=None, virtual_vram_gb=4.0, - donor_device="cpu", expert_mode_allocations="", keep_loaded=True, **kwargs): - settings_str = f"{compute_device}{virtual_vram_gb}{donor_device}{expert_mode_allocations}{keep_loaded}" - current_hash = hashlib.sha256(settings_str.encode()).hexdigest() - - if not hasattr(cls, '_last_hash'): - cls._last_hash = current_hash - logger.mgpu_mm_log(f"IS_CHANGED first call: {current_hash[:8]}") - elif cls._last_hash != current_hash: - cls._last_hash = current_hash - logger.mgpu_mm_log(f"IS_CHANGED CHANGED: {current_hash[:8]} ← settings changed") - return current_hash - - def override(self, *args, compute_device=None, virtual_vram_gb=4.0, - donor_device="cpu", expert_mode_allocations="", keep_loaded=True, **kwargs): - - unload_distorch_model = not keep_loaded - - from . import set_current_device - if compute_device is not None: - set_current_device(compute_device) - - # Register our patched ModelPatcher - register_patched_safetensor_modelpatcher() - - # Build allocation string - vram_string = "" - if virtual_vram_gb > 0: - vram_string = f"{compute_device};{virtual_vram_gb};{donor_device}" - elif expert_mode_allocations: # Only include compute device if there's an expert string - vram_string = compute_device - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - fn = getattr(super(), cls.FUNCTION) - - # Load the model and get hash, then store allocation for future runs - out = fn(*args, **kwargs) - - model_to_check = None - if hasattr(out[0], 'model'): - model_to_check = out[0] - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_to_check = out[0].patcher - - if model_to_check: - model_hash = create_safetensor_model_hash(model_to_check, "override_store") - settings_str = f"{compute_device}{virtual_vram_gb}{donor_device}{expert_mode_allocations}" - settings_hash = hashlib.sha256(settings_str.encode()).hexdigest() - - # Store allocation for next run - this enables DisTorch for subsequent loads - safetensor_allocation_store[model_hash] = full_allocation - safetensor_settings_store[model_hash] = settings_hash - logger.debug(f"[MultiGPU DisTorch V2] Stored allocation for model {model_hash[:8]}: {full_allocation}") - - logger.info(f"[MultiGPU DisTorch V2] Full allocation string: {full_allocation}") - - logger.mgpu_mm_log(f"[FLAG_SET_START] Setting '_mgpu_unload_distorch_model' to: {unload_distorch_model} (keep_loaded={keep_loaded})") - - # DIAGNOSTIC: Log full object chain at SET time - if hasattr(out[0], 'model'): - mp = out[0] # This is the ModelPatcher - mp_id = id(mp) - inner_model = getattr(mp, 'model', None) - inner_model_id = id(inner_model) if inner_model else None - inner_model_name = type(inner_model).__name__ if inner_model else "None" - - # Format inner_model_id properly for f-string - inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" - - logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] ModelPatcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") - - # FIX: Store flag on ModelPatcher itself (not inner model) - # This aligns with where it will be READ in model_management_mgpu.py - mp._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") - - # Also set on inner model for backwards compatibility during transition - if inner_model: - inner_model._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") - - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - mp = out[0].patcher # This is the ModelPatcher - mp_id = id(mp) - inner_model = getattr(mp, 'model', None) - inner_model_id = id(inner_model) if inner_model else None - inner_model_name = type(inner_model).__name__ if inner_model else "None" - - # Format inner_model_id properly for f-string - inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" - - logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] ModelPatcher via patcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") - - # FIX: Store flag on ModelPatcher itself - mp._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") - - # Also set on inner model for backwards compatibility - if inner_model: - inner_model._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") - - if unload_distorch_model: - logger.mgpu_mm_log("[FLAG_TRIGGER] unload_distorch_model=True, triggering full system cleanup") - force_full_system_cleanup(reason="policy_every_load", force=True) - - return out - - return NodeOverrideDisTorchSafetensorV2 - - -def override_class_with_distorch_safetensor_v2_clip(cls): - """DisTorch 2.0 wrapper for safetensor CLIP models""" - - class NodeOverrideDisTorchSafetensorV2Clip(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) # Changed from compute_device - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 128.0, "step": 0.1}) - inputs["optional"]["donor_device"] = (devices, {"default": "cpu"}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) - inputs["optional"]["keep_loaded"] = ("BOOLEAN", {"default": True}) - return inputs - - CATEGORY = "multigpu/distorch_2" - FUNCTION = "override" - TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch2)" - - @classmethod - def IS_CHANGED(s, *args, device=None, virtual_vram_gb=4.0, # Changed from compute_device - donor_device="cpu", expert_mode_allocations="", keep_loaded=True, **kwargs): - # Create a hash of our specific settings - settings_str = f"{device}{virtual_vram_gb}{donor_device}{expert_mode_allocations}{keep_loaded}" # Changed from compute_device - current_hash = hashlib.sha256(settings_str.encode()).hexdigest() - - if not hasattr(cls, '_last_hash'): - cls._last_hash = current_hash - logger.mgpu_mm_log(f"IS_CHANGED first call: {current_hash[:8]}") - elif cls._last_hash != current_hash: - cls._last_hash = current_hash - logger.mgpu_mm_log(f"IS_CHANGED CHANGED: {current_hash[:8]} ← settings changed") - return current_hash - - def override(self, *args, device=None, virtual_vram_gb=4.0, # Changed from compute_device - donor_device="cpu", expert_mode_allocations="", keep_loaded=True, **kwargs): - - unload_distorch_model = not keep_loaded - - from . import set_current_text_encoder_device # Use text encoder device setter - if device is not None: - set_current_text_encoder_device(device) - - kwargs['device'] = 'default' # Hardcode device setting like in standard clip wrapper - - # Register our patched ModelPatcher - register_patched_safetensor_modelpatcher() - - # Call original function - fn = getattr(super(), cls.FUNCTION) - - # Call the main function once - out = fn(*args, **kwargs) - - logger.mgpu_mm_log(f"[FLAG_SET_START] Setting '_mgpu_unload_distorch_model' to: {unload_distorch_model} (keep_loaded={keep_loaded})") - - # DIAGNOSTIC: Log full object chain at SET time - if hasattr(out[0], 'model'): - mp = out[0] # This is the ModelPatcher - mp_id = id(mp) - inner_model = getattr(mp, 'model', None) - inner_model_id = id(inner_model) if inner_model else None - inner_model_name = type(inner_model).__name__ if inner_model else "None" - - # Format inner_model_id properly for f-string - inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" - - logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] CLIP ModelPatcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") - - # FIX: Store flag on ModelPatcher itself - mp._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") - - # Also set on inner model for backwards compatibility - if inner_model: - inner_model._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") - - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - mp = out[0].patcher # This is the ModelPatcher - mp_id = id(mp) - inner_model = getattr(mp, 'model', None) - inner_model_id = id(inner_model) if inner_model else None - inner_model_name = type(inner_model).__name__ if inner_model else "None" - - # Format inner_model_id properly for f-string - inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" - - logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] CLIP ModelPatcher via patcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") - - # FIX: Store flag on ModelPatcher itself - mp._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") - - # Also set on inner model for backwards compatibility - if inner_model: - inner_model._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") - - vram_string = "" - if virtual_vram_gb > 0: - vram_string = f"{device};{virtual_vram_gb};{donor_device}" - elif expert_mode_allocations: - vram_string = device - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - logger.info(f"[MultiGPU DisTorch V2] Full allocation string: {full_allocation}") - - # Store allocation AFTER loading for next time - model_to_check = None - if hasattr(out[0], 'model'): - model_to_check = out[0] - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_to_check = out[0].patcher - - if model_to_check: - model_hash = create_safetensor_model_hash(model_to_check, "override_store") - settings_str = f"{device}{virtual_vram_gb}{donor_device}{expert_mode_allocations}" - settings_hash = hashlib.sha256(settings_str.encode()).hexdigest() - - # Store allocation for next time - safetensor_allocation_store[model_hash] = full_allocation - safetensor_settings_store[model_hash] = settings_hash - - if unload_distorch_model: - logger.mgpu_mm_log("[FLAG_TRIGGER] unload_distorch_model=True, triggering full system cleanup") - force_full_system_cleanup(reason="policy_every_load", force=True) - - return out - - return NodeOverrideDisTorchSafetensorV2Clip - -def override_class_with_distorch_safetensor_v2_clip_no_device(cls): - """DisTorch 2.0 wrapper for safetensor CLIP models""" - - class NodeOverrideDisTorchSafetensorV2ClipNoDevice(cls): - @classmethod - def INPUT_TYPES(s): - inputs = copy.deepcopy(cls.INPUT_TYPES()) - devices = get_device_list() - default_device = devices[1] if len(devices) > 1 else devices[0] - - inputs["optional"] = inputs.get("optional", {}) - inputs["optional"]["device"] = (devices, {"default": default_device}) # Changed from compute_device - inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 128.0, "step": 0.1}) - inputs["optional"]["donor_device"] = (devices, {"default": "cpu"}) - inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) - inputs["optional"]["keep_loaded"] = ("BOOLEAN", {"default": True}) - return inputs - - CATEGORY = "multigpu/distorch_2" - FUNCTION = "override" - TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch2)" - - @classmethod - def IS_CHANGED(s, *args, device=None, virtual_vram_gb=4.0, # Changed from compute_device - donor_device="cpu", expert_mode_allocations="", keep_loaded=True, **kwargs): - # Create a hash of our specific settings - settings_str = f"{device}{virtual_vram_gb}{donor_device}{expert_mode_allocations}{keep_loaded}" # Changed from compute_device - current_hash = hashlib.sha256(settings_str.encode()).hexdigest() - - if not hasattr(cls, '_last_hash'): - cls._last_hash = current_hash - logger.mgpu_mm_log(f"IS_CHANGED first call: {current_hash[:8]}") - elif cls._last_hash != current_hash: - cls._last_hash = current_hash - logger.mgpu_mm_log(f"IS_CHANGED CHANGED: {current_hash[:8]} ← settings changed") - return current_hash - def override(self, *args, device=None, virtual_vram_gb=4.0, # Changed from compute_device - donor_device="cpu", expert_mode_allocations="", keep_loaded=True, **kwargs): - - unload_distorch_model = not keep_loaded - - from . import set_current_text_encoder_device # Use text encoder device setter - if device is not None: - set_current_text_encoder_device(device) - - # Register our patched ModelPatcher - register_patched_safetensor_modelpatcher() - - # Call original function - fn = getattr(super(), cls.FUNCTION) - - # Call the main function once - out = fn(*args, **kwargs) - - logger.mgpu_mm_log(f"[FLAG_SET_START] Setting '_mgpu_unload_distorch_model' to: {unload_distorch_model} (keep_loaded={keep_loaded})") - - # DIAGNOSTIC: Log full object chain at SET time - if hasattr(out[0], 'model'): - mp = out[0] # This is the ModelPatcher - mp_id = id(mp) - inner_model = getattr(mp, 'model', None) - inner_model_id = id(inner_model) if inner_model else None - inner_model_name = type(inner_model).__name__ if inner_model else "None" - - # Format inner_model_id properly for f-string - inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" - - logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] CLIP_NoDevice ModelPatcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") - - # FIX: Store flag on ModelPatcher itself - mp._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") - - # Also set on inner model for backwards compatibility - if inner_model: - inner_model._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") - - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - mp = out[0].patcher # This is the ModelPatcher - mp_id = id(mp) - inner_model = getattr(mp, 'model', None) - inner_model_id = id(inner_model) if inner_model else None - inner_model_name = type(inner_model).__name__ if inner_model else "None" - - # Format inner_model_id properly for f-string - inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" - - logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] CLIP_NoDevice ModelPatcher via patcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") - - # FIX: Store flag on ModelPatcher itself - mp._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") - - # Also set on inner model for backwards compatibility - if inner_model: - inner_model._mgpu_unload_distorch_model = unload_distorch_model - logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") - - vram_string = "" - if virtual_vram_gb > 0: - vram_string = f"{device};{virtual_vram_gb};{donor_device}" - elif expert_mode_allocations: - vram_string = device - - full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" - - logger.info(f"[MultiGPU DisTorch V2] Full allocation string: {full_allocation}") - - # Store allocation AFTER loading for next time - model_to_check = None - if hasattr(out[0], 'model'): - model_to_check = out[0] - elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): - model_to_check = out[0].patcher - - if model_to_check: - model_hash = create_safetensor_model_hash(model_to_check, "override_store") - settings_str = f"{device}{virtual_vram_gb}{donor_device}{expert_mode_allocations}" - settings_hash = hashlib.sha256(settings_str.encode()).hexdigest() - - # Store allocation for next time - safetensor_allocation_store[model_hash] = full_allocation - safetensor_settings_store[model_hash] = settings_hash - - if unload_distorch_model: - logger.mgpu_mm_log("[FLAG_TRIGGER] unload_distorch_model=True, triggering full system cleanup") - force_full_system_cleanup(reason="policy_every_load", force=True) - - return out - - return NodeOverrideDisTorchSafetensorV2ClipNoDevice +# NOTE: All wrapper functions have been moved to wrappers.py for better organization. +# This file (distorch_2.py) now contains ONLY backend logic: +# - register_patched_safetensor_modelpatcher() +# - analyze_safetensor_loading() and analyze_safetensor_loading_clip() +# - calculate_safetensor_vvram_allocation() +# - Allocation stores and model hash functions diff --git a/nodes.py b/nodes.py index 447a02c..fac43f8 100644 --- a/nodes.py +++ b/nodes.py @@ -551,32 +551,4 @@ class UNetLoaderLP: elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): out[0].patcher.model._distorch_high_precision_loras = False - return out - - -class FullCleanupMultiGPU: - @classmethod - def INPUT_TYPES(s): - return { - "required": { - "image": ("IMAGE",), - "reason": ("STRING", {"default": "inline_node", "multiline": False}), - }, - "optional": { - "force": ("BOOLEAN", {"default": True}), - } - } - - RETURN_TYPES = ("IMAGE",) - RETURN_NAMES = ("image",) - FUNCTION = "cleanup" - CATEGORY = "multigpu/maintenance" - TITLE = "Full System Cleanup (MultiGPU)" - - def cleanup(self, image, reason, force=True): - """ - Trigger the full system cleanup to match ComfyUI's 'Free model and node cache'. - Passthroughs the input image unchanged; summary is logged via MultiGPU logger. - """ - _ = force_full_system_cleanup(reason=reason, force=force) - return (image,) + return out \ No newline at end of file diff --git a/wrappers.py b/wrappers.py new file mode 100644 index 0000000..0da36ff --- /dev/null +++ b/wrappers.py @@ -0,0 +1,531 @@ +""" +ComfyUI-MultiGPU Wrapper Functions +All node override/wrapper generation functions consolidated in one location +""" + +import copy +import hashlib +import logging +from .device_utils import get_device_list + +logger = logging.getLogger("MultiGPU") + + +# ============================================================================ +# DISTORCH V2 SAFETENSOR WRAPPERS (DisTorch2 for .safetensors and .gguf) +# ============================================================================ + +def _create_distorch_safetensor_v2_override(cls, device_param_name, device_setter_func, apply_device_kwarg_workaround): + """ + Internal factory function - creates DisTorch 2.0 override class with parameterized behavior. + + Args: + cls: The base class to override + device_param_name: Parameter name ("compute_device" or "device") + device_setter_func: Function to call for device setting + apply_device_kwarg_workaround: If True, sets kwargs['device'] = 'default' for ComfyUI compatibility + + Returns: + Override class with specified behavior + """ + from .distorch_2 import ( + register_patched_safetensor_modelpatcher, + safetensor_allocation_store, + safetensor_settings_store, + create_safetensor_model_hash + ) + from .model_management_mgpu import force_full_system_cleanup + + class NodeOverrideDisTorchSafetensorV2(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"][device_param_name] = (devices, {"default": default_device}) + inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 128.0, "step": 0.1}) + inputs["optional"]["donor_device"] = (devices, {"default": "cpu"}) + inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) + inputs["optional"]["keep_loaded"] = ("BOOLEAN", {"default": True}) + return inputs + + CATEGORY = "multigpu/distorch_2" + FUNCTION = "override" + TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch2)" + + @classmethod + def IS_CHANGED(s, *args, virtual_vram_gb=4.0, donor_device="cpu", + expert_mode_allocations="", keep_loaded=True, **kwargs): + device_value = kwargs.get(device_param_name) + settings_str = f"{device_value}{virtual_vram_gb}{donor_device}{expert_mode_allocations}{keep_loaded}" + current_hash = hashlib.sha256(settings_str.encode()).hexdigest() + + if not hasattr(cls, '_last_hash'): + cls._last_hash = current_hash + logger.mgpu_mm_log(f"IS_CHANGED first call: {current_hash[:8]}") + elif cls._last_hash != current_hash: + cls._last_hash = current_hash + logger.mgpu_mm_log(f"IS_CHANGED CHANGED: {current_hash[:8]} ← settings changed") + return current_hash + + def override(self, *args, virtual_vram_gb=4.0, donor_device="cpu", + expert_mode_allocations="", keep_loaded=True, **kwargs): + + device_value = kwargs.get(device_param_name) + unload_distorch_model = not keep_loaded + + if device_value is not None: + device_setter_func(device_value) + + # Strip MultiGPU-specific parameters before calling original function + clean_kwargs = {k: v for k, v in kwargs.items() + if k not in [device_param_name, 'virtual_vram_gb', + 'donor_device', 'expert_mode_allocations', + 'keep_loaded']} + + if apply_device_kwarg_workaround: + clean_kwargs['device'] = 'default' + + register_patched_safetensor_modelpatcher() + + vram_string = "" + if virtual_vram_gb > 0: + vram_string = f"{device_value};{virtual_vram_gb};{donor_device}" + elif expert_mode_allocations: + vram_string = device_value + + full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" + + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **clean_kwargs) + + model_to_check = None + if hasattr(out[0], 'model'): + model_to_check = out[0] + elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): + model_to_check = out[0].patcher + + if model_to_check: + model_hash = create_safetensor_model_hash(model_to_check, "override_store") + settings_str = f"{device_value}{virtual_vram_gb}{donor_device}{expert_mode_allocations}" + settings_hash = hashlib.sha256(settings_str.encode()).hexdigest() + + safetensor_allocation_store[model_hash] = full_allocation + safetensor_settings_store[model_hash] = settings_hash + logger.debug(f"[MultiGPU DisTorch V2] Stored allocation for model {model_hash[:8]}: {full_allocation}") + + logger.info(f"[MultiGPU DisTorch V2] Full allocation string: {full_allocation}") + logger.mgpu_mm_log(f"[FLAG_SET_START] Setting '_mgpu_unload_distorch_model' to: {unload_distorch_model} (keep_loaded={keep_loaded})") + + if hasattr(out[0], 'model'): + mp = out[0] + mp_id = id(mp) + inner_model = getattr(mp, 'model', None) + inner_model_id = id(inner_model) if inner_model else None + inner_model_name = type(inner_model).__name__ if inner_model else "None" + inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" + + logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] ModelPatcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") + + mp._mgpu_unload_distorch_model = unload_distorch_model + logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") + + if inner_model: + inner_model._mgpu_unload_distorch_model = unload_distorch_model + logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") + + elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): + mp = out[0].patcher + mp_id = id(mp) + inner_model = getattr(mp, 'model', None) + inner_model_id = id(inner_model) if inner_model else None + inner_model_name = type(inner_model).__name__ if inner_model else "None" + inner_id_str = f"0x{inner_model_id:x}" if inner_model_id is not None else "None" + + logger.mgpu_mm_log(f"[OBJECT_CHAIN_SET] ModelPatcher via patcher: mp_id=0x{mp_id:x}, inner_model_id={inner_id_str}, inner_model_type={inner_model_name}") + + mp._mgpu_unload_distorch_model = unload_distorch_model + logger.mgpu_mm_log(f"[FLAG_SET_LOCATION] Set on ModelPatcher (mp_id=0x{mp_id:x}): mp._mgpu_unload_distorch_model = {unload_distorch_model}") + + if inner_model: + inner_model._mgpu_unload_distorch_model = unload_distorch_model + logger.mgpu_mm_log(f"[FLAG_SET_COMPAT] Also set on inner model (inner_model_id=0x{inner_model_id:x}) for compatibility") + + if unload_distorch_model: + logger.mgpu_mm_log("[FLAG_TRIGGER] unload_distorch_model=True, triggering full system cleanup") + force_full_system_cleanup(reason="policy_every_load", force=True) + + return out + + return NodeOverrideDisTorchSafetensorV2 + + +def override_class_with_distorch_safetensor_v2(cls): + """DisTorch 2.0 wrapper for safetensor UNet/VAE models""" + from . import set_current_device + return _create_distorch_safetensor_v2_override( + cls, + device_param_name="compute_device", + device_setter_func=set_current_device, + apply_device_kwarg_workaround=False + ) + + +def override_class_with_distorch_safetensor_v2_clip(cls): + """DisTorch 2.0 wrapper for safetensor CLIP models (with device kwarg workaround)""" + from . import set_current_text_encoder_device + return _create_distorch_safetensor_v2_override( + cls, + device_param_name="device", + device_setter_func=set_current_text_encoder_device, + apply_device_kwarg_workaround=True + ) + + +def override_class_with_distorch_safetensor_v2_clip_no_device(cls): + """DisTorch 2.0 wrapper for safetensor Triple/Quad CLIP models (no device kwarg workaround)""" + from . import set_current_text_encoder_device + return _create_distorch_safetensor_v2_override( + cls, + device_param_name="device", + device_setter_func=set_current_text_encoder_device, + apply_device_kwarg_workaround=False + ) + + +# ============================================================================ +# DISTORCH V1 LEGACY WRAPPERS (Rewritten to call V2 backend) +# ============================================================================ + +def override_class_with_distorch_gguf(cls): + """DisTorch V1 Legacy wrapper - maintains V1 UI but calls V2 backend""" + from . import set_current_device + from .distorch_2 import register_patched_safetensor_modelpatcher, safetensor_allocation_store, create_safetensor_model_hash + + class NodeOverrideDisTorchGGUFLegacy(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["device"] = (devices, {"default": default_device}) + inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 24.0, "step": 0.1}) + inputs["optional"]["use_other_vram"] = ("BOOLEAN", {"default": False}) + inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) + return inputs + + CATEGORY = "multigpu/legacy" + FUNCTION = "override" + TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (Legacy)" + + def override(self, *args, device=None, expert_mode_allocations="", use_other_vram=False, virtual_vram_gb=0.0, **kwargs): + if device is not None: + set_current_device(device) + + # Strip MultiGPU-specific parameters before calling original function + clean_kwargs = {k: v for k, v in kwargs.items() + if k not in ['device', 'virtual_vram_gb', 'use_other_vram', + 'expert_mode_allocations']} + + register_patched_safetensor_modelpatcher() + + vram_string = "" + if virtual_vram_gb > 0: + if use_other_vram: + available_devices = [d for d in get_device_list() if d != "cpu"] + other_devices = [d for d in available_devices if d != device] + other_devices.sort(key=lambda x: int(x.split(':')[1] if ':' in x else x[-1]), reverse=False) + device_string = ','.join(other_devices + ['cpu']) + vram_string = f"{device};{virtual_vram_gb};{device_string}" + else: + vram_string = f"{device};{virtual_vram_gb};cpu" + + full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" + + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **clean_kwargs) + + if hasattr(out[0], 'model'): + model_hash = create_safetensor_model_hash(out[0], "v1_compat") + safetensor_allocation_store[model_hash] = full_allocation + elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): + model_hash = create_safetensor_model_hash(out[0].patcher, "v1_compat") + safetensor_allocation_store[model_hash] = full_allocation + + return out + + return NodeOverrideDisTorchGGUFLegacy + + +def override_class_with_distorch_gguf_v2(cls): + """DisTorch V2 wrapper for GGUF models""" + from . import set_current_device + from .distorch_2 import register_patched_safetensor_modelpatcher, safetensor_allocation_store, create_safetensor_model_hash + + class NodeOverrideDisTorchGGUFv2(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + compute_device = devices[1] if len(devices) > 1 else devices[0] + + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["compute_device"] = (devices, {"default": compute_device}) + inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 128.0, "step": 0.1}) + inputs["optional"]["donor_device"] = (devices, {"default": "cpu"}) + inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) + return inputs + + CATEGORY = "multigpu/distorch_2" + FUNCTION = "override" + TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch2)" + + def override(self, *args, compute_device=None, virtual_vram_gb=4.0, donor_device="cpu", expert_mode_allocations="", **kwargs): + if compute_device is not None: + set_current_device(compute_device) + + # Strip MultiGPU-specific parameters before calling original function + clean_kwargs = {k: v for k, v in kwargs.items() + if k not in ['compute_device', 'virtual_vram_gb', + 'donor_device', 'expert_mode_allocations']} + + register_patched_safetensor_modelpatcher() + + vram_string = "" + if virtual_vram_gb > 0: + vram_string = f"{compute_device};{virtual_vram_gb};{donor_device}" + elif expert_mode_allocations: + vram_string = compute_device + + full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" + + logger.info(f"[MultiGPU DisTorch V2] Full allocation string: {full_allocation}") + + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **clean_kwargs) + + if hasattr(out[0], 'model'): + model_hash = create_safetensor_model_hash(out[0], "v2_gguf") + safetensor_allocation_store[model_hash] = full_allocation + elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): + model_hash = create_safetensor_model_hash(out[0].patcher, "v2_gguf") + safetensor_allocation_store[model_hash] = full_allocation + + return out + + return NodeOverrideDisTorchGGUFv2 + + +def override_class_with_distorch_clip(cls): + """DisTorch V1 wrapper for CLIP models - calls V2 backend""" + from . import set_current_text_encoder_device + from .distorch_2 import register_patched_safetensor_modelpatcher, safetensor_allocation_store, create_safetensor_model_hash + + class NodeOverrideDisTorchClip(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["device"] = (devices, {"default": default_device}) + inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 24.0, "step": 0.1}) + inputs["optional"]["use_other_vram"] = ("BOOLEAN", {"default": False}) + inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) + return inputs + + CATEGORY = "multigpu" + FUNCTION = "override" + TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch)" + + def override(self, *args, device=None, expert_mode_allocations="", use_other_vram=False, virtual_vram_gb=0.0, **kwargs): + if device is not None: + set_current_text_encoder_device(device) + + # Strip MultiGPU-specific parameters before calling original function + clean_kwargs = {k: v for k, v in kwargs.items() + if k not in ['device', 'virtual_vram_gb', 'use_other_vram', + 'expert_mode_allocations']} + + register_patched_safetensor_modelpatcher() + + vram_string = "" + if virtual_vram_gb > 0: + if use_other_vram: + available_devices = [d for d in get_device_list() if d != "cpu"] + other_devices = [d for d in available_devices if d != device] + other_devices.sort(key=lambda x: int(x.split(':')[1] if ':' in x else x[-1]), reverse=False) + device_string = ','.join(other_devices + ['cpu']) + vram_string = f"{device};{virtual_vram_gb};{device_string}" + else: + vram_string = f"{device};{virtual_vram_gb};cpu" + + full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" + + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **clean_kwargs) + + if hasattr(out[0], 'model'): + model_hash = create_safetensor_model_hash(out[0], "v1_clip") + safetensor_allocation_store[model_hash] = full_allocation + elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): + model_hash = create_safetensor_model_hash(out[0].patcher, "v1_clip") + safetensor_allocation_store[model_hash] = full_allocation + + return out + + return NodeOverrideDisTorchClip + + +def override_class_with_distorch_clip_no_device(cls): + """DisTorch V1 wrapper for Triple/Quad CLIP models - calls V2 backend""" + from . import set_current_text_encoder_device + from .distorch_2 import register_patched_safetensor_modelpatcher, safetensor_allocation_store, create_safetensor_model_hash + + class NodeOverrideDisTorchClipNoDevice(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["device"] = (devices, {"default": default_device}) + inputs["optional"]["virtual_vram_gb"] = ("FLOAT", {"default": 4.0, "min": 0.0, "max": 24.0, "step": 0.1}) + inputs["optional"]["use_other_vram"] = ("BOOLEAN", {"default": False}) + inputs["optional"]["expert_mode_allocations"] = ("STRING", {"multiline": False, "default": ""}) + return inputs + + CATEGORY = "multigpu" + FUNCTION = "override" + TITLE = f"{cls.TITLE if hasattr(cls, 'TITLE') else cls.__name__} (DisTorch)" + + def override(self, *args, device=None, expert_mode_allocations="", use_other_vram=False, virtual_vram_gb=0.0, **kwargs): + if device is not None: + set_current_text_encoder_device(device) + + # Strip MultiGPU-specific parameters before calling original function + clean_kwargs = {k: v for k, v in kwargs.items() + if k not in ['device', 'virtual_vram_gb', 'use_other_vram', + 'expert_mode_allocations']} + + register_patched_safetensor_modelpatcher() + + vram_string = "" + if virtual_vram_gb > 0: + if use_other_vram: + available_devices = [d for d in get_device_list() if d != "cpu"] + other_devices = [d for d in available_devices if d != device] + other_devices.sort(key=lambda x: int(x.split(':')[1] if ':' in x else x[-1]), reverse=False) + device_string = ','.join(other_devices + ['cpu']) + vram_string = f"{device};{virtual_vram_gb};{device_string}" + else: + vram_string = f"{device};{virtual_vram_gb};cpu" + + full_allocation = f"{expert_mode_allocations}#{vram_string}" if expert_mode_allocations or vram_string else "" + + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **clean_kwargs) + + if hasattr(out[0], 'model'): + model_hash = create_safetensor_model_hash(out[0], "v1_clip_nodev") + safetensor_allocation_store[model_hash] = full_allocation + elif hasattr(out[0], 'patcher') and hasattr(out[0].patcher, 'model'): + model_hash = create_safetensor_model_hash(out[0].patcher, "v1_clip_nodev") + safetensor_allocation_store[model_hash] = full_allocation + + return out + + return NodeOverrideDisTorchClipNoDevice + + +# Backward compatibility alias +override_class_with_distorch = override_class_with_distorch_gguf + + +# ============================================================================ +# STANDARD MULTIGPU WRAPPERS (Device selection without DisTorch) +# ============================================================================ + +def override_class(cls): + """Standard MultiGPU device override for UNet/VAE models""" + from . import set_current_device + + class NodeOverride(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["device"] = (devices, {"default": default_device}) + return inputs + + CATEGORY = "multigpu" + FUNCTION = "override" + + def override(self, *args, device=None, **kwargs): + if device is not None: + set_current_device(device) + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **kwargs) + return out + + return NodeOverride + + +def override_class_clip(cls): + """Standard MultiGPU device override for CLIP models (with device kwarg workaround)""" + from . import set_current_text_encoder_device + + class NodeOverride(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["device"] = (devices, {"default": default_device}) + return inputs + + CATEGORY = "multigpu" + FUNCTION = "override" + + def override(self, *args, device=None, **kwargs): + if device is not None: + set_current_text_encoder_device(device) + kwargs['device'] = 'default' + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **kwargs) + return out + + return NodeOverride + + +def override_class_clip_no_device(cls): + """Standard MultiGPU device override for Triple/Quad CLIP models (no device kwarg workaround)""" + from . import set_current_text_encoder_device + + class NodeOverride(cls): + @classmethod + def INPUT_TYPES(s): + inputs = copy.deepcopy(cls.INPUT_TYPES()) + devices = get_device_list() + default_device = devices[1] if len(devices) > 1 else devices[0] + inputs["optional"] = inputs.get("optional", {}) + inputs["optional"]["device"] = (devices, {"default": default_device}) + return inputs + + CATEGORY = "multigpu" + FUNCTION = "override" + + def override(self, *args, device=None, **kwargs): + if device is not None: + set_current_text_encoder_device(device) + fn = getattr(super(), cls.FUNCTION) + out = fn(*args, **kwargs) + return out + + return NodeOverride