From ce6a4b95fecf5f15e2491fb0c4ae0f86dcd31c8b Mon Sep 17 00:00:00 2001 From: doubletwisted <118005273+doubletwisted@users.noreply.github.com> Date: Wed, 6 Aug 2025 12:37:58 +0200 Subject: [PATCH] Update README and plugin files - Shortened README.md to remove fluff and make it more concise - Updated ComfyUI.py plugin file with latest changes - Added deadline_api.py for API functionality - Updated __init__.py and deadline_submit.py --- README.md | 175 ++++++------------------- __init__.py | 32 ++++- deadline_api.py | 229 +++++++++++++++++++++++++++++++++ deadline_submit.py | 77 +++++++++-- plugins/ComfyUI/ComfyUI.py | 254 ++++++++++++++++++++++++++++++++++--- 5 files changed, 600 insertions(+), 167 deletions(-) create mode 100644 deadline_api.py diff --git a/README.md b/README.md index cf9cf0d..8c0d0c5 100644 --- a/README.md +++ b/README.md @@ -1,155 +1,58 @@ # ComfyUI Deadline Plugin -A comprehensive plugin for integrating ComfyUI workflows with Thinkbox Deadline render farm management. +Submit ComfyUI workflows to Thinkbox Deadline render farm. -## 🚀 Features +## Features -- **Seamless Integration**: Submit ComfyUI workflows directly to Deadline from within ComfyUI -- **Distributed Rendering**: Leverage your render farm to process ComfyUI workflows at scale -- **Real-time Progress Monitoring**: Track rendering progress through Deadline Monitor -- **Seed Variation Control**: Automatically vary seeds across tasks for batch rendering -- **Flexible Configuration**: Support for pools, groups, priorities, and custom output directories +- Submit ComfyUI workflows directly to Deadline +- Batch rendering with seed variation +- Real-time progress monitoring via Deadline Monitor +- Configurable pools, groups, and priorities -## 📦 Repository Structure +## Installation -This repository contains two main components: +### ComfyUI Manager (Recommended) +1. Open ComfyUI Manager → Install Custom Nodes +2. Search "ComfyUI Deadline Submission" → Install +3. Restart ComfyUI -``` -ComfyUI-Deadline-Plugin/ -├── deadline_submit.py # Main ComfyUI custom node -├── __init__.py # Node exports for ComfyUI -├── plugins/ComfyUI/ # Deadline plugin (manual install required) -│ ├── ComfyUI.py -│ └── ComfyUI.param -├── example_extra_model_paths.yaml # Example model paths configuration -├── requirements.txt # Python dependencies -├── pyproject.toml # Modern Python packaging -├── LICENSE # MIT License -└── README.md # This file +### Manual Installation +```bash +cd ComfyUI/custom_nodes +git clone https://github.com/YOUR_USERNAME/ComfyUI-Deadline-Plugin.git ``` -## 🔧 Installation +### Deadline Plugin Setup (Required) +Copy `plugins/ComfyUI/` to your Deadline Repository's `custom/plugins/` directory and restart Deadline services. -### Option 1: ComfyUI Manager (Recommended) +## Usage -1. **Install Custom Nodes** (Automatic): - - Open ComfyUI Manager - - Go to "Install Custom Nodes" - - Search for "ComfyUI Deadline Submission" - - Click Install - - Restart ComfyUI +1. Add "Submit to Deadline" node to your workflow +2. Configure job settings (name, priority, pool, etc.) +3. Execute workflow +4. Monitor progress in Deadline Monitor -2. **Install Deadline Plugin** (Manual - Required): - - Copy the `plugins/ComfyUI/` directory to your Deadline Repository's `custom/plugins/` directory - - Restart Deadline services or Deadline Monitor +### Key Settings -### Option 2: Manual Installation +- **batch_count**: Number of tasks (1-100) +- **change_seeds_per_task**: Randomize seeds for different outputs +- **priority**: Job priority (0-100) +- **pool/group**: Deadline worker assignment -1. **Clone Repository**: - ```bash - cd ComfyUI/custom_nodes - git clone https://github.com/YOUR_USERNAME/ComfyUI-Deadline-Plugin.git - ``` +## Configuration -2. **Install Deadline Plugin**: - - Copy `ComfyUI-Deadline-Plugin/plugins/ComfyUI/` to `[Deadline Repository]/custom/plugins/ComfyUI/` - - Restart Deadline services +### Model Paths (Optional) +For render farms with shared storage, copy `example_extra_model_paths.yaml` to your ComfyUI installation as `extra_model_paths.yaml` and update paths. -4. **Restart ComfyUI** +## How It Works -### Worker Machine Setup +1. Captures current ComfyUI workflow +2. Submits to Deadline with proper configuration +3. Workers execute workflow via ComfyUI API +4. Progress reported through Deadline Monitor -No additional Python dependencies required - the plugin uses only standard library modules. +## Requirements -## 🎯 Quick Start - -### Using ComfyUI Custom Nodes - -1. **Add Submit Node**: In ComfyUI, add a "Submit to Deadline" node to your workflow -2. **Configure Settings**: Set job name, priority, pool, group, etc. -3. **Execute Workflow**: Run your workflow normally -4. **Monitor Progress**: Check job status in Deadline Monitor - -### Direct Workflow Submission - -1. **Export Workflow**: Save your workflow as JSON from ComfyUI -2. **Submit via Deadline**: Use the ComfyUI plugin in Deadline Monitor -3. **Configure Job**: Set rendering parameters and submit - -## ⚙️ Configuration - -### Deadline Plugin Configuration - -Configure these settings in Deadline Monitor: - -- **ComfyUI Instalation Path**: Path toComfyUI_windows_portable folder - -### Model Paths Configuration (Optional) - -ComfyUI is very slow to read models from network. For render farms with shared model storage, you can configure ComfyUI to use centralized model paths: - -1. **Copy the example file**: `ComfyUI\custom_nodes\deadline_submission\example_extra_model_paths.yaml` to your ComfyUI installation -2. **Rename it**: `extra_model_paths.yaml` -3. **Edit paths**: Update the paths to match your network storage setup - -The example configuration shows a local-first, network-fallback setup: -- **Local path**: `C:/AI` (fast access on each worker) -- **Network path**: `X:/AI` (centralized storage, fallback) - -This ensures workers use local models when available, falling back to network storage when needed. - -## 📋 Usage Guide - -### Submit to Deadline Node - -**Required Inputs:** -- `workflow_file`: Override workflow file path (leave empty for auto-detection) -- `auto_detect_workflow`: Use current workflow (recommended: ON) -- `batch_count`: Number of tasks to create (1-100) -- `chunk_size`: Frames per task (1-16) -- `change_seeds_per_task`: Vary seeds across tasks for different outputs -- `priority`: Job priority (0-100) -- `pool`: Deadline pool to use -- `group`: Deadline group to use -- `job_name`: Name for the Deadline job -- `bypass`: Skip submission (for testing) -- `skip_local_execution`: Submit only vs. submit and run locally - -**Optional Inputs:** -- `output_directory`: Custom output directory for all workers -- `comment`: Job comment -- `department`: Department name - -**Outputs:** -- `job_id`: Deadline job ID for tracking - - -### Seed Variation Feature - -Control how seeds are handled across batch tasks: - -- **ON**: Each task gets randomized seeds → Different outputs -- **OFF**: All tasks use original seeds → Identical outputs - -Compatible with: -- KSampler nodes (seed parameter) -- RandomNoise nodes (noise_seed parameter) -- Any node with seed-like parameters - - -## 🏗️ Technical Details - -### How It Works - -1. **Workflow Capture**: Automatically captures current ComfyUI workflow -2. **Deadline Submission**: Creates Deadline job with proper configuration -3. **Worker Execution**: - - Starts ComfyUI server on worker - - Submits workflow via API - - Randomizes seeds for batch - - Monitors progress -4. **Progress Reporting**: Progress updates through Deadline Monitor - ---- - -**Note**: This plugin requires both ComfyUI and Thinkbox Deadline to be properly installed and configured. The ComfyUI custom nodes can be installed automatically via ComfyUI Manager, but the Deadline plugin must be manually copied to your Deadline Repository. \ No newline at end of file +- ComfyUI installation on worker machines +- Thinkbox Deadline +- No additional Python dependencies (uses standard library) \ No newline at end of file diff --git a/__init__.py b/__init__.py index 0061609..99f8941 100644 --- a/__init__.py +++ b/__init__.py @@ -4,8 +4,36 @@ ComfyUI Deadline Plugin A comprehensive plugin for integrating ComfyUI workflows with Thinkbox Deadline render farm management. """ +import logging + # Import the custom nodes -from .deadline_submit import NODE_CLASS_MAPPINGS, NODE_DISPLAY_NAME_MAPPINGS +from .deadline_submit import NODE_CLASS_MAPPINGS as SUBMIT_MAPPINGS, NODE_DISPLAY_NAME_MAPPINGS as SUBMIT_DISPLAY_MAPPINGS + +# Use the mappings from deadline_submit +NODE_CLASS_MAPPINGS = SUBMIT_MAPPINGS +NODE_DISPLAY_NAME_MAPPINGS = SUBMIT_DISPLAY_MAPPINGS + +# Setup web extensions +WEB_DIRECTORY = "./web" # Export for ComfyUI Manager -__all__ = ["NODE_CLASS_MAPPINGS", "NODE_DISPLAY_NAME_MAPPINGS"] \ No newline at end of file +__all__ = ["NODE_CLASS_MAPPINGS", "NODE_DISPLAY_NAME_MAPPINGS", "WEB_DIRECTORY"] + +# Initialize API when server is available +def init_server(server): + """Initialize Deadline API with ComfyUI server""" + try: + from .deadline_api import integrate_with_comfyui + integrate_with_comfyui(server) + logging.info("Deadline API initialized") + except Exception as e: + logging.error(f"Failed to initialize Deadline API: {e}") + +# ComfyUI will call this if it exists +try: + import server + if hasattr(server, "PromptServer") and server.PromptServer.instance: + init_server(server.PromptServer.instance) +except: + # Server not available yet, will be initialized later + pass \ No newline at end of file diff --git a/deadline_api.py b/deadline_api.py new file mode 100644 index 0000000..475b0c3 --- /dev/null +++ b/deadline_api.py @@ -0,0 +1,229 @@ +""" +Deadline API endpoints for ComfyUI +Provides REST and WebSocket endpoints for the Deadline panel +""" + +import json +import asyncio +from typing import Dict, List, Optional, Set +from aiohttp import web +import logging + +logger = logging.getLogger(__name__) + +class DeadlineAPIHandler: + """Handles API requests for Deadline integration""" + + def __init__(self): + self.workers: Dict[str, Dict] = {} + self.active_jobs: Dict[str, Dict] = {} + self.websocket_clients: Set[web.WebSocketResponse] = set() + self.update_lock = asyncio.Lock() + + async def get_workers(self, request: web.Request) -> web.Response: + """GET /deadline/workers - Return list of active workers""" + try: + response_data = { + "workers": list(self.workers.values()), + "activeWorkers": len([w for w in self.workers.values() if w.get("status") == "active"]), + "totalJobs": len(self.active_jobs) + } + return web.json_response(response_data) + except Exception as e: + logger.error(f"Error getting workers: {e}") + return web.json_response({"error": str(e)}, status=500) + + async def submit_job(self, request: web.Request) -> web.Response: + """POST /deadline/submit - Submit workflow to Deadline""" + try: + data = await request.json() + workflow = data.get("workflow") + is_distributed = data.get("isDistributed", False) + master_ws = data.get("masterWs", "localhost:8188") + + # Here you would integrate with the actual Deadline submission + # For now, return a mock response + job_id = f"job_{len(self.active_jobs) + 1:04d}" + + self.active_jobs[job_id] = { + "id": job_id, + "status": "submitted", + "isDistributed": is_distributed, + "masterWs": master_ws + } + + # Notify WebSocket clients + await self._broadcast({ + "type": "job_submitted", + "jobId": job_id + }) + + return web.json_response({"jobId": job_id, "status": "submitted"}) + except Exception as e: + logger.error(f"Error submitting job: {e}") + return web.json_response({"error": str(e)}, status=500) + + async def stop_worker(self, request: web.Request) -> web.Response: + """POST /deadline/workers/{workerId}/stop - Stop specific worker""" + try: + worker_id = request.match_info.get("workerId") + + if worker_id in self.workers: + self.workers[worker_id]["status"] = "stopping" + + # Here you would send actual stop command to Deadline + + # Remove worker after a delay + asyncio.create_task(self._remove_worker_delayed(worker_id)) + + await self._broadcast({ + "type": "worker_stopping", + "workerId": worker_id + }) + + return web.json_response({"status": "stopping"}) + else: + return web.json_response({"error": "Worker not found"}, status=404) + except Exception as e: + logger.error(f"Error stopping worker: {e}") + return web.json_response({"error": str(e)}, status=500) + + async def stop_all_workers(self, request: web.Request) -> web.Response: + """POST /deadline/workers/stop-all - Stop all workers""" + try: + for worker_id in list(self.workers.keys()): + self.workers[worker_id]["status"] = "stopping" + + # Here you would send actual stop commands to Deadline + + await self._broadcast({ + "type": "all_workers_stopping" + }) + + # Clear workers after a delay + asyncio.create_task(self._clear_workers_delayed()) + + return web.json_response({"status": "stopping all"}) + except Exception as e: + logger.error(f"Error stopping all workers: {e}") + return web.json_response({"error": str(e)}, status=500) + + async def websocket_handler(self, request: web.Request) -> web.WebSocketResponse: + """WebSocket endpoint for real-time updates""" + ws = web.WebSocketResponse() + await ws.prepare(request) + + # Add to clients set + self.websocket_clients.add(ws) + + try: + # Send initial state + await ws.send_json({ + "type": "initial_state", + "workers": list(self.workers.values()), + "activeWorkers": len([w for w in self.workers.values() if w.get("status") == "active"]), + "totalJobs": len(self.active_jobs) + }) + + # Keep connection alive + async for msg in ws: + if msg.type == web.WSMsgType.TEXT: + try: + data = json.loads(msg.data) + # Handle incoming messages if needed + if data.get("type") == "ping": + await ws.send_json({"type": "pong"}) + except json.JSONDecodeError: + pass + elif msg.type == web.WSMsgType.ERROR: + logger.error(f"WebSocket error: {ws.exception()}") + finally: + # Remove from clients set + self.websocket_clients.discard(ws) + + return ws + + async def register_worker(self, worker_info: Dict) -> None: + """Register a new worker""" + async with self.update_lock: + worker_id = worker_info.get("id") + self.workers[worker_id] = worker_info + + await self._broadcast({ + "type": "worker_registered", + "worker": worker_info + }) + + async def update_worker_status(self, worker_id: str, status: Dict) -> None: + """Update worker status""" + async with self.update_lock: + if worker_id in self.workers: + self.workers[worker_id].update(status) + + await self._broadcast({ + "type": "worker_update", + "workers": list(self.workers.values()) + }) + + async def unregister_worker(self, worker_id: str) -> None: + """Unregister a worker""" + async with self.update_lock: + if worker_id in self.workers: + del self.workers[worker_id] + + await self._broadcast({ + "type": "worker_unregistered", + "workerId": worker_id + }) + + async def _broadcast(self, message: Dict) -> None: + """Broadcast message to all WebSocket clients""" + if not self.websocket_clients: + return + + # Create tasks for sending to all clients + tasks = [] + for ws in list(self.websocket_clients): + if not ws.closed: + tasks.append(ws.send_json(message)) + + # Send to all clients concurrently + if tasks: + await asyncio.gather(*tasks, return_exceptions=True) + + async def _remove_worker_delayed(self, worker_id: str, delay: float = 2.0) -> None: + """Remove worker after a delay""" + await asyncio.sleep(delay) + await self.unregister_worker(worker_id) + + async def _clear_workers_delayed(self, delay: float = 2.0) -> None: + """Clear all workers after a delay""" + await asyncio.sleep(delay) + async with self.update_lock: + self.workers.clear() + await self._broadcast({ + "type": "workers_cleared" + }) + + +# Global handler instance +deadline_api = DeadlineAPIHandler() + + +def setup_routes(app: web.Application) -> None: + """Setup routes for Deadline API""" + app.router.add_get("/deadline/workers", deadline_api.get_workers) + app.router.add_post("/deadline/submit", deadline_api.submit_job) + app.router.add_post("/deadline/workers/{workerId}/stop", deadline_api.stop_worker) + app.router.add_post("/deadline/workers/stop-all", deadline_api.stop_all_workers) + app.router.add_get("/deadline", deadline_api.websocket_handler) + + +# Integration with ComfyUI +def integrate_with_comfyui(server): + """Integrate Deadline API with ComfyUI server""" + if hasattr(server, "app"): + setup_routes(server.app) + logger.info("Deadline API routes registered") + else: + logger.warning("Could not register Deadline API routes - server.app not found") \ No newline at end of file diff --git a/deadline_submit.py b/deadline_submit.py index 83bf60a..ad50b37 100644 --- a/deadline_submit.py +++ b/deadline_submit.py @@ -354,11 +354,8 @@ class DeadlineJobSubmitter: f.write("DefaultCudaDeviceZero=True\n") - # Map boolean to appropriate SeedMode value - if config.get('change_seeds_per_task', True): - f.write("SeedMode=change\n") - else: - f.write("SeedMode=fixed\n") + # Note: Seed handling is now managed by DeadlineSeed nodes + f.write("SeedMode=fixed\n") if config['batch_count'] > 1: f.write("BatchMode=True\n") @@ -427,6 +424,60 @@ class ExecutionInterruptor: return True +class DeadlineSeed: + """ + Distributes seed values across Deadline tasks. + On first task: passes through the original seed. + On subsequent tasks: adds offset based on task ID. + """ + + @classmethod + def INPUT_TYPES(cls): + return { + "required": { + "seed": ("INT", { + "default": 1125899906842, + "min": 0, + "max": 1125899906842624, + "forceInput": False # Widget by default, can be converted to input + }), + }, + "hidden": { + "task_id": ("INT", {"default": 0}), + "batch_mode": ("BOOLEAN", {"default": False}), + }, + } + + RETURN_TYPES = ("INT",) + RETURN_NAMES = ("seed",) + FUNCTION = "distribute" + CATEGORY = "deadline" + + def distribute(self, seed, task_id=0, batch_mode=False): + """ + Distribute seeds across Deadline tasks. + + Args: + seed: Base seed value + task_id: Current task ID (injected by Deadline) + batch_mode: Whether this is running in batch mode + """ + # Ensure task_id is an integer + try: + task_id = int(task_id) + except (ValueError, TypeError): + task_id = 0 + + if not batch_mode or task_id == 0: + # First task or not in batch mode: pass through original seed + print(f"Deadline Seed: Task {task_id} using original seed {seed}") + return (seed,) + else: + # Subsequent tasks: add offset based on task ID + new_seed = seed + task_id + print(f"Deadline Seed: Task {task_id} using modified seed {new_seed} (original: {seed})") + return (new_seed,) + # Node implementation class DeadlineSubmitNode: """Submit the current ComfyUI workflow to Thinkbox Deadline""" @@ -460,11 +511,7 @@ class DeadlineSubmitNode: "max": NodeDefaults.MAX_CHUNK_SIZE, "step": 1 }), - "change_seeds_per_task": ("BOOLEAN", { - "default": True, - "label_on": "Vary seeds across tasks", - "label_off": "Keep original seeds" - }), + "priority": ("INT", { "default": NodeDefaults.PRIORITY, "min": 0, @@ -488,6 +535,7 @@ class DeadlineSubmitNode: }), "comment": ("STRING", {"default": ""}), "department": ("STRING", {"default": ""}), + }, "hidden": { "prompt": "PROMPT", @@ -529,7 +577,7 @@ class DeadlineSubmitNode: return [NodeDefaults.GROUP] def submit_to_deadline(self, workflow_file, auto_detect_workflow, batch_count, chunk_size, - change_seeds_per_task, priority, pool, group, job_name, bypass, + priority, pool, group, job_name, bypass, skip_local_execution=True, output_directory="", comment="", department="", prompt=None, extra_pnginfo=None): """Submit the workflow to Deadline for rendering""" @@ -549,7 +597,7 @@ class DeadlineSubmitNode: # Create job configuration job_config = self._create_job_config( job_name, priority, pool, group, batch_count, chunk_size, - change_seeds_per_task, output_directory, comment, department + output_directory, comment, department ) # Submit to Deadline @@ -591,7 +639,7 @@ class DeadlineSubmitNode: raise Exception(f"Could not read workflow file: {str(e)}") def _create_job_config(self, job_name: str, priority: int, pool: str, group: str, - batch_count: int, chunk_size: int, change_seeds_per_task: bool, + batch_count: int, chunk_size: int, output_directory: str, comment: str, department: str) -> Dict: """Create job configuration dictionary""" return { @@ -601,7 +649,6 @@ class DeadlineSubmitNode: 'group': group, 'batch_count': batch_count, 'chunk_size': chunk_size, - 'change_seeds_per_task': change_seeds_per_task, 'output_directory': output_directory, 'comment': comment, 'department': department @@ -610,9 +657,11 @@ class DeadlineSubmitNode: # Register the nodes NODE_CLASS_MAPPINGS = { "DeadlineSubmit": DeadlineSubmitNode, + "DeadlineSeed": DeadlineSeed, } # Add display names for the nodes NODE_DISPLAY_NAME_MAPPINGS = { "DeadlineSubmit": "Submit to Deadline", + "DeadlineSeed": "Deadline Seed", } \ No newline at end of file diff --git a/plugins/ComfyUI/ComfyUI.py b/plugins/ComfyUI/ComfyUI.py index de51124..2843afd 100644 --- a/plugins/ComfyUI/ComfyUI.py +++ b/plugins/ComfyUI/ComfyUI.py @@ -14,6 +14,8 @@ import urllib.error import urllib.parse import traceback import random +import platform +from typing import Tuple """ ComfyUI Deadline Plugin @@ -40,6 +42,27 @@ SEED_PARAMETER_NAMES = ["seed", "noise_seed", "value"] # Output node types that indicate the workflow will produce output OUTPUT_NODE_TYPES = ["SaveImage", "PreviewImage", "SaveVideo"] +def get_distributed_config_for_plugin(plugin) -> Tuple[bool, bool, bool]: + """Get distributed configuration with plugin info priority, fallback to environment""" + # Priority 1: Plugin info entries (preferred) + worker_mode = plugin.GetBooleanPluginInfoEntryWithDefault("WorkerMode", False) + distributed_mode = plugin.GetBooleanPluginInfoEntryWithDefault("DistributedMode", False) + force_new_instance = plugin.GetBooleanPluginInfoEntryWithDefault("ForceNewInstance", False) + + # Priority 2: Environment variables (fallback for backwards compatibility) + if not worker_mode and not distributed_mode and not force_new_instance: + worker_mode = os.environ.get('COMFY_WORKER_MODE', '0').lower() in ('1', 'true', 'yes') + distributed_mode = os.environ.get('DEADLINE_DIST_MODE', '0').lower() in ('1', 'true', 'yes') + force_new_instance = os.environ.get('COMFY_FORCE_NEW_INSTANCE', '0').lower() in ('1', 'true', 'yes') + + if worker_mode or distributed_mode or force_new_instance: + plugin.LogWarning("Using environment variables for distributed config. Consider updating to plugin info entries.") + + # Log the configuration + plugin.LogInfo(f"Distributed config - WorkerMode: {worker_mode}, DistributedMode: {distributed_mode}, ForceNewInstance: {force_new_instance}") + + return worker_mode, distributed_mode, force_new_instance + def GetDeadlinePlugin(): return ComfyUI() @@ -71,6 +94,8 @@ class ComfyUI(DeadlinePlugin): """Setup stdout handlers for ComfyUI output parsing""" # Server startup handlers self.AddStdoutHandlerCallback(".*Starting server.*").HandleCallback += self.HandleServerStarted + # Also catch the GUI message as backup + self.AddStdoutHandlerCallback(".*To see the GUI go to.*").HandleCallback += self.HandleServerStarted self.AddStdoutHandlerCallback(".*Error:.*").HandleCallback += self.HandleStdoutError self.AddStdoutHandlerCallback(".*Exception:.*").HandleCallback += self.HandleStdoutError @@ -265,6 +290,26 @@ class ComfyUI(DeadlinePlugin): def _determine_final_port(self, base_port: int) -> str: """Determine final port to use, checking for existing instances""" + # Check if we should force a new instance (for distributed workers) + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + if force_new_instance or worker_mode or distributed_mode: + self.LogInfo("Worker/Distributed mode: Will start new ComfyUI instance (not reusing existing)") + self.use_existing_comfyui = False + + # For workers, use dynamic port allocation to avoid conflicts + if worker_mode or distributed_mode: + worker_port = self._calculate_worker_port(base_port) + self.comfyui_port = self._find_available_port(worker_port) + self.LogInfo(f"Worker mode: Using port {self.comfyui_port}") + else: + self.comfyui_port = self._find_available_port(base_port) + self.LogInfo(f"Force new instance: Using port {self.comfyui_port}") + + self.comfyui_api_url = f"http://127.0.0.1:{self.comfyui_port}" + return self.comfyui_port + + # Normal batch mode logic self.LogInfo(f"Checking if ComfyUI is already running on port {base_port}") if self._is_port_in_use(base_port): @@ -282,6 +327,25 @@ class ComfyUI(DeadlinePlugin): return self.comfyui_port + def _calculate_worker_port(self, base_port: int) -> int: + """Calculate unique worker port based on task ID to avoid conflicts""" + try: + # Get task ID from Deadline environment + task_id = int(os.environ.get('DEADLINE_TASK_ID', '1')) + + # Calculate unique worker port: base_port + 100 + task_id + # This ensures workers get ports like 8289, 8290, 8291, etc. + # For single PC testing, this allows multiple workers on one GPU + worker_port = base_port + 100 + task_id + + self.LogInfo(f"Calculated worker port: {worker_port} (base: {base_port}, task: {task_id})") + return worker_port + except ValueError: + # Fallback if task ID is not a valid integer + fallback_port = base_port + 100 + self.LogInfo(f"Could not parse task ID, using fallback port: {fallback_port}") + return fallback_port + def PreRenderTasks(self): """Setup tasks before rendering""" self.LogInfo("ComfyUI PreRenderTasks started.") @@ -441,6 +505,18 @@ print('Dummy command timeout reached or task completed.') cuda_arg = self._get_cuda_device_arg() if cuda_arg: args_list.append(cuda_arg) + + # Add ComfyUI flags for worker mode + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + if worker_mode or distributed_mode: + args_list.append("--listen") # Allow external connections + args_list.append("--enable-cors-header") # Enable CORS for API access + ##args_list.append("--dont-print-server") # Reduce startup output for workers + + # Always add windows standalone build flag if on Windows + if platform.system().lower() == 'windows': + args_list.append("--windows-standalone-build") # Add output directory if custom one was specified if self.custom_output_dir_specified and self.comfyui_output_dir: @@ -455,14 +531,30 @@ print('Dummy command timeout reached or task completed.') def HandleServerStarted(self): """Called when the ComfyUI server has started""" + # Prevent multiple triggers from stdout handlers + if self.workflow_submitted: + self.LogInfo("ComfyUI server startup detected, but workflow already submitted - ignoring") + return + self.LogInfo("ComfyUI server has started") self.server_started = True + # Check if we're in worker/distributed mode + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + self.LogInfo(f"Distributed config check - WorkerMode: {worker_mode}, DistributedMode: {distributed_mode}") + self.LogInfo(f"use_existing_comfyui={self.use_existing_comfyui}, workflow_submitted={self.workflow_submitted}") + # Start workflow submission if not using existing instance if not self.use_existing_comfyui and not self.workflow_submitted: + # Mark as submitted IMMEDIATELY to prevent race conditions + self.workflow_submitted = True + self.LogInfo("Starting workflow submission thread...") workflow_thread = threading.Thread(target=self.submit_workflow) workflow_thread.daemon = True workflow_thread.start() + else: + self.LogInfo("Skipping workflow submission - using existing instance or already submitted") def http_request(self, url: str, method: str = "GET", data=None, headers=None, verbose: bool = True) -> dict: """Make an HTTP request to the ComfyUI API""" @@ -578,6 +670,42 @@ print('Dummy command timeout reached or task completed.') else: return random.randint(0, MAX_SEED_VALUE) + def inject_deadline_seed_parameters(self, workflow_data: dict) -> bool: + """ + Inject task_id and batch_mode into DeadlineDistributedSeed nodes. + + Args: + workflow_data: The workflow data + + Returns: + bool: True if any nodes were modified, False otherwise + """ + try: + task_id = int(self.GetCurrentTaskId()) + batch_mode = self.batch_mode + nodes_modified = False + + for node_id, node in workflow_data.items(): + if not isinstance(node, dict): + continue + + if node.get("class_type") == "DeadlineSeed": + if "inputs" not in node: + node["inputs"] = {} + + # Inject task_id and batch_mode as hidden parameters + node["inputs"]["task_id"] = task_id + node["inputs"]["batch_mode"] = batch_mode + + self.LogInfo(f"Injected task_id={task_id}, batch_mode={batch_mode} into DeadlineSeed node {node_id}") + nodes_modified = True + + return nodes_modified + + except Exception as e: + self.LogWarning(f"Error injecting deadline seed parameters: {e}") + return False + def load_and_validate_workflow(self) -> dict: """Load workflow file and validate its structure""" workflow_file = self._get_workflow_file_path() @@ -590,14 +718,20 @@ print('Dummy command timeout reached or task completed.') workflow_data = self._load_workflow_from_file(workflow_file) workflow_data = self.validate_workflow(workflow_data) - # Apply seed manipulation - task_id = self.GetCurrentTaskId() - seeds_modified = self.modify_workflow_seeds(workflow_data, task_id) + # Inject parameters for DeadlineDistributedSeed nodes + deadline_seeds_injected = self.inject_deadline_seed_parameters(workflow_data) - if seeds_modified: - self.LogInfo(f"Applied seed manipulation for task ID {task_id}") + # Apply seed manipulation only if no DeadlineDistributedSeed nodes are present + if not deadline_seeds_injected: + task_id = self.GetCurrentTaskId() + seeds_modified = self.modify_workflow_seeds(workflow_data, task_id) + + if seeds_modified: + self.LogInfo(f"Applied seed manipulation for task ID {task_id}") + else: + self.LogInfo(f"No seed manipulation applied for task ID {task_id}") else: - self.LogInfo(f"No seed manipulation applied for task ID {task_id}") + self.LogInfo("DeadlineSeed nodes detected - skipping automatic seed modification") return workflow_data except Exception as e: @@ -607,7 +741,19 @@ print('Dummy command timeout reached or task completed.') def _get_workflow_file_path(self) -> str: """Get the workflow file path from plugin settings""" - workflow_file = self.GetPluginInfoEntryWithDefault("ComfyWorkflowFile", self.GetDataFilename()) + # Check if we're in distributed/worker mode + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + if distributed_mode or worker_mode: + # For distributed workers, use the WorkflowFile from plugin info (should be dummy workflow) + workflow_file = self.GetPluginInfoEntryWithDefault("WorkflowFile", "") + if not workflow_file: + # Fallback to ComfyWorkflowFile if WorkflowFile is not set + workflow_file = self.GetPluginInfoEntryWithDefault("ComfyWorkflowFile", self.GetDataFilename()) + else: + # Normal batch mode + workflow_file = self.GetPluginInfoEntryWithDefault("ComfyWorkflowFile", self.GetDataFilename()) + workflow_file = RepositoryUtils.CheckPathMapping(workflow_file) self.LogInfo(f"Workflow file setting from plugin info: '{workflow_file}'") return workflow_file @@ -834,10 +980,28 @@ print('Dummy command timeout reached or task completed.') for i in range(1, self.chunk_size): prompt_workflow = copy.deepcopy(workflow_data) - # Modify seeds if needed - if self.GetPluginInfoEntryWithDefault("SeedMode", "fixed") != "fixed": - self.modify_workflow_seeds(prompt_workflow, i) - self.LogInfo(f"Modified seeds for additional prompt {i}") + # Check if workflow has DeadlineSeed nodes + has_deadline_seeds = any( + node.get("class_type") == "DeadlineSeed" + for node in prompt_workflow.values() + if isinstance(node, dict) + ) + + if has_deadline_seeds: + # Update task_id for DeadlineSeed nodes (chunk-local indexing) + for node_id, node in prompt_workflow.items(): + if isinstance(node, dict) and node.get("class_type") == "DeadlineSeed": + if "inputs" not in node: + node["inputs"] = {} + # Use i as the offset for chunks within the same task + base_task_id = int(node["inputs"].get("task_id", 0)) + node["inputs"]["task_id"] = base_task_id + i + self.LogInfo(f"Updated DeadlineSeed node {node_id} task_id to {base_task_id + i}") + else: + # Modify seeds using the old method + if self.GetPluginInfoEntryWithDefault("SeedMode", "fixed") != "fixed": + self.modify_workflow_seeds(prompt_workflow, i) + self.LogInfo(f"Modified seeds for additional prompt {i}") # Queue the workflow data = {"prompt": prompt_workflow, "client_id": self.client_id} @@ -910,8 +1074,16 @@ print('Dummy command timeout reached or task completed.') self.SetProgress(100) self.SetStatusMessage("Finished Render") self.task_completed = True - self.signal_task_completion() - self.LogInfo(f"All {self.chunk_size} prompts in chunk completed, task marked as complete") + + # Check if we're in distributed worker mode + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + if worker_mode and distributed_mode: + self.LogInfo("Distributed worker mode: Registration completed, entering keep-alive mode") + self._enter_distributed_keep_alive_mode() + else: + self.signal_task_completion() + self.LogInfo(f"All {self.chunk_size} prompts in chunk completed, task marked as complete") def _move_to_next_prompt(self): """Move to tracking the next prompt""" @@ -1013,7 +1185,15 @@ print('Dummy command timeout reached or task completed.') while time.time() - start_time < DEFAULT_TIMEOUT and self.thread_running: if self.task_completed: self.LogInfo("Task already marked as complete by stdout handler") - self.signal_task_completion() + + # Check if we're in distributed worker mode + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + if worker_mode and distributed_mode: + self.LogInfo("Distributed worker mode: Registration completed, entering keep-alive mode") + self._enter_distributed_keep_alive_mode() + else: + self.signal_task_completion() return True if self.prompt_id: @@ -1032,10 +1212,54 @@ print('Dummy command timeout reached or task completed.') return False if self.task_completed: - self.signal_task_completion() + # Check if we're in distributed worker mode + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + + if worker_mode and distributed_mode: + self.LogInfo("Distributed worker mode: Registration workflow completed, entering keep-alive mode") + self.LogInfo("Task will remain active to process distributed workflows from master") + + # Don't signal completion - enter keep-alive mode instead + self._enter_distributed_keep_alive_mode() + else: + # Normal mode - complete the task + self.signal_task_completion() return self.task_completed + def _enter_distributed_keep_alive_mode(self): + """Enter keep-alive mode for distributed workers""" + import time + import threading + + self.LogInfo("🔄 Entering distributed worker keep-alive mode...") + self.LogInfo("Worker will remain active until manually stopped or job is cancelled") + + def keep_alive_loop(): + """Keep the task alive indefinitely""" + try: + while True: + self.LogInfo("🔄 Distributed worker is alive and ready for workflows...") + time.sleep(300) # Log every 5 minutes + except KeyboardInterrupt: + self.LogInfo("🛑 Keep-alive interrupted by user") + except Exception as e: + self.LogInfo(f"❌ Keep-alive error: {e}") + + # Start keep-alive in daemon thread + keep_alive_thread = threading.Thread(target=keep_alive_loop, daemon=True) + keep_alive_thread.start() + + try: + # Block main thread indefinitely + self.LogInfo("🔄 Main thread entering infinite wait...") + while True: + time.sleep(60) # Check every minute + except KeyboardInterrupt: + self.LogInfo("🛑 Distributed worker keep-alive interrupted") + except Exception as e: + self.LogInfo(f"❌ Distributed worker keep-alive error: {e}") + def _poll_prompt_status(self, poll_count: int) -> bool: """Poll the status of the current prompt""" try: