From 847fc78cf8fb3da6e57ce30b136276f5a036f665 Mon Sep 17 00:00:00 2001 From: Dominik Bargiel Date: Mon, 15 Jun 2026 19:17:09 +0200 Subject: [PATCH] Improve Deadline submission workflow --- README.md | 115 ++- __init__.py | 42 +- deadline_api.py | 229 ----- deadline_submit.py | 1008 +++++++++++----------- plugins/ComfyUI/ComfyUI.py | 398 ++++++--- pyproject.toml | 2 +- scripts/copy_comfy_to_network.ps1 | 41 +- scripts/deploy_deadline_plugin.ps1 | 181 ++++ scripts/maintenance/ComfyModelsSync.py | 440 ++++++++++ scripts/maintenance/ComfyUISync.py | 176 ++++ scripts/maintenance/README.md | 39 + scripts/maintenance/submit_comfy_sync.py | 348 ++++++++ 12 files changed, 2054 insertions(+), 965 deletions(-) delete mode 100644 deadline_api.py create mode 100644 scripts/deploy_deadline_plugin.ps1 create mode 100644 scripts/maintenance/ComfyModelsSync.py create mode 100644 scripts/maintenance/ComfyUISync.py create mode 100644 scripts/maintenance/README.md create mode 100644 scripts/maintenance/submit_comfy_sync.py diff --git a/README.md b/README.md index 46e4ecf..a0106f0 100644 --- a/README.md +++ b/README.md @@ -1,73 +1,96 @@ # ComfyUI Deadline Plugin -Submit ComfyUI workflows to Thinkbox Deadline render farm. - -Check out a quick demo showing how ComfyUI can play nice with Deadline. -Please check out the video: - -

ComfyUI x Deadline: Leverage Your GPUs

+Submit ComfyUI workflows to Thinkbox Deadline. ## Features -- Submit ComfyUI workflows directly to Deadline -- Batch rendering with seed variation -- Real-time progress monitoring via Deadline Monitor -- Configurable pools, groups, and priorities +- Submit the current ComfyUI API prompt directly to Deadline +- Submit-only local execution: normal output nodes are not rendered on the submitter +- Deadline variation jobs using `batch_count` and `chunk_size` +- Deterministic seed variation through the `DeadlineSeed` node +- Stages referenced default ComfyUI input assets beside the output directory +- Stores the normal ComfyUI `workflow.json` with node placement in the Deadline job files +- Launches isolated portable Windows ComfyUI worker instances +- Preserves compatibility flags used by `ComfyUI-Deadline-Distributed` ## Installation -### ComfyUI Manager (Recommended) -1. Open ComfyUI Manager → Install Custom Nodes -2. Search "ComfyUI Deadline Submission" → Install -3. Restart ComfyUI +### ComfyUI -### Manual Installation ```bash cd ComfyUI/custom_nodes -git clone https://github.com/YOUR_USERNAME/ComfyUI-Deadline-Plugin.git +git clone https://github.com/doubletwisted/ComfyUI-Deadline-Plugin.git ``` -### Deadline Plugin Setup (Required) -Copy `plugins/ComfyUI/` to your Deadline Repository's `custom/plugins/` directory and restart Deadline services. +Restart ComfyUI after installing or updating. + +### Deadline + +Deploy `plugins/ComfyUI/` into your Deadline Repository `custom/plugins/` directory, then restart Deadline services or reload the repository plugin. + +For render-farm maintenance, this repo includes publishable templates in `scripts/maintenance`. Copy those scripts to a shared path reachable by Workers, update their default paths or pass overrides, then submit them as Deadline maintenance jobs: + +```powershell +python \\YOUR-SERVER\share\scripts\maintenance\submit_comfy_sync.py --type both +``` + +The maintenance submitter can submit the ComfyUI install sync, the model sync, or both: + +```powershell +python \\YOUR-SERVER\share\scripts\maintenance\submit_comfy_sync.py --type installation +python \\YOUR-SERVER\share\scripts\maintenance\submit_comfy_sync.py --type models +``` + +For Deadline repository plugin deploys, this repo includes a direct deploy helper: + +```powershell +powershell.exe -ExecutionPolicy Bypass -File .\scripts\deploy_deadline_plugin.ps1 +``` + +The deploy script discovers the repository with `deadlinecommand -GetRepositoryPath`. You can also pass either the repository root or the custom plugins folder explicitly: + +```powershell +powershell.exe -ExecutionPolicy Bypass -File .\scripts\deploy_deadline_plugin.ps1 -RepositoryPath "\\YOUR-SERVER\Repository\custom\plugins" +``` + +In Deadline Monitor, configure the ComfyUI plugin `ComfyUI Installation Paths` setting. Use one portable Windows ComfyUI root per line: + +```text +C:\ComfyUI_windows_portable +D:\Apps\ComfyUI +\\YOUR-SERVER\software\ComfyUI +``` + +Workers try each entry in order after Deadline path mapping and pick the first path containing both `ComfyUI\main.py` and `python_embeded\python.exe`. ## Usage -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 +1. Add `Submit to Deadline` to the workflow. +2. Set a farm-visible `output_directory`. +3. Use `DeadlineSeed` anywhere a seed value should vary between Deadline variations. +4. Run the workflow in ComfyUI. +5. Monitor the submitted job in Deadline Monitor. -### Key Settings +`batch_count` is the total number of Deadline variations. `chunk_size` is how many variations a single Deadline task should queue into its worker ComfyUI instance. These are not animation frame ranges; Deadline frames are used internally as variation indices. -- **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 +For a `DeadlineSeed` base seed of `1000`, variation `0` uses `1000`, variation `1` uses `1001`, and so on. The worker rewrites the queued prompt before execution, so saved image metadata contains the actual seed used. Additional Deadline metadata is written under `extra_pnginfo.deadline`. -## Configuration +Deadline jobs include two workflow files when ComfyUI provides the UI workflow metadata: `prompt_to_execute.json` is the API prompt used by the worker, and `workflow.json` is the standard ComfyUI workflow with node positions. Patched workers embed the standard workflow metadata into outputs, so dropping a generated image back into ComfyUI opens the normal graph layout, not the API prompt format. -### 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. +## Input Staging -### Multiple ComfyUI Install Paths -In Deadline Monitor → Repository Options → ComfyUI plugin configuration, the `ComfyUI Installation Paths` field now accepts multiple entries. Add one path per line (drive letters, UNC shares, etc.). Workers will walk the list from top to bottom, apply Deadline path mapping, expand environment variables, and pick the first path that contains both `ComfyUI/main.py` and `python_embeded/python.exe`. This mirrors the multi-path behavior used by Deadline's Houdini plugin. +Default ComfyUI upload nodes such as `Load Image`, `Load Audio`, and `Load Video` store files in the local `ComfyUI/input` folder. On submission, this plugin copies referenced input assets to a shared sibling input folder beside the output directory: -Example: -``` -C:\ComfyUI_windows_portable -D:\Apps\ComfyUI -\\NAS01\software\ComfyUI +```text +\input\ ``` -## How It Works +Workers launch ComfyUI with `--input-directory` pointing at that staged folder. The embedded standard workflow metadata is rewritten to use absolute staged asset paths, so another ComfyUI session on a different machine can reopen the generated image as long as that shared input path is visible there. Missing or invalid referenced input files fail submission before the Deadline job is sent. -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 +Only files referenced by the submitted prompt are copied. Existing identical files are reused; conflicting filenames get a submission suffix. Absolute path loader nodes are left unchanged and must already point to farm-visible storage. -## Requirements +## Notes -- ComfyUI installation on worker machines -- Thinkbox Deadline -- No additional Python dependencies (uses standard library) +- This V2 path targets portable Windows ComfyUI workers. +- Deadline owns render timeout policy; configure timeouts in Deadline Monitor. +- The base plugin no longer registers mock `/deadline/*` routes. Distributed-worker routes remain owned by `ComfyUI-Deadline-Distributed`. diff --git a/__init__.py b/__init__.py index 99f8941..f345951 100644 --- a/__init__.py +++ b/__init__.py @@ -1,39 +1,15 @@ """ -ComfyUI Deadline Plugin +ComfyUI Deadline Plugin. -A comprehensive plugin for integrating ComfyUI workflows with Thinkbox Deadline render farm management. +Provides Deadline submission and seed nodes for ComfyUI. """ -import logging +from .deadline_submit import ( + NODE_CLASS_MAPPINGS, + NODE_DISPLAY_NAME_MAPPINGS, + register_on_prompt_handler, +) -# Import the custom nodes -from .deadline_submit import NODE_CLASS_MAPPINGS as SUBMIT_MAPPINGS, NODE_DISPLAY_NAME_MAPPINGS as SUBMIT_DISPLAY_MAPPINGS +register_on_prompt_handler() -# 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", "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 +__all__ = ["NODE_CLASS_MAPPINGS", "NODE_DISPLAY_NAME_MAPPINGS"] diff --git a/deadline_api.py b/deadline_api.py deleted file mode 100644 index 475b0c3..0000000 --- a/deadline_api.py +++ /dev/null @@ -1,229 +0,0 @@ -""" -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 ad50b37..7d031c3 100644 --- a/deadline_submit.py +++ b/deadline_submit.py @@ -1,445 +1,376 @@ -# deadline_submit.py - """ -ComfyUI Deadline Submission Node -by Dominik Bargiel dominikbargiel97@gmail.com +ComfyUI Deadline submission nodes. -A ComfyUI custom node for submitting workflows to Thinkbox Deadline render farm. +This module owns the ComfyUI-side submission flow. It packages the current API +prompt, stages referenced input assets, and submits a render job to Deadline. """ -import os -import sys +import copy +import filecmp import json -import tempfile -import subprocess -import uuid -import time +import os import re -from typing import Optional, Dict, List, Any, Union, Tuple +import shutil +import subprocess +import sys +import tempfile +import time +import uuid +from typing import Any, Dict, List, Optional, Tuple + -# Configuration constants DEADLINE_COMMAND_PATHS = { - 'windows': "C:\\Program Files\\Thinkbox\\Deadline10\\bin\\deadlinecommand.exe", - 'linux': "/opt/Thinkbox/Deadline10/bin/deadlinecommand" + "windows": "C:\\Program Files\\Thinkbox\\Deadline10\\bin\\deadlinecommand.exe", + "linux": "/opt/Thinkbox/Deadline10/bin/deadlinecommand", } -# Node configuration constants +DEADLINE_SUBMIT_NODE_TYPES = {"DeadlineSubmit", "SaveAndSubmitNode"} +OUTPUT_NODE_TYPES = {"SaveImage", "PreviewImage", "SaveVideo", "VHS_VideoCombine"} +INPUT_LOADER_FIELDS = { + "LoadImage": ("image",), + "LoadImageMask": ("image",), + "LoadAudio": ("audio",), + "LoadVideo": ("file",), +} +MEDIA_EXTENSIONS = { + ".png", ".jpg", ".jpeg", ".webp", ".bmp", ".gif", ".tif", ".tiff", ".exr", + ".mp4", ".mov", ".avi", ".mkv", ".webm", ".m4v", + ".wav", ".mp3", ".flac", ".ogg", ".m4a", ".aac", +} + + class NodeDefaults: - JOB_NAME = "ComfyUI via DeadlineNode" + JOB_NAME = "ComfyUI via Deadline" PRIORITY = 50 POOL = "none" GROUP = "none" BATCH_COUNT = 1 CHUNK_SIZE = 1 - MAX_BATCH_COUNT = 100 - MAX_CHUNK_SIZE = 16 + MAX_BATCH_COUNT = 10000 + MAX_CHUNK_SIZE = 256 MAX_PRIORITY = 100 + class DeadlineCommandHelper: - """Helper class for interacting with Deadline command line""" - @staticmethod def get_deadline_command() -> str: - """Get the path to the deadlinecommand executable""" - deadline_bin = "" - try: - deadline_bin = os.environ.get('DEADLINE_PATH', '') - except KeyError: - pass + deadline_bin = os.environ.get("DEADLINE_PATH", "") if not deadline_bin and os.path.exists("/Users/Shared/Thinkbox/DEADLINE_PATH"): try: - with open("/Users/Shared/Thinkbox/DEADLINE_PATH") as f: - deadline_bin = f.read().strip() + with open("/Users/Shared/Thinkbox/DEADLINE_PATH", "r", encoding="utf-8") as handle: + deadline_bin = handle.read().strip() except Exception: - pass + deadline_bin = "" + candidates = [] if deadline_bin: - deadline_command = os.path.join(deadline_bin, "deadlinecommand") - if os.path.exists(deadline_command): - return deadline_command + candidates.append(os.path.join(deadline_bin, "deadlinecommand.exe" if os.name == "nt" else "deadlinecommand")) + candidates.append(DEADLINE_COMMAND_PATHS["windows"] if sys.platform.startswith("win") else DEADLINE_COMMAND_PATHS["linux"]) - # Try platform-specific default paths - if sys.platform.startswith('win'): - default_path = DEADLINE_COMMAND_PATHS['windows'] - else: - default_path = DEADLINE_COMMAND_PATHS['linux'] - - if os.path.exists(default_path): - return default_path - + for candidate in candidates: + if os.path.exists(candidate): + return candidate return "" @staticmethod - def call_deadline_command(arguments: List[str], hide_window: bool = True, read_stdout: bool = True) -> str: - """Call deadlinecommand with the given arguments""" + def call_deadline_command(arguments: List[str], hide_window: bool = True) -> str: deadline_command = DeadlineCommandHelper.get_deadline_command() if not deadline_command: - raise Exception("Deadline command not found") - + raise RuntimeError("Deadline command not found. Set DEADLINE_PATH or install Deadline Client.") + startupinfo = None creationflags = 0 - - if os.name == 'nt': - if hide_window: - try: - startupinfo = subprocess.STARTUPINFO() - if hasattr(subprocess, '_subprocess') and hasattr(subprocess._subprocess, 'STARTF_USESHOWWINDOW'): - startupinfo.dwFlags |= subprocess._subprocess.STARTF_USESHOWWINDOW - elif hasattr(subprocess, 'STARTF_USESHOWWINDOW'): - startupinfo.dwFlags |= subprocess.STARTF_USESHOWWINDOW - except: - pass - else: - CREATE_NO_WINDOW = 0x08000000 - creationflags = CREATE_NO_WINDOW - - full_arguments = [deadline_command] + arguments - - proc = subprocess.Popen( - full_arguments, - stdin=subprocess.PIPE, - stdout=subprocess.PIPE, - stderr=subprocess.PIPE, - startupinfo=startupinfo, - creationflags=creationflags + if os.name == "nt" and hide_window: + try: + startupinfo = subprocess.STARTUPINFO() + startupinfo.dwFlags |= subprocess.STARTF_USESHOWWINDOW + except Exception: + startupinfo = None + elif os.name == "nt": + creationflags = 0x08000000 + + process = subprocess.Popen( + [deadline_command] + arguments, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + startupinfo=startupinfo, + creationflags=creationflags, ) - - output = "" - if read_stdout: - output, errors = proc.communicate() - - if sys.version_info[0] >= 3 and isinstance(output, bytes): - output = output.decode(errors="replace") - + stdout, stderr = process.communicate() + output = stdout.decode(errors="replace") if isinstance(stdout, bytes) else str(stdout) + errors = stderr.decode(errors="replace") if isinstance(stderr, bytes) else str(stderr) + + if process.returncode != 0: + raise RuntimeError(f"deadlinecommand failed with code {process.returncode}: {errors or output}") return output @staticmethod def get_job_id_from_submission(submission_results: str) -> str: - """Parse the job ID from the submission results""" - for line in submission_results.split(): - if line.startswith("JobID="): - return line.replace("JobID=", "").strip() + for token in submission_results.replace("\r", "\n").split(): + if token.startswith("JobID="): + return token.split("=", 1)[1].strip() return "" + class WorkflowProcessor: - """Handles workflow data processing and validation""" - @staticmethod - def normalize_workflow(workflow_data: Union[Dict, List]) -> Optional[Dict]: - """Normalize workflow data to ensure compatibility""" - if not workflow_data: - print("Deadline Submission: Error - Empty workflow data.") - return None - - # If workflow is already in UI format (dictionary with node IDs as keys) - if isinstance(workflow_data, dict): - is_ui_format = any(isinstance(key, str) and key.isdigit() for key in workflow_data.keys()) - if is_ui_format: - return workflow_data - - # If it's the API format (list of nodes) - if isinstance(workflow_data, list): - return WorkflowProcessor._convert_api_to_ui_format(workflow_data) - - # Not recognized format - print(f"Deadline Submission: Warning - Unrecognized workflow format. Attempting to use as-is.") - return workflow_data if isinstance(workflow_data, dict) else None + def normalize_prompt(prompt: Dict[str, Any]) -> Dict[str, Any]: + if not isinstance(prompt, dict) or not prompt: + raise ValueError("ComfyUI did not provide a valid API prompt.") + return copy.deepcopy(prompt) @staticmethod - def _convert_api_to_ui_format(workflow_list: List) -> Dict: - """Convert API format workflow to UI format""" - ui_format = {} - for node in workflow_list: - if isinstance(node, list) and len(node) >= 3: - node_id = str(node[0]) - ui_format[node_id] = { - "class_type": node[1], - "inputs": node[2] - } - return ui_format + def prepare_for_worker(prompt: Dict[str, Any]) -> Dict[str, Any]: + prepared = WorkflowProcessor.normalize_prompt(prompt) + removed = [] + for node_id, node in list(prepared.items()): + if isinstance(node, dict) and node.get("class_type") in DEADLINE_SUBMIT_NODE_TYPES: + removed.append(node_id) + del prepared[node_id] + + if removed: + print(f"Deadline Submission: Removed submit node(s) from worker prompt: {', '.join(map(str, removed))}") + + WorkflowProcessor.validate_worker_prompt(prepared) + return prepared @staticmethod - def validate_workflow(workflow_data: Dict) -> bool: - """Basic validation that workflow contains important nodes""" - if not workflow_data: - return False - - has_output_node = False - has_checkpoint = False - - output_node_types = ["SaveImage", "PreviewImage", "SaveVideo"] - checkpoint_types = ["CheckpointLoaderSimple", "CheckpointLoader", "UNETLoader"] - - for node_id, node in workflow_data.items(): - if not isinstance(node, dict) or "class_type" not in node: + def validate_worker_prompt(prompt: Dict[str, Any]) -> None: + if not prompt: + raise ValueError("Worker prompt is empty after removing Deadline submit nodes.") + + has_output = any( + isinstance(node, dict) and node.get("class_type") in OUTPUT_NODE_TYPES + for node in prompt.values() + ) + if not has_output: + print("Deadline Submission: Warning - worker prompt has no known output node.") + + +class InputAssetStager: + def __init__(self, output_directory: str, job_name: str, submission_id: str): + self.output_directory = os.path.abspath(output_directory) + self.job_name = job_name + self.submission_id = submission_id + + def stage_referenced_assets(self, prompt: Dict[str, Any]) -> Tuple[str, str, List[Dict[str, Any]]]: + input_dir = self._get_local_input_directory() + references = self._collect_references(prompt, input_dir) + staging_dir = self._staging_directory() + manifest_path = os.path.join(staging_dir, f"deadline_input_manifest_{self.submission_id}.json") + + if not references: + os.makedirs(staging_dir, exist_ok=True) + manifest = { + "submission_id": self.submission_id, + "input_directory": staging_dir, + "assets": [], + } + self._write_manifest(manifest_path, manifest) + return staging_dir, manifest_path, [] + + os.makedirs(staging_dir, exist_ok=True) + assets = [] + for original_rel_path, source_path in sorted(references.items()): + staged_rel_path, destination_path = self._resolve_destination(staging_dir, original_rel_path, source_path) + os.makedirs(os.path.dirname(destination_path), exist_ok=True) + if not os.path.exists(destination_path): + shutil.copy2(source_path, destination_path) + assets.append({ + "original_relative_path": original_rel_path.replace("\\", "/"), + "staged_relative_path": staged_rel_path.replace("\\", "/"), + "relative_path": staged_rel_path.replace("\\", "/"), + "source": source_path, + "destination": destination_path, + "size": os.path.getsize(destination_path), + }) + + manifest = { + "submission_id": self.submission_id, + "input_directory": staging_dir, + "assets": assets, + } + self._write_manifest(manifest_path, manifest) + print(f"Deadline Submission: Staged {len(assets)} input asset(s) to {staging_dir}") + return staging_dir, manifest_path, assets + + def _get_local_input_directory(self) -> str: + try: + import folder_paths + return os.path.abspath(folder_paths.get_input_directory()) + except Exception as exc: + raise RuntimeError(f"Could not resolve ComfyUI input directory: {exc}") + + def _collect_references(self, prompt: Dict[str, Any], input_dir: str) -> Dict[str, str]: + references: Dict[str, str] = {} + for node in prompt.values(): + if not isinstance(node, dict): continue - + class_type = node.get("class_type", "") - - if class_type in output_node_types: - has_output_node = True - - if class_type in checkpoint_types: - has_checkpoint = True - - if not has_output_node: - print("Deadline Submission: Warning - No output nodes found in workflow.") - - if not has_checkpoint: - print("Deadline Submission: Warning - No checkpoint loader found in workflow.") - - return True + inputs = node.get("inputs", {}) + if not isinstance(inputs, dict): + continue - @staticmethod - def save_workflow_file(workflow_data: Dict, file_path: Optional[str] = None) -> Optional[str]: - """Save workflow data to a file for submission""" - if not workflow_data: - print("Deadline Submission: No workflow data to save.") + candidate_values: List[Tuple[Any, bool]] = [] + for field_name in INPUT_LOADER_FIELDS.get(class_type, ()): + if field_name in inputs: + candidate_values.append((inputs[field_name], True)) + + for value in inputs.values(): + if isinstance(value, str): + candidate_values.append((value, False)) + + for value, strict in candidate_values: + if not isinstance(value, str): + continue + resolved = self._resolve_input_file(value, input_dir, strict) + if not resolved: + continue + rel_path, source_path = resolved + references[rel_path] = source_path + return references + + def _resolve_destination(self, staging_dir: str, rel_path: str, source_path: str) -> Tuple[str, str]: + destination_path = os.path.abspath(os.path.join(staging_dir, rel_path)) + if not os.path.exists(destination_path): + return rel_path, destination_path + + if os.path.isfile(destination_path) and filecmp.cmp(source_path, destination_path, shallow=False): + return rel_path, destination_path + + stem, extension = os.path.splitext(rel_path) + staged_rel_path = f"{stem}_{self.submission_id}{extension}" + return staged_rel_path, os.path.abspath(os.path.join(staging_dir, staged_rel_path)) + + def _resolve_input_file(self, value: str, input_dir: str, strict: bool) -> Optional[Tuple[str, str]]: + clean_value, annotation = self._strip_annotation(value) + if annotation in {"output", "temp"}: return None - - if not file_path: - temp_dir = tempfile.gettempdir() - file_path = os.path.join(temp_dir, f"comfyui_workflow_for_deadline_{uuid.uuid4()}.json") - - try: - with open(file_path, 'w') as f: - json.dump(workflow_data, f, indent=2) - - # Create a metadata file for debugging - WorkflowProcessor._create_metadata_file(file_path) - - print(f"Deadline Submission: Successfully saved workflow for submission to: {file_path}") - return file_path - except Exception as e: - print(f"Deadline Submission: Error saving workflow file: {e}") + if not clean_value or os.path.isabs(clean_value): + return None + if os.path.splitext(clean_value)[1].lower() not in MEDIA_EXTENSIONS: return None - @staticmethod - def _create_metadata_file(workflow_path: str): - """Create a metadata file alongside the workflow""" + candidate = os.path.abspath(os.path.join(input_dir, clean_value)) try: - with open(f"{workflow_path}.metadata", 'w') as f: - metadata = { - "generator": "ComfyUI Deadline Submission Plugin", - "captured_at": time.strftime("%Y-%m-%d %H:%M:%S"), - "notes": "This workflow was captured and prepared for Deadline rendering." - } - json.dump(metadata, f, indent=2) - except Exception as e: - print(f"Deadline Submission: Warning - Could not create metadata file: {e}") + common = os.path.commonpath([input_dir, candidate]) + except ValueError: + common = "" + if common != input_dir: + raise ValueError(f"Input asset escapes ComfyUI input directory: {value}") + if not os.path.isfile(candidate): + if not strict: + return None + raise FileNotFoundError(f"Referenced input asset was not found: {value} ({candidate})") + + rel_path = os.path.relpath(candidate, input_dir) + return rel_path, candidate + + def _strip_annotation(self, value: str) -> Tuple[str, Optional[str]]: + match = re.match(r"^(.*)\s+\[(input|output|temp)\]\s*$", value) + if not match: + return value.strip().replace("/", os.sep), None + return match.group(1).strip().replace("/", os.sep), match.group(2) + + def _staging_directory(self) -> str: + parent = os.path.dirname(self.output_directory.rstrip("\\/")) + return os.path.join(parent, "input") + + def _write_manifest(self, manifest_path: str, manifest: Dict[str, Any]) -> None: + with open(manifest_path, "w", encoding="utf-8") as handle: + json.dump(manifest, handle, indent=2) - @staticmethod - def prepare_workflow_for_submission(workflow_data: Dict) -> Dict: - """Prepare workflow by setting DeadlineSubmit nodes to bypassed""" - normalized_workflow = WorkflowProcessor.normalize_workflow(workflow_data) - if not normalized_workflow: - raise Exception("Failed to normalize workflow") - - # Set any DeadlineSubmit nodes to bypassed - deadline_node_types = ["DeadlineSubmit", "SaveAndSubmitNode"] - for node_id, node in normalized_workflow.items(): - if isinstance(node, dict) and node.get("class_type") in deadline_node_types: - print(f"Deadline Submission: Setting node {node_id} to bypassed") - if "inputs" not in node: - node["inputs"] = {} - node["inputs"]["bypass"] = True - - WorkflowProcessor.validate_workflow(normalized_workflow) - return normalized_workflow class DeadlineJobSubmitter: - """Handles submission of jobs to Deadline""" - - def __init__(self, workflow_data: Dict, job_config: Dict): + def __init__(self, workflow_data: Dict[str, Any], job_config: Dict[str, Any]): self.workflow_data = workflow_data self.job_config = job_config def submit_job(self) -> Tuple[bool, str]: - """Submit the job to Deadline and return success status and job ID or error message""" try: - workflow_path = self._save_workflow() - if not workflow_path: - return False, "Failed to save workflow for submission" - - job_id = self._submit_to_deadline(workflow_path) - if job_id: - return True, job_id - else: - return False, "Job submitted but no JobID returned" - - except Exception as e: - return False, f"Error submitting to Deadline: {str(e)}" - - def _save_workflow(self) -> Optional[str]: - """Save the workflow to a temporary file""" - return WorkflowProcessor.save_workflow_file(self.workflow_data) - - def _submit_to_deadline(self, workflow_path: str) -> str: - """Submit the workflow to Deadline and return job ID""" - submission_temp_dir = tempfile.mkdtemp(prefix="comfy_deadline_job_") - - try: - job_info_file, plugin_info_file, workflow_copy = self._create_submission_files( - submission_temp_dir, workflow_path - ) - - command_args = [job_info_file, plugin_info_file, workflow_copy] - result = DeadlineCommandHelper.call_deadline_command(command_args) - + submission_dir = tempfile.mkdtemp(prefix="comfy_deadline_job_") + job_info_file, plugin_info_file, auxiliary_files = self._create_submission_files(submission_dir) + result = DeadlineCommandHelper.call_deadline_command([job_info_file, plugin_info_file] + auxiliary_files) job_id = DeadlineCommandHelper.get_job_id_from_submission(result) - if job_id: - print(f"Deadline Submission: Successfully submitted job. JobID: {job_id}") - return job_id - else: - print(f"Deadline Submission: Job submitted but JobID not found. Result: {result}") - return "" - - except Exception as e: - print(f"Deadline Submission: Error during submission: {e}") - raise + if not job_id: + return False, f"Deadline submission did not return a JobID. Output: {result}" + return True, job_id + except Exception as exc: + return False, str(exc) - def _create_submission_files(self, temp_dir: str, workflow_path: str) -> Tuple[str, str, str]: - """Create job info and plugin info files for submission""" - job_info_file = os.path.join(temp_dir, "job_info.txt") - plugin_info_file = os.path.join(temp_dir, "plugin_info.txt") - - # Copy workflow to submission directory - workflow_copy = os.path.join(temp_dir, "workflow_to_submit.json") - try: - import shutil - shutil.copy2(workflow_path, workflow_copy) - except Exception: - workflow_copy = workflow_path + def _create_submission_files(self, submission_dir: str) -> Tuple[str, str, List[str]]: + job_info_file = os.path.join(submission_dir, "job_info.txt") + plugin_info_file = os.path.join(submission_dir, "plugin_info.txt") + prompt_file = os.path.join(submission_dir, "prompt_to_execute.json") + standard_workflow_file = os.path.join(submission_dir, "workflow.json") - self._create_job_info_file(job_info_file) - self._create_plugin_info_file(plugin_info_file) - - return job_info_file, plugin_info_file, workflow_copy + with open(prompt_file, "w", encoding="utf-8") as handle: + json.dump(self.workflow_data, handle, indent=2) - def _create_job_info_file(self, job_info_file: str): - """Create the job info file""" + auxiliary_files = [prompt_file] + standard_workflow = self.job_config.get("standard_workflow") + if standard_workflow: + with open(standard_workflow_file, "w", encoding="utf-8") as handle: + json.dump(standard_workflow, handle, indent=2) + auxiliary_files.append(standard_workflow_file) + + self._write_job_info(job_info_file) + self._write_plugin_info(plugin_info_file) + return job_info_file, plugin_info_file, auxiliary_files + + def _write_job_info(self, path: str) -> None: config = self.job_config - - with open(job_info_file, 'w') as f: - f.write(f"Plugin=ComfyUI\n") - f.write(f"Name={config['job_name']}\n") - f.write(f"Comment={config.get('comment', '')}\n") - f.write(f"Department={config.get('department', '')}\n") - f.write(f"Pool={config['pool'] if config['pool'] != 'none' else ''}\n") - f.write(f"Group={config['group'] if config['group'] != 'none' else ''}\n") - f.write(f"Priority={config['priority']}\n") - - # Add frame range if batch count > 1 - if config['batch_count'] > 1: - f.write(f"Frames=0-{config['batch_count'] - 1}\n") - f.write(f"ChunkSize={config['chunk_size']}\n") - else: - f.write(f"Frames=0\n") - f.write(f"ChunkSize=1\n") - - # Add output directory if specified - if config.get('output_directory'): - abs_output_dir = os.path.abspath(config['output_directory'].strip()) - f.write(f"OutputDirectory0={abs_output_dir}\n") + batch_count = int(config["batch_count"]) + chunk_size = max(1, int(config["chunk_size"])) + output_dir = os.path.abspath(config["output_directory"]) - def _create_plugin_info_file(self, plugin_info_file: str): - """Create the plugin info file""" + with open(path, "w", encoding="utf-8") as handle: + handle.write("Plugin=ComfyUI\n") + handle.write(f"Name={config['job_name']}\n") + handle.write(f"Comment={config.get('comment', '')}\n") + handle.write(f"Department={config.get('department', '')}\n") + handle.write(f"Pool={'' if config['pool'] == 'none' else config['pool']}\n") + handle.write(f"Group={'' if config['group'] == 'none' else config['group']}\n") + handle.write(f"Priority={int(config['priority'])}\n") + handle.write(f"Frames=0-{batch_count - 1}\n") + handle.write(f"ChunkSize={chunk_size}\n") + handle.write(f"OutputDirectory0={output_dir}\n") + + def _write_plugin_info(self, path: str) -> None: config = self.job_config - - with open(plugin_info_file, 'w') as f: - if config.get('output_directory'): - abs_output_dir = os.path.abspath(config['output_directory'].strip()) - f.write(f"JobOutputDirectory={abs_output_dir}\n") - - f.write("DefaultCudaDeviceZero=True\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") + entries = { + "StandardWorkflowFile": "workflow.json" if config.get("standard_workflow") else "", + "JobOutputDirectory": os.path.abspath(config["output_directory"]), + "JobInputDirectory": os.path.abspath(config["input_directory"]), + "InputManifestFile": os.path.abspath(config["input_manifest"]), + "SubmissionId": config["submission_id"], + "BatchCount": str(int(config["batch_count"])), + "BatchMode": "True", + "DefaultCudaDeviceZero": "True", + "SeedMode": "fixed", + "WorkerMode": "False", + "DistributedMode": "False", + "ForceNewInstance": "True", + } -class ExecutionInterruptor: - """Handles interrupting local ComfyUI execution""" - - @staticmethod - def interrupt_local_execution(): - """Attempt to interrupt local ComfyUI execution""" - try: - import sys - - # Try to interrupt using the known working approach - if ExecutionInterruptor._try_nodes_interrupt(): - print("Deadline Submission: Successfully interrupted via nodes module") - elif ExecutionInterruptor._try_comfy_graph_interrupt(): - print("Deadline Submission: Successfully interrupted via comfy.graph module") - else: - print("Deadline Submission: No interruption mechanism found, local execution may still occur") - - except Exception as e: - print(f"Deadline Submission: Unable to prevent local execution (safe to ignore): {str(e)}") + with open(path, "w", encoding="utf-8") as handle: + for key, value in entries.items(): + if value != "": + handle.write(f"{key}={value}\n") - @staticmethod - def _try_nodes_interrupt() -> bool: - """Try to interrupt using the nodes module""" - import sys - - if 'nodes' not in sys.modules: - return False - - nodes_module = sys.modules['nodes'] - if not hasattr(nodes_module, 'interrupt_processing'): - return False - - interrupt_attr = getattr(nodes_module, 'interrupt_processing') - if callable(interrupt_attr): - interrupt_attr(True) - else: - nodes_module.interrupt_processing = True - - return True - - @staticmethod - def _try_comfy_graph_interrupt() -> bool: - """Try to interrupt using the comfy.graph module""" - import sys - - if 'comfy' not in sys.modules: - return False - - comfy_module = sys.modules['comfy'] - if not hasattr(comfy_module, 'graph'): - return False - - graph_module = comfy_module.graph - if not hasattr(graph_module, 'interrupt_processing'): - return False - - interrupt_attr = getattr(graph_module, 'interrupt_processing') - if callable(interrupt_attr): - interrupt_attr(True) - else: - graph_module.interrupt_processing = True - - 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, + "default": 1125899906842, "min": 0, "max": 1125899906842624, - "forceInput": False # Widget by default, can be converted to input + "forceInput": False, }), }, "hidden": { @@ -447,100 +378,64 @@ class DeadlineSeed: "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): + except (TypeError, ValueError): 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,) + seed = int(seed) + if batch_mode and task_id: + return (seed + task_id,) + return (seed,) + -# Node implementation class DeadlineSubmitNode: - """Submit the current ComfyUI workflow to Thinkbox Deadline""" - @classmethod def INPUT_TYPES(cls): pools = cls._get_deadline_pools() groups = cls._get_deadline_groups() - return { "required": { - "workflow_file": ("STRING", { - "default": "", - "multiline": False, - "placeholder": "(Optional) Override if auto-detect is OFF" - }), - "auto_detect_workflow": ("BOOLEAN", { - "default": True, - "label_on": "Use current (recommended)", - "label_off": "Use 'workflow_file' input" + "output_directory": ("STRING", { + "default": "", + "multiline": False, + "placeholder": "Farm-visible output directory", }), "batch_count": ("INT", { - "default": NodeDefaults.BATCH_COUNT, - "min": 1, - "max": NodeDefaults.MAX_BATCH_COUNT, - "step": 1 + "default": NodeDefaults.BATCH_COUNT, + "min": 1, + "max": NodeDefaults.MAX_BATCH_COUNT, + "step": 1, }), "chunk_size": ("INT", { - "default": NodeDefaults.CHUNK_SIZE, - "min": 1, - "max": NodeDefaults.MAX_CHUNK_SIZE, - "step": 1 + "default": NodeDefaults.CHUNK_SIZE, + "min": 1, + "max": NodeDefaults.MAX_CHUNK_SIZE, + "step": 1, }), - "priority": ("INT", { - "default": NodeDefaults.PRIORITY, - "min": 0, - "max": NodeDefaults.MAX_PRIORITY + "default": NodeDefaults.PRIORITY, + "min": 0, + "max": NodeDefaults.MAX_PRIORITY, }), "pool": (pools, {"default": NodeDefaults.POOL}), "group": (groups, {"default": NodeDefaults.GROUP}), "job_name": ("STRING", {"default": NodeDefaults.JOB_NAME}), - "bypass": ("BOOLEAN", {"default": False}), - "skip_local_execution": ("BOOLEAN", { - "default": True, - "label_on": "Submit Only", - "label_off": "Submit and Run Locally" - }), }, "optional": { - "output_directory": ("STRING", { - "default": "", - "multiline": False, - "placeholder": "(Optional) Output directory on worker" - }), "comment": ("STRING", {"default": ""}), "department": ("STRING", {"default": ""}), - }, "hidden": { "prompt": "PROMPT", "extra_pnginfo": "EXTRA_PNGINFO", - } + }, } RETURN_TYPES = ("STRING",) @@ -551,117 +446,210 @@ class DeadlineSubmitNode: @classmethod def IS_CHANGED(cls, **kwargs): - """Return a unique value each time to force execution""" - return f"deadline_submit_{time.time()}" + return f"deadline_submit_{time.time()}_{uuid.uuid4()}" @classmethod def _get_deadline_pools(cls) -> List[str]: - """Get available Deadline pools""" try: - result = DeadlineCommandHelper.call_deadline_command(["-pools"], hide_window=True) - pools = [line.strip() for line in result.splitlines() if line.strip()] - return pools if pools else [NodeDefaults.POOL] - except Exception as e: - print(f"Deadline Submission: Error getting Deadline pools: {e}") + output = DeadlineCommandHelper.call_deadline_command(["-pools"]) + pools = [line.strip() for line in output.splitlines() if line.strip()] + return pools or [NodeDefaults.POOL] + except Exception as exc: + print(f"Deadline Submission: Could not query Deadline pools: {exc}") return [NodeDefaults.POOL] @classmethod def _get_deadline_groups(cls) -> List[str]: - """Get available Deadline groups""" try: - result = DeadlineCommandHelper.call_deadline_command(["-groups"], hide_window=True) - groups = [line.strip() for line in result.splitlines() if line.strip()] - return groups if groups else [NodeDefaults.GROUP] - except Exception as e: - print(f"Deadline Submission: Error getting Deadline groups: {e}") + output = DeadlineCommandHelper.call_deadline_command(["-groups"]) + groups = [line.strip() for line in output.splitlines() if line.strip()] + return groups or [NodeDefaults.GROUP] + except Exception as exc: + print(f"Deadline Submission: Could not query Deadline groups: {exc}") return [NodeDefaults.GROUP] - def submit_to_deadline(self, workflow_file, auto_detect_workflow, batch_count, chunk_size, - 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""" - if bypass: - print("Deadline Submission: Bypass enabled. Submission skipped.") - return ("Bypassed",) - - print(f"Deadline Submission: Node execution triggered. Auto-detect: {auto_detect_workflow}") - + def submit_to_deadline( + self, + output_directory, + batch_count, + chunk_size, + priority, + pool, + group, + job_name, + comment="", + department="", + prompt=None, + extra_pnginfo=None, + **_legacy_inputs, + ): try: - # Get workflow data - workflow_data = self._get_workflow_data(auto_detect_workflow, workflow_file, prompt) - - # Prepare workflow for submission - prepared_workflow = WorkflowProcessor.prepare_workflow_for_submission(workflow_data) - - # Create job configuration - job_config = self._create_job_config( - job_name, priority, pool, group, batch_count, chunk_size, - output_directory, comment, department - ) - - # Submit to Deadline - submitter = DeadlineJobSubmitter(prepared_workflow, job_config) - success, result = submitter.submit_job() - - if success: - if skip_local_execution: - ExecutionInterruptor.interrupt_local_execution() - return (result,) - else: - return (f"Error: {result}",) - - except Exception as e: - print(f"Deadline Submission: Error during submission: {e}") - return (f"Error: {str(e)}",) + output_directory = self._prepare_output_directory(output_directory) + batch_count = max(1, int(batch_count)) + chunk_size = max(1, min(int(chunk_size), batch_count)) + submission_id = uuid.uuid4().hex[:12] - def _get_workflow_data(self, auto_detect_workflow: bool, workflow_file: str, prompt) -> Dict: - """Get workflow data from either auto-detection or file""" - if auto_detect_workflow: - print("Deadline Submission: Auto-detect ON. Checking for workflow...") - if prompt is None: - raise Exception("ComfyUI did not inject PROMPT parameter") - - print("Deadline Submission: Found workflow from ComfyUI's PROMPT parameter injection") - return prompt - else: - print(f"Deadline Submission: Auto-detect OFF. Using specified workflow_file: '{workflow_file}'.") - user_workflow_path = workflow_file.strip() - - if not user_workflow_path or not os.path.exists(user_workflow_path): - raise Exception(f"Specified workflow file not found: '{user_workflow_path}'") - + raise ValueError("ComfyUI did not inject the current API prompt.") + + worker_prompt = WorkflowProcessor.prepare_for_worker(prompt) + stager = InputAssetStager(output_directory, job_name, submission_id) + input_dir, manifest_file, assets = stager.stage_referenced_assets(worker_prompt) + self._rewrite_prompt_asset_references(worker_prompt, assets) + standard_workflow = self._extract_standard_workflow(extra_pnginfo) + self._rewrite_standard_workflow_assets(standard_workflow, assets) + + job_config = { + "submission_id": submission_id, + "job_name": job_name.strip() or NodeDefaults.JOB_NAME, + "priority": int(priority), + "pool": pool, + "group": group, + "batch_count": batch_count, + "chunk_size": chunk_size, + "output_directory": output_directory, + "input_directory": input_dir, + "input_manifest": manifest_file, + "standard_workflow": standard_workflow, + "comment": comment, + "department": department, + } + + submitter = DeadlineJobSubmitter(worker_prompt, job_config) + success, result = submitter.submit_job() + if not success: + raise RuntimeError(result) + + print(f"Deadline Submission: Submitted job {result} with {batch_count} variation(s), chunk size {chunk_size}, {len(assets)} staged asset(s).") + return (result,) + except Exception as exc: + print(f"Deadline Submission: Error during submission: {exc}") + raise + + def _rewrite_prompt_asset_references(self, prompt: Dict[str, Any], assets: List[Dict[str, Any]]) -> None: + staged_by_original = self._asset_map(assets, "staged_relative_path") + if not staged_by_original: + return + + for node in prompt.values(): + if not isinstance(node, dict): + continue + inputs = node.get("inputs") + if not isinstance(inputs, dict): + continue + + class_type = node.get("class_type", "") + for field_name in INPUT_LOADER_FIELDS.get(class_type, ()): + value = inputs.get(field_name) + normalized = self._normalize_asset_reference(value) + if normalized in staged_by_original: + inputs[field_name] = staged_by_original[normalized] + + for field_name, value in list(inputs.items()): + normalized = self._normalize_asset_reference(value) + if normalized in staged_by_original: + inputs[field_name] = staged_by_original[normalized] + + def _rewrite_standard_workflow_assets(self, workflow: Optional[Dict[str, Any]], assets: List[Dict[str, Any]]) -> None: + staged_absolute_by_original = self._asset_map(assets, "destination") + if not workflow or not staged_absolute_by_original: + return + + nodes = workflow.get("nodes", []) + if not isinstance(nodes, list): + return + + for node in nodes: + if not isinstance(node, dict): + continue + widgets = node.get("widgets_values") + if not isinstance(widgets, list): + continue + + for index, value in enumerate(widgets): + normalized = self._normalize_asset_reference(value) + if normalized in staged_absolute_by_original: + widgets[index] = staged_absolute_by_original[normalized] + + def _asset_map(self, assets: List[Dict[str, Any]], target_key: str) -> Dict[str, str]: + mapping = {} + for asset in assets: + original = self._normalize_asset_reference(asset.get("original_relative_path")) + target = asset.get(target_key) + if original and target: + mapping[original] = target + return mapping + + def _normalize_asset_reference(self, value: Any) -> Optional[str]: + if not isinstance(value, str) or not value: + return None + cleaned = re.sub(r"\s+\[(input|output|temp)\]\s*$", "", value.strip()) + if os.path.isabs(cleaned): + return None + return os.path.normpath(cleaned.replace("/", os.sep)).replace("\\", "/") + + def _extract_standard_workflow(self, extra_pnginfo: Any) -> Optional[Dict[str, Any]]: + if not isinstance(extra_pnginfo, dict): + return None + + workflow = extra_pnginfo.get("workflow") + if isinstance(workflow, dict): + return copy.deepcopy(workflow) + if isinstance(workflow, str): try: - with open(user_workflow_path, 'r') as f: - return json.load(f) - except Exception as e: - raise Exception(f"Could not read workflow file: {str(e)}") + parsed = json.loads(workflow) + return parsed if isinstance(parsed, dict) else None + except json.JSONDecodeError: + return None + return None + + def _prepare_output_directory(self, output_directory: str) -> str: + output_directory = (output_directory or "").strip().strip("\"") + if not output_directory: + raise ValueError("output_directory is required and must be farm-visible.") + + output_directory = os.path.abspath(os.path.expandvars(output_directory)) + os.makedirs(output_directory, exist_ok=True) + if not os.path.isdir(output_directory): + raise ValueError(f"Output path is not a directory: {output_directory}") + return output_directory + + +def on_prompt(json_data: Dict[str, Any]) -> Dict[str, Any]: + prompt = json_data.get("prompt") + if not isinstance(prompt, dict): + return json_data + + submit_ids = [ + str(node_id) + for node_id, node in prompt.items() + if isinstance(node, dict) and node.get("class_type") in DEADLINE_SUBMIT_NODE_TYPES + ] + + if submit_ids: + json_data["partial_execution_targets"] = submit_ids[:1] + if len(submit_ids) > 1: + print(f"Deadline Submission: Multiple submit nodes detected; only node {submit_ids[0]} will execute locally.") + return json_data + + +def register_on_prompt_handler() -> None: + try: + import server + instance = getattr(getattr(server, "PromptServer", None), "instance", None) + if instance and hasattr(instance, "add_on_prompt_handler"): + instance.add_on_prompt_handler(on_prompt) + print("Deadline Submission: Registered submit-only prompt handler.") + except Exception as exc: + print(f"Deadline Submission: Could not register prompt handler: {exc}") - def _create_job_config(self, job_name: str, priority: int, pool: str, group: str, - batch_count: int, chunk_size: int, - output_directory: str, comment: str, department: str) -> Dict: - """Create job configuration dictionary""" - return { - 'job_name': job_name, - 'priority': priority, - 'pool': pool, - 'group': group, - 'batch_count': batch_count, - 'chunk_size': chunk_size, - 'output_directory': output_directory, - 'comment': comment, - 'department': department - } -# 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 a5f29e9..342338d 100644 --- a/plugins/ComfyUI/ComfyUI.py +++ b/plugins/ComfyUI/ComfyUI.py @@ -15,6 +15,8 @@ import urllib.parse import traceback import random import platform +import copy +import uuid from typing import Tuple """ @@ -30,7 +32,6 @@ Supports both existing ComfyUI instances and launching new ones. DEFAULT_PORT = 8188 PORT_OFFSET_PER_GPU = 100 MAX_PORT_SEARCH_RANGE = 100 -DEFAULT_TIMEOUT = 6000 # 100 minutes DEFAULT_POLLING_INTERVAL = 10 # seconds MAX_SEED_VALUE = 2147483647 PROGRESS_LOG_INTERVAL = 10 # Log every 10 polls @@ -102,7 +103,7 @@ class ComfyUI(DeadlinePlugin): # Progress handlers self.AddStdoutHandlerCallback(r"\s*([0-9]+)%\|.*\|\s*([0-9]+)/([0-9]+).*").HandleCallback += self.HandleStdoutProgressBar self.AddStdoutHandlerCallback(r"Progress: ([0-9.]+)%.*").HandleCallback += self.HandleStdoutProgressPercent - self.AddStdoutHandlerCallback(r"Prompt executed in ([0-9.]+) seconds").HandleCallback += self.HandleStdoutPromptExecuted + # Completion is tracked through /history so Deadline's own timeout policy remains authoritative. def _initialize_member_variables(self): """Initialize all member variables""" @@ -119,11 +120,18 @@ class ComfyUI(DeadlinePlugin): self.progress_value = 0 self.thread_running = True self.custom_output_dir_specified = False + self.comfyui_input_dir = "" + self.input_manifest_file = "" + self.submission_id = "" + self.standard_workflow_file = "" + self.standard_workflow = None self.comfyui_install_path = None self.comfyui_path_candidates = [] # Batch processing variables self.chunk_size = 1 + self.batch_count = 1 + self.assigned_variation_indices = [0] self.prompts_executed = 0 self.batch_mode = False @@ -216,7 +224,7 @@ class ComfyUI(DeadlinePlugin): def InitializeProcess(self): """Initialize process settings""" - self.SingleFramesOnly = True + self.SingleFramesOnly = False self.PluginType = PluginType.Simple self.ProcessPriority = ProcessPriorityClass.BelowNormal self.UseProcessTree = True @@ -282,15 +290,44 @@ class ComfyUI(DeadlinePlugin): def _setup_batch_processing(self): """Setup batch processing configuration""" self.batch_mode = self.GetBooleanPluginInfoEntryWithDefault("BatchMode", False) + self.batch_count = int(self.GetPluginInfoEntryWithDefault("BatchCount", "1")) if self.batch_mode: self.chunk_size = int(self.GetJob().ChunkSize) - self.LogInfo(f"Batch mode enabled. Chunk size: {self.chunk_size}") + self.assigned_variation_indices = self._get_assigned_variation_indices() + self.chunk_size = len(self.assigned_variation_indices) + self.LogInfo(f"Batch mode enabled. Assigned variation indices: {self.assigned_variation_indices}") else: self.chunk_size = 1 + self.assigned_variation_indices = [0] self.LogInfo("Batch mode disabled. Processing single task.") self.prompts_executed = 0 + def _get_assigned_variation_indices(self): + """Return Deadline frame numbers for the current task; frames are variation indices.""" + chunk_size = max(1, int(self.GetJob().ChunkSize)) + task_id = int(self.GetCurrentTaskId()) + start_from_task = task_id * chunk_size + end_from_task = min(start_from_task + chunk_size - 1, self.batch_count - 1) + indices_from_task = list(range(start_from_task, end_from_task + 1)) + + try: + start_frame = int(self.GetStartFrame()) + end_frame = int(self.GetEndFrame()) + indices = list(range(start_frame, end_frame + 1)) + except Exception as e: + self.LogWarning(f"Could not read Deadline task frame range, using task/chunk math: {e}") + indices = indices_from_task + + if len(indices) != len(indices_from_task): + self.LogWarning( + f"Deadline frame range reported {indices}, but task/chunk math expects " + f"{indices_from_task}; using task/chunk math for variation assignment." + ) + indices = indices_from_task + + return [index for index in indices if 0 <= index < self.batch_count] or [0] + def _setup_output_directory(self): """Setup output directory configuration""" job_output_directory_plugin = self.GetPluginInfoEntryWithDefault("JobOutputDirectory", "") @@ -330,6 +367,28 @@ class ComfyUI(DeadlinePlugin): if not os.path.exists(self.comfyui_output_dir): self._create_directory(self.comfyui_output_dir, "default output") + def _setup_input_directory(self): + """Setup optional staged ComfyUI input directory for default loader nodes.""" + input_dir = self.GetPluginInfoEntryWithDefault("JobInputDirectory", "").strip() + self.input_manifest_file = self.GetPluginInfoEntryWithDefault("InputManifestFile", "").strip() + self.submission_id = self.GetPluginInfoEntryWithDefault("SubmissionId", "").strip() + + if not input_dir: + self.comfyui_input_dir = "" + self.LogInfo("No staged input directory specified; ComfyUI will use its default input folder.") + return + + try: + input_dir = RepositoryUtils.CheckPathMapping(input_dir) + except Exception as e: + self.LogWarning(f"Path mapping failed for JobInputDirectory '{input_dir}': {e}") + + self.comfyui_input_dir = os.path.abspath(os.path.expandvars(input_dir)) + if not os.path.isdir(self.comfyui_input_dir): + raise ComfyUIError(f"Staged input directory does not exist: {self.comfyui_input_dir}") + + self.LogInfo(f"ComfyUI will use staged input directory: {self.comfyui_input_dir}") + def _create_directory(self, directory_path: str, description: str): """Create a directory with error handling""" try: @@ -361,40 +420,18 @@ 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): - self.LogInfo(f"ComfyUI is already running on port {base_port}, will use existing instance") - self.use_existing_comfyui = True - self.comfyui_port = str(base_port) - self.comfyui_api_url = f"http://127.0.0.1:{self.comfyui_port}" - self.server_started = True + + self.use_existing_comfyui = False + 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/distributed mode: starting isolated ComfyUI on port {self.comfyui_port}") else: - self.LogInfo(f"No ComfyUI instance detected on port {base_port}") - self.use_existing_comfyui = False self.comfyui_port = self._find_available_port(base_port) - self.LogInfo(f"Will use port {self.comfyui_port} for ComfyUI") - self.comfyui_api_url = f"http://127.0.0.1:{self.comfyui_port}" + self.LogInfo(f"Normal render mode: starting isolated ComfyUI on port {self.comfyui_port}") + + self.comfyui_api_url = f"http://127.0.0.1:{self.comfyui_port}" return self.comfyui_port @@ -427,6 +464,7 @@ class ComfyUI(DeadlinePlugin): if not comfyui_path: raise ComfyUIError("ComfyUI installation path could not be resolved.") self._setup_output_directory() + self._setup_input_directory() self._setup_temp_directory() self._calculate_comfyui_port() self.task_completed = False @@ -608,6 +646,10 @@ print('Dummy command timeout reached or task completed.') else: self.LogInfo("Not passing --output-directory to ComfyUI, it will use its default.") + if self.comfyui_input_dir: + args_list.append(f'--input-directory "{self.comfyui_input_dir}"') + self.LogInfo(f"Passing --input-directory \"{self.comfyui_input_dir}\" to ComfyUI.") + args = " ".join(args_list) self.LogInfo(f"Render Arguments: {args}") return args @@ -756,6 +798,9 @@ print('Dummy command timeout reached or task completed.') def _set_deadline_environment_variables(self): """Set Deadline-specific environment variables for ComfyUI process""" try: + os.environ['GIT_PYTHON_REFRESH'] = 'quiet' + self.LogInfo("Set GIT_PYTHON_REFRESH=quiet so workers do not require git.exe for ComfyUI-Manager startup") + # Get the actual Deadline worker name and set it as environment variable slave_name = self.GetSlaveName() if slave_name: @@ -831,21 +876,19 @@ print('Dummy command timeout reached or task completed.') try: workflow_data = self._load_workflow_from_file(workflow_file) workflow_data = self.validate_workflow(workflow_data) - - # Inject parameters for DeadlineDistributedSeed nodes - deadline_seeds_injected = self.inject_deadline_seed_parameters(workflow_data) - - # 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}") + + worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) + if worker_mode or distributed_mode: + deadline_seeds_injected = self.inject_deadline_seed_parameters(workflow_data) + 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 legacy seed manipulation for distributed task ID {task_id}") else: - self.LogInfo(f"No seed manipulation applied for task ID {task_id}") + self.LogInfo("DeadlineSeed nodes detected for distributed worker workflow") else: - self.LogInfo("DeadlineSeed nodes detected - skipping automatic seed modification") + self.LogInfo("Normal V2 job: seed variation will be applied per queued prompt via DeadlineSeed only") return workflow_data except Exception as e: @@ -865,13 +908,77 @@ print('Dummy command timeout reached or task completed.') # 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()) + # Normal V2 jobs execute the API prompt, while the first auxiliary file may be the standard UI workflow. + workflow_file = self.GetPluginInfoEntryWithDefault("ComfyWorkflowFile", "") + if not workflow_file: + workflow_file = self.GetDataFilename() - workflow_file = RepositoryUtils.CheckPathMapping(workflow_file) + workflow_file = self._resolve_workflow_file_path(workflow_file) self.LogInfo(f"Workflow file setting from plugin info: '{workflow_file}'") return workflow_file + def _resolve_workflow_file_path(self, workflow_file: str) -> str: + """Resolve absolute or auxiliary-file-relative workflow paths.""" + workflow_file = (workflow_file or "").strip().strip('"') + if not workflow_file: + return "" + + try: + mapped = RepositoryUtils.CheckPathMapping(workflow_file) + except Exception: + mapped = workflow_file + + if os.path.isabs(mapped): + return mapped + + candidates = [] + data_filename = self.GetDataFilename() + if data_filename: + candidates.append(os.path.join(os.path.dirname(data_filename), mapped)) + + try: + candidates.append(os.path.join(os.path.dirname(os.path.abspath(__file__)), mapped)) + except Exception: + pass + + if self.temp_dir: + candidates.append(os.path.join(self.temp_dir, mapped)) + + candidates.append(os.path.abspath(mapped)) + + for candidate in candidates: + try: + candidate = RepositoryUtils.CheckPathMapping(candidate) + except Exception: + pass + if os.path.exists(candidate): + return candidate + + self.LogWarning(f"Could not resolve relative workflow file '{workflow_file}'. Tried: {candidates}") + return candidates[0] if candidates else mapped + + def _load_standard_workflow_metadata(self): + """Load the optional UI workflow used for image metadata drag-and-drop.""" + standard_workflow_file = self.GetPluginInfoEntryWithDefault("StandardWorkflowFile", "").strip() + if not standard_workflow_file: + self.standard_workflow = None + return + + try: + standard_workflow_file = self._resolve_workflow_file_path(standard_workflow_file) + if not os.path.exists(standard_workflow_file): + self.LogWarning(f"Standard workflow metadata file not found: {standard_workflow_file}") + self.standard_workflow = None + return + + with open(standard_workflow_file, "r") as f: + self.standard_workflow = json.load(f) + self.standard_workflow_file = standard_workflow_file + self.LogInfo(f"Loaded standard workflow metadata: {standard_workflow_file}") + except Exception as e: + self.standard_workflow = None + self.LogWarning(f"Could not load standard workflow metadata: {e}") + def _load_workflow_from_file(self, workflow_file: str) -> dict: """Load workflow data from JSON file""" with open(workflow_file, 'r') as f: @@ -1034,7 +1141,11 @@ print('Dummy command timeout reached or task completed.') # A matching Triton is not available warning "A matching Triton is not available, some optimizations will not be enabled", # xformers version warnings - "WARNING: You need pytorch with cu130 or higher to use optimized CUDA operations" + "WARNING: You need pytorch with cu130 or higher to use optimized CUDA operations", + # ComfyUI-Manager imports GitPython on startup, but render workers do not need git.exe. + "ImportError: Bad git executable", + "ImportError: Failed to initialize: Bad git executable", + "Cannot import C:\\AI\\ComfyUI_windows_portable4\\ComfyUI\\custom_nodes\\ComfyUI-Manager module for custom nodes: Failed to initialize: Bad git executable", ] for pattern in non_critical_patterns: @@ -1068,7 +1179,7 @@ print('Dummy command timeout reached or task completed.') self.FailRender(f"Error connecting to ComfyUI API: {response['status_code']}") return False - self.client_id = response['json']().get('client_id', '') + self.client_id = f"deadline-{uuid.uuid4().hex}" self.LogInfo(f"Got client ID: {self.client_id}") return True except Exception as e: @@ -1080,15 +1191,13 @@ print('Dummy command timeout reached or task completed.') """Submit workflow to ComfyUI queue""" try: self._reset_prompt_tracking() - - # Queue initial prompt - if not self._queue_single_prompt(workflow_data): - return False - - # Queue additional prompts for batch mode - if self.batch_mode and self.chunk_size > 1: - self._queue_batch_prompts(workflow_data) - + + for variation_index in self.assigned_variation_indices: + prompt_workflow, metadata, workflow_metadata = self._prepare_variation_prompt(workflow_data, variation_index) + if not self._queue_prompt(prompt_workflow, metadata, workflow_metadata): + return False + time.sleep(0.2) + self.LogInfo(f"Queued total of {len(self.prompt_ids)} prompts: {self.prompt_ids}") return True except Exception as e: @@ -1102,9 +1211,19 @@ print('Dummy command timeout reached or task completed.') self.completed_prompts = set() self.current_tracking_index = 0 - def _queue_single_prompt(self, workflow_data: dict) -> bool: - """Queue a single prompt to ComfyUI""" - data = {"prompt": workflow_data, "client_id": self.client_id} + def _queue_prompt(self, workflow_data: dict, deadline_metadata: dict, workflow_metadata: dict = None) -> bool: + """Queue one prepared variation prompt to ComfyUI.""" + extra_pnginfo = {"deadline": deadline_metadata} + if workflow_metadata is not None: + extra_pnginfo["workflow"] = workflow_metadata + + data = { + "prompt": workflow_data, + "client_id": self.client_id, + "extra_data": { + "extra_pnginfo": extra_pnginfo + }, + } response = self.http_request(f"{self.comfyui_api_url}/prompt", method="POST", data=data) if response['status_code'] != 200: @@ -1114,55 +1233,73 @@ print('Dummy command timeout reached or task completed.') self.prompt_id = response['json']()['prompt_id'] self.prompt_ids.append(self.prompt_id) - self.LogInfo(f"Queued prompt with ID: {self.prompt_id}") + self.LogInfo(f"Queued variation {deadline_metadata['variation_index']} with prompt ID: {self.prompt_id}") self.workflow_submitted = True return True - def _queue_batch_prompts(self, workflow_data: dict): - """Queue additional prompts for batch processing""" - self.LogInfo(f"Batch mode with chunk size {self.chunk_size}. Queueing additional prompts...") - - import copy - - for i in range(1, self.chunk_size): - prompt_workflow = copy.deepcopy(workflow_data) - - # 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} - response = self.http_request(f"{self.comfyui_api_url}/prompt", method="POST", data=data) - - if response['status_code'] != 200: - self.LogWarning(f"Error queuing additional prompt {i}: {response['text']}") - break - - prompt_id = response['json']()['prompt_id'] - self.prompt_ids.append(prompt_id) - self.LogInfo(f"Queued additional prompt {i} with ID: {prompt_id}") - - time.sleep(0.5) # Small delay between submissions + def _prepare_variation_prompt(self, workflow_data: dict, variation_index: int): + """Deep-copy and rewrite DeadlineSeed nodes for one global variation index.""" + prompt_workflow = copy.deepcopy(workflow_data) + seeds = [] + + for node_id, node in prompt_workflow.items(): + if not isinstance(node, dict) or node.get("class_type") != "DeadlineSeed": + continue + inputs = node.setdefault("inputs", {}) + base_seed = int(inputs.get("seed", 0)) + actual_seed = base_seed + int(variation_index) + inputs["seed"] = actual_seed + inputs["task_id"] = 0 + inputs["batch_mode"] = False + seeds.append({ + "node_id": str(node_id), + "base_seed": base_seed, + "actual_seed": actual_seed, + }) + self.LogInfo(f"DeadlineSeed node {node_id}: base {base_seed}, variation {variation_index}, actual {actual_seed}") + + metadata = { + "job_id": getattr(self.GetJob(), "JobId", ""), + "task_id": str(self.GetCurrentTaskId()), + "variation_index": int(variation_index), + "submission_id": self.submission_id, + "output_directory": self.comfyui_output_dir, + "input_directory": self.comfyui_input_dir, + "input_manifest": self.input_manifest_file, + "seeds": seeds, + } + if seeds: + metadata["base_seed"] = seeds[0]["base_seed"] + metadata["actual_seed"] = seeds[0]["actual_seed"] + + workflow_metadata = self._prepare_standard_workflow_metadata(seeds) + return prompt_workflow, metadata, workflow_metadata + + def _prepare_standard_workflow_metadata(self, seeds: list): + """Patch the UI workflow metadata so dropped output images reopen the actual variation.""" + if self.standard_workflow is None: + return None + + workflow_metadata = copy.deepcopy(self.standard_workflow) + seed_by_node_id = { + str(seed_info["node_id"]): seed_info["actual_seed"] + for seed_info in seeds + } + + nodes = workflow_metadata.get("nodes", []) + if isinstance(nodes, list): + for node in nodes: + if not isinstance(node, dict): + continue + node_id = str(node.get("id", "")) + node_type = node.get("type") or node.get("class_type") + if node_type != "DeadlineSeed" or node_id not in seed_by_node_id: + continue + widgets = node.get("widgets_values") + if isinstance(widgets, list) and widgets: + widgets[0] = seed_by_node_id[node_id] + + return workflow_metadata def process_history_data(self, history_data: dict) -> bool: """Process history data and update task status""" @@ -1221,16 +1358,7 @@ print('Dummy command timeout reached or task completed.') self.SetProgress(100) self.SetStatusMessage("Finished Render") self.task_completed = True - - # 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") + self.LogInfo(f"All {self.chunk_size} prompt(s) in this task completed") def _move_to_next_prompt(self): """Move to tracking the next prompt""" @@ -1263,17 +1391,8 @@ print('Dummy command timeout reached or task completed.') """Handle prompt execution errors""" error_msg = status.get('error', 'Unknown error') self.LogWarning(f"ComfyUI reported error for prompt {self.prompt_id}: {error_msg}") - - if self.chunk_size > 1: - # Continue with remaining prompts in batch - self.LogWarning(f"Continuing with remaining prompts in chunk") - self.completed_prompts.add(self.prompt_id) - self._move_to_next_prompt() - return False - else: - # Single prompt mode, fail the task - self.FailRender(f"ComfyUI workflow failed: {error_msg}") - return True + self.FailRender(f"ComfyUI workflow failed: {error_msg}") + return True def _update_execution_progress(self, progress: float): """Update progress from execution information""" @@ -1320,7 +1439,6 @@ print('Dummy command timeout reached or task completed.') def monitor_workflow_execution(self) -> bool: """Poll history endpoint and wait for workflow completion""" - start_time = time.time() self.LogInfo(f"Beginning to monitor workflow execution for chunk size {self.chunk_size}") self.LogInfo(f"Monitoring prompts in this order: {self.prompt_ids}") @@ -1329,9 +1447,9 @@ print('Dummy command timeout reached or task completed.') poll_count = 0 - while time.time() - start_time < DEFAULT_TIMEOUT and self.thread_running: + while self.thread_running: if self.task_completed: - self.LogInfo("Task already marked as complete by stdout handler") + self.LogInfo("Task already marked as complete") # Check if we're in distributed worker mode worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) @@ -1352,12 +1470,6 @@ print('Dummy command timeout reached or task completed.') poll_count += 1 time.sleep(DEFAULT_POLLING_INTERVAL) - # Check for timeout - if not self.task_completed and time.time() - start_time >= DEFAULT_TIMEOUT: - self.LogWarning(f"Timeout waiting for workflow to complete after {DEFAULT_TIMEOUT} seconds") - self.FailRender(f"Timeout waiting for workflow to complete") - return False - if self.task_completed: # Check if we're in distributed worker mode worker_mode, distributed_mode, force_new_instance = get_distributed_config_for_plugin(self) @@ -1460,6 +1572,8 @@ print('Dummy command timeout reached or task completed.') workflow_data = self.load_and_validate_workflow() if not workflow_data: return + + self._load_standard_workflow_metadata() if not self.initialize_api_connection(): return @@ -1472,4 +1586,4 @@ print('Dummy command timeout reached or task completed.') except Exception as e: self.LogWarning(f"Error during workflow submission: {e}") traceback.print_exc() - self.FailRender(f"Error during workflow submission: {str(e)}") \ No newline at end of file + self.FailRender(f"Error during workflow submission: {str(e)}") diff --git a/pyproject.toml b/pyproject.toml index bbcc59e..8625616 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -33,7 +33,7 @@ Issues = "https://github.com/doubletwisted/ComfyUI-Deadline-Plugin/issues" [tool.setuptools] # For flat structure with .py files in root -py-modules = ["deadline_submit", "deadline_api", "__init__"] +py-modules = ["deadline_submit", "__init__"] [comfyui] PublisherId = "doubletwisted" diff --git a/scripts/copy_comfy_to_network.ps1 b/scripts/copy_comfy_to_network.ps1 index ac0889d..1f90b35 100644 --- a/scripts/copy_comfy_to_network.ps1 +++ b/scripts/copy_comfy_to_network.ps1 @@ -7,6 +7,8 @@ $source = "C:\AI\ComfyUI_windows_portable4" $destination = "X:\AI\ComfyUI_windows_portable4" +$publishLock = "$destination.sync_in_progress" + $logDir = "X:\scripts\copy_comfy_to_network_logs" # Get current date and time in YYYY-MM-DD_HH-MM-SS format @@ -27,6 +29,15 @@ if (!(Test-Path -Path $logDir)) { Start-Transcript -Path $logFile -Append +function Assert-RobocopySuccess { + param([string]$Operation) + + $exitCode = $LASTEXITCODE + if ($exitCode -ge 8) { + throw "Robocopy failed during $Operation with exit code $exitCode" + } +} + # Create destination if it doesn't exist if (!(Test-Path -Path $destination)) { @@ -35,6 +46,10 @@ if (!(Test-Path -Path $destination)) { } +Set-Content -Path $publishLock -Value "Publishing from $env:COMPUTERNAME at $(Get-Date -Format o)" -Force + +try { + # --- Extra exclusions for robocopy --- $extraExcludeDirs = @() @@ -42,13 +57,21 @@ $extraExcludeDirs = @() # Precise absolute excludes for specific Triton and SageAttention folders $extraExcludeDirs += "$source\python_embeded\Lib\site-packages\triton" $extraExcludeDirs += "$source\python_embeded\Lib\site-packages\triton-3.2.0.dist-info" +$extraExcludeDirs += "$source\python_embeded\Lib\site-packages\sageattention" +$extraExcludeDirs += "$source\python_embeded\Lib\site-packages\sageattention-2.1.1.dist-info" $extraExcludeDirs += "$source\SageAttention" -$allExcludeDirs = @('__pycache__', 'output', 'input') + $extraExcludeDirs +$allExcludeDirs = @( + '__pycache__', + "$source\ComfyUI\input", + "$source\ComfyUI\output", + "$source\ComfyUI\temp" +) + $extraExcludeDirs # Main robocopy with all exclusions, mirror mode, and junction exclusion robocopy $source $destination /MIR /XD $allExcludeDirs /XF "*.md5" "*.log" "*.tmp" /XJ /R:5 /W:5 /NFL /NDL /NP +Assert-RobocopySuccess "main mirror" # Copy the specific example.png file @@ -66,7 +89,7 @@ if (Test-Path $exampleFile) { } - Copy-Item $exampleFile $inputDestination -Force + Copy-Item $exampleFile $inputDestination -Force -ErrorAction Stop } @@ -79,12 +102,14 @@ if (Test-Path $inputSource) { # Copy only subdirectories and their contents from input folder robocopy $inputSource $inputDestination /S /XF * /R:2 /W:2 /NFL /NDL /NP + Assert-RobocopySuccess "input folder structure" # Then copy contents of subdirectories Get-ChildItem $inputSource -Directory | ForEach-Object { robocopy $_.FullName "$inputDestination\$($_.Name)" /E /R:2 /W:2 /NFL /NDL /NP + Assert-RobocopySuccess "input subfolder $($_.FullName)" } @@ -101,6 +126,7 @@ if (Test-Path $outputSource) { # Copy directory structure without files robocopy $outputSource $outputDestination /E /XF * /R:2 /W:2 /NFL /NDL /NP + Assert-RobocopySuccess "output folder structure" } @@ -112,6 +138,13 @@ Write-Host "- Copied all input subfolders with contents" Write-Host "- Created empty output folder structure" -# Stop transcript +} +finally { -Stop-Transcript + Remove-Item -LiteralPath $publishLock -Force -ErrorAction SilentlyContinue + + # Stop transcript + + Stop-Transcript + +} diff --git a/scripts/deploy_deadline_plugin.ps1 b/scripts/deploy_deadline_plugin.ps1 new file mode 100644 index 0000000..37611e1 --- /dev/null +++ b/scripts/deploy_deadline_plugin.ps1 @@ -0,0 +1,181 @@ +<# +.SYNOPSIS +Deploys this repository's Deadline render plugin to the configured Deadline Repository. + +.DESCRIPTION +Finds the Deadline Repository path via deadlinecommand -GetRepositoryPath, then copies +plugins\ComfyUI into \custom\plugins\ComfyUI. + +You can also pass -RepositoryPath explicitly. The value may be either the repository +root or the custom plugins directory. + +.EXAMPLE +powershell.exe -ExecutionPolicy Bypass -File .\scripts\deploy_deadline_plugin.ps1 + +.EXAMPLE +powershell.exe -ExecutionPolicy Bypass -File .\scripts\deploy_deadline_plugin.ps1 -RepositoryPath "\\YOUR-SERVER\Repository" + +.EXAMPLE +powershell.exe -ExecutionPolicy Bypass -File .\scripts\deploy_deadline_plugin.ps1 -RepositoryPath "\\YOUR-SERVER\Repository\custom\plugins" +#> + +[CmdletBinding(SupportsShouldProcess = $true)] +param( + [string]$RepositoryPath = "", + [string]$DeadlineCommand = "", + [string]$PluginName = "ComfyUI", + [switch]$Clean, + [switch]$NoBackup +) + +$ErrorActionPreference = "Stop" + +function Write-Step { + param([string]$Message) + Write-Host "[deploy-deadline-plugin] $Message" +} + +function Resolve-DeadlineCommand { + param([string]$ExplicitPath) + + $candidates = @() + if ($ExplicitPath) { + $candidates += $ExplicitPath + } + if ($env:DEADLINE_PATH) { + $candidates += (Join-Path $env:DEADLINE_PATH "deadlinecommand.exe") + $candidates += (Join-Path $env:DEADLINE_PATH "deadlinecommand") + } + $candidates += "C:\Program Files\Thinkbox\Deadline10\bin\deadlinecommand.exe" + $candidates += "/opt/Thinkbox/Deadline10/bin/deadlinecommand" + + foreach ($candidate in $candidates) { + if ($candidate -and (Test-Path -LiteralPath $candidate)) { + return (Resolve-Path -LiteralPath $candidate).Path + } + } + + throw "Could not find deadlinecommand. Pass -DeadlineCommand or set DEADLINE_PATH." +} + +function Get-DeadlineRepositoryPath { + param( + [string]$ExplicitRepositoryPath, + [string]$DeadlineCommandPath + ) + + if ($ExplicitRepositoryPath) { + return $ExplicitRepositoryPath.Trim().Trim('"') + } + + foreach ($argument in @("-GetRepositoryPath", "-GetRepositoryRoot")) { + $output = & $DeadlineCommandPath $argument 2>&1 + if ($LASTEXITCODE -eq 0) { + $path = (@($output) | ForEach-Object { "$_" } | Where-Object { ![string]::IsNullOrWhiteSpace($_) } | Select-Object -First 1) + if ($path) { + $path = $path.ToString().Trim() + } + if ($path -and $path -notmatch "Bad submission arguments") { + return $path + } + } + } + + throw "deadlinecommand did not return a repository path." +} + +function Resolve-CustomPluginsPath { + param([string]$RepositoryOrPluginsPath) + + $path = $RepositoryOrPluginsPath.Trim().Trim('"') + if ([string]::IsNullOrWhiteSpace($path)) { + throw "Repository path is empty." + } + + $leaf = Split-Path -Leaf $path + $parent = Split-Path -Parent $path + $parentLeaf = if ($parent) { Split-Path -Leaf $parent } else { "" } + + if ($leaf -ieq "plugins" -and $parentLeaf -ieq "custom") { + return $path + } + + return (Join-Path $path "custom\plugins") +} + +function Copy-PluginDirectory { + param( + [string]$SourceDirectory, + [string]$DestinationDirectory, + [bool]$ShouldClean + ) + + if (!(Test-Path -LiteralPath $SourceDirectory)) { + throw "Source plugin directory does not exist: $SourceDirectory" + } + + $destinationParent = Split-Path -Parent $DestinationDirectory + if (!(Test-Path -LiteralPath $destinationParent)) { + if ($PSCmdlet.ShouldProcess($destinationParent, "Create custom plugins directory")) { + New-Item -ItemType Directory -Path $destinationParent -Force | Out-Null + } + } + + if ($ShouldClean -and (Test-Path -LiteralPath $DestinationDirectory)) { + $resolvedDestination = (Resolve-Path -LiteralPath $DestinationDirectory).Path + if ($resolvedDestination -match "\\custom\\plugins\\[^\\]+$") { + if ($PSCmdlet.ShouldProcess($resolvedDestination, "Remove existing plugin directory before deploy")) { + Remove-Item -LiteralPath $resolvedDestination -Recurse -Force + } + } + else { + throw "Refusing to clean unexpected destination path: $resolvedDestination" + } + } + + if (!(Test-Path -LiteralPath $DestinationDirectory)) { + if ($PSCmdlet.ShouldProcess($DestinationDirectory, "Create plugin directory")) { + New-Item -ItemType Directory -Path $DestinationDirectory -Force | Out-Null + } + } + + if ($PSCmdlet.ShouldProcess($DestinationDirectory, "Copy Deadline plugin files from $SourceDirectory")) { + robocopy $SourceDirectory $DestinationDirectory /E /XF "*.pyc" /XD "__pycache__" /R:3 /W:2 /NFL /NDL /NP + $exitCode = $LASTEXITCODE + if ($exitCode -gt 7) { + throw "robocopy failed with exit code $exitCode" + } + } +} + +$deadlineCommandPath = Resolve-DeadlineCommand -ExplicitPath $DeadlineCommand +$repositoryPathResolved = Get-DeadlineRepositoryPath -ExplicitRepositoryPath $RepositoryPath -DeadlineCommandPath $deadlineCommandPath +$customPluginsPath = Resolve-CustomPluginsPath -RepositoryOrPluginsPath $repositoryPathResolved + +$scriptDirectory = Split-Path -Parent $MyInvocation.MyCommand.Path +$projectRoot = Split-Path -Parent $scriptDirectory +$sourcePluginDirectory = Join-Path $projectRoot "plugins\$PluginName" +$destinationPluginDirectory = Join-Path $customPluginsPath $PluginName + +Write-Step "deadlinecommand: $deadlineCommandPath" +Write-Step "Repository/custom plugins path: $customPluginsPath" +Write-Step "Source: $sourcePluginDirectory" +Write-Step "Destination: $destinationPluginDirectory" + +if (!(Test-Path -LiteralPath $customPluginsPath)) { + throw "Custom plugins path does not exist or is not reachable: $customPluginsPath" +} + +if (!$NoBackup -and (Test-Path -LiteralPath $destinationPluginDirectory)) { + $timestamp = Get-Date -Format "yyyyMMdd_HHmmss" + $backupDirectory = Join-Path $customPluginsPath "$PluginName.backup_$timestamp" + if ($PSCmdlet.ShouldProcess($backupDirectory, "Create backup of current $PluginName Deadline plugin")) { + Copy-Item -LiteralPath $destinationPluginDirectory -Destination $backupDirectory -Recurse -Force + Write-Step "Backup created: $backupDirectory" + } +} + +Copy-PluginDirectory -SourceDirectory $sourcePluginDirectory -DestinationDirectory $destinationPluginDirectory -ShouldClean ([bool]$Clean) + +Write-Step "Deploy complete." +Write-Step "If workers have cached the old plugin, restart them or submit a new job so Deadline copies the updated plugin sandbox." diff --git a/scripts/maintenance/ComfyModelsSync.py b/scripts/maintenance/ComfyModelsSync.py new file mode 100644 index 0000000..23beeea --- /dev/null +++ b/scripts/maintenance/ComfyModelsSync.py @@ -0,0 +1,440 @@ +from __future__ import absolute_import, print_function + +r""" +ComfyUI models sync maintenance script. + +Syncs AI models from shared storage to a local drive based on computer name. +- Computers in E_DRIVE_COMPUTERS copy to E:\AI\models +- Computers in D_DRIVE_COMPUTERS copy to D:\AI\models +- All other computers copy to C:\AI\models +Uses the same model list (modellist.txt) for all computers. + +Features: +- Computer-specific drive selection +- Delta-aware space checking with 15% safety buffer +- Delta-aware copying (only copies changed/missing files) +- Cleans up files not in the model list +- Uses Robocopy for reliable file operations +- Comprehensive logging for Deadline integration + +Designed to be launched on Deadline Workers via the CommandLine or +DeadlineCommand plugin, +so all stdout/stderr (or Deadline logging if available) ends up in job logs. +""" + +import os +import shutil +import subprocess +import sys +from datetime import datetime + +# Configuration +SOURCE_DIR = os.environ.get("COMFY_MODELS_SOURCE", r"\\YOUR-SERVER\share\AI\models") +MODEL_LIST_PATH = os.environ.get("COMFY_MODEL_LIST", r"\\YOUR-SERVER\share\scripts\modellist.txt") + +# Drive configuration by computer +# Add Worker hostnames here for site-specific drive placement. +D_DRIVE_COMPUTERS = set() +E_DRIVE_COMPUTERS = set() +ROBOCOPY_EXECUTABLE = "robocopy" +ROBOCOPY_FLAGS = [ + "/R:10", # Retry 10 times + "/W:30", # Wait 30 seconds between retries + "/V", # Verbose output + "/TS", # Include source time stamps + "/FP", # Include full path names + "/NP", # No progress indicator + "/MT:1" # Single-threaded for reliability +] + +try: + from Deadline.Scripting import ClientUtils # type: ignore +except Exception: # pragma: no cover - Deadline libs unavailable outside Worker + ClientUtils = None + + +def log(message): + """Log to Deadline if available, otherwise stdout.""" + timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") + line = "[ComfyModelsSync] {0} {1}".format(timestamp, message) + if ClientUtils: + ClientUtils.LogText(line) + else: + print(line) + + +def get_drive_space_info(dest_dir): + """Get available space information for destination drive.""" + import shutil + + try: + # Get drive letter and root path + drive_letter = dest_dir.split(':')[0].upper() + drive_root = "{0}:\\".format(drive_letter) + + # Get disk usage stats using cross-platform method + total_bytes, used_bytes, free_bytes = shutil.disk_usage(drive_root) + + free_gb = free_bytes / (1024**3) + total_gb = total_bytes / (1024**3) + + return { + 'free_bytes': free_bytes, + 'free_gb': round(free_gb, 2), + 'total_gb': round(total_gb, 2) + } + except Exception as e: + log("WARNING: Could not get drive space info: {0}".format(e)) + return None + + +def calculate_required_space(valid_model_paths, dest_dir): + """Calculate space required for files that need copying (delta-aware).""" + total_delta_bytes = 0 + files_to_copy = [] + + log("Analyzing space requirements (delta-aware)...") + + for source_path in valid_model_paths: + # Calculate destination path + if source_path.startswith(SOURCE_DIR): + relative_path = source_path[len(SOURCE_DIR):].lstrip(os.sep) + else: + relative_path = os.path.basename(source_path) + + dest_path = os.path.join(dest_dir, relative_path) + + # Check if file needs copying + needs_copy_flag, _ = needs_copy(source_path, dest_path) + + if needs_copy_flag: + try: + file_size = os.path.getsize(source_path) + total_delta_bytes += file_size + files_to_copy.append(source_path) + except OSError: + log("WARNING: Could not get size for {0}".format(source_path)) + + total_delta_gb = round(total_delta_bytes / (1024**3), 2) + + log("Space analysis complete:") + log(" Files requiring copy: {0}".format(len(files_to_copy))) + log(" Total space required (delta): {0} GB".format(total_delta_gb)) + + return { + 'total_delta_bytes': total_delta_bytes, + 'total_delta_gb': total_delta_gb, + 'files_to_copy': files_to_copy + } + + +def get_computer_config(): + """Determine destination drive based on computer name.""" + import socket + + # Get computer name + computer_name = socket.gethostname().upper() + log("Running on computer: {0}".format(computer_name)) + + # Determine drive + if computer_name in E_DRIVE_COMPUTERS: + dest_drive = "E" + elif computer_name in D_DRIVE_COMPUTERS: + dest_drive = "D" + else: + dest_drive = "C" + + dest_dir = "{0}:\\AI\\models".format(dest_drive) + + log("Using destination drive: {0} (directory: {1})".format(dest_drive, dest_dir)) + + return dest_dir + + +def read_model_list(): + """Read and parse the model list file.""" + if not os.path.exists(MODEL_LIST_PATH): + raise RuntimeError("Model list file not found: {0}".format(MODEL_LIST_PATH)) + + with open(MODEL_LIST_PATH, 'r') as f: + lines = f.readlines() + + # Parse paths, skip empty lines + model_paths = [] + for line in lines: + line = line.strip() + if line and not line.startswith('#'): + model_paths.append(line) + + log("Loaded {0} model paths from list".format(len(model_paths))) + return model_paths + + +def ensure_destination(): + """Ensure destination directory exists.""" + if not os.path.isdir(DEST_DIR): + log("Destination directory does not exist, creating: {0}".format(DEST_DIR)) + os.makedirs(DEST_DIR) + + +def get_file_size_mb(file_path): + """Get file size in MB.""" + try: + size_bytes = os.path.getsize(file_path) + return size_bytes / (1024 * 1024) + except OSError: + return 0 + + +def analyze_models(model_paths): + """Analyze model files and return valid/missing lists.""" + valid_files = [] + missing_files = [] + + for model_path in model_paths: + if os.path.exists(model_path): + valid_files.append(model_path) + else: + missing_files.append(model_path) + log("WARNING: Model file not found: {0}".format(model_path)) + + total_size_mb = sum(get_file_size_mb(path) for path in valid_files) + + log("Analysis complete:") + log(" Total models in list: {0}".format(len(model_paths))) + log(" Valid source files: {0}".format(len(valid_files))) + log(" Missing source files: {0}".format(len(missing_files))) + log(" Total size of valid files: {0:.2f} MB".format(total_size_mb)) + + return valid_files, missing_files + + +def cleanup_extra_files(valid_model_paths): + """Remove files from destination that are not in the model list.""" + if not os.path.exists(DEST_DIR): + log("Destination directory does not exist yet - no cleanup needed") + return 0, 0 + + # Build set of expected relative paths + expected_relative_paths = set() + for model_path in valid_model_paths: + if model_path.startswith(SOURCE_DIR): + relative_path = model_path[len(SOURCE_DIR):].lstrip(os.sep) + expected_relative_paths.add(relative_path) + + # Find all files in destination + removed_count = 0 + removed_size_mb = 0 + + for root, dirs, files in os.walk(DEST_DIR): + for file in files: + full_path = os.path.join(root, file) + relative_path = os.path.relpath(full_path, DEST_DIR) + + if relative_path not in expected_relative_paths: + try: + size_mb = get_file_size_mb(full_path) + os.remove(full_path) + log("Removed extra file: {0} ({1:.2f} MB)".format(relative_path, size_mb)) + removed_count += 1 + removed_size_mb += size_mb + except OSError as e: + log("WARNING: Failed to remove {0}: {1}".format(relative_path, e)) + + # Clean up empty directories + for root, dirs, files in os.walk(DEST_DIR, topdown=False): + for dir_name in dirs: + dir_path = os.path.join(root, dir_name) + try: + if not os.listdir(dir_path): + os.rmdir(dir_path) + log("Removed empty directory: {0}".format(os.path.relpath(dir_path, DEST_DIR))) + except OSError: + pass # Directory not empty or other error + + if removed_count > 0: + log("Cleanup completed: {0} files removed ({1:.2f} MB freed)".format(removed_count, removed_size_mb)) + else: + log("Cleanup completed: No extra files found") + + return removed_count, removed_size_mb + + +def needs_copy(source_path, dest_path): + """Check if file needs to be copied.""" + if not os.path.exists(dest_path): + return True, "missing" + + try: + source_stat = os.stat(source_path) + dest_stat = os.stat(dest_path) + + if source_stat.st_mtime > dest_stat.st_mtime: + return True, "newer" + elif source_stat.st_size != dest_stat.st_size: + return True, "different size" + else: + return False, "up to date" + except OSError: + return True, "error checking" + + +def copy_model_file(source_path, dest_path): + """Copy a model file using robocopy.""" + source_dir = os.path.dirname(source_path) + dest_dir = os.path.dirname(dest_path) + file_name = os.path.basename(source_path) + + # Ensure destination directory exists (including subdirectories) + os.makedirs(dest_dir, exist_ok=True) + + cmd = [ROBOCOPY_EXECUTABLE, source_dir, dest_dir, file_name] + ROBOCOPY_FLAGS + + log("Running: {0}".format(" ".join(cmd))) + process = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True) + stdout, stderr = process.communicate() + + if process.returncode >= 8: + error_msg = "Robocopy failed with exit code {0}".format(process.returncode) + if stdout: + error_msg += "\nSTDOUT: {0}".format(stdout.strip()) + if stderr: + error_msg += "\nSTDERR: {0}".format(stderr.strip()) + raise RuntimeError(error_msg) + + return process.returncode + + +def sync_models(valid_model_paths): + """Sync all valid model files.""" + copied_count = 0 + skipped_count = 0 + error_count = 0 + total_processed = 0 + + for source_path in valid_model_paths: + total_processed += 1 + + # Calculate destination path + if source_path.startswith(SOURCE_DIR): + relative_path = source_path[len(SOURCE_DIR):].lstrip(os.sep) + else: + relative_path = os.path.basename(source_path) + + dest_path = os.path.join(DEST_DIR, relative_path) + + # Check if copy is needed + needs_copy_flag, reason = needs_copy(source_path, dest_path) + + if needs_copy_flag: + try: + size_mb = get_file_size_mb(source_path) + log("Copying ({0}): {1} ({2:.2f} MB)".format(reason, relative_path, size_mb)) + + exit_code = copy_model_file(source_path, dest_path) + + if exit_code <= 7: # Success + log("SUCCESS: {0} (exit code: {1})".format(relative_path, exit_code)) + copied_count += 1 + else: + log("WARNING: {0} completed with exit code {1}".format(relative_path, exit_code)) + copied_count += 1 + + except Exception as e: + log("ERROR: Failed to copy {0}: {1}".format(relative_path, e)) + error_count += 1 + else: + log("Skipped (up to date): {0}".format(relative_path)) + skipped_count += 1 + + # Progress update + if total_processed % 10 == 0 or total_processed == len(valid_model_paths): + percent = (total_processed * 100) / len(valid_model_paths) + log("Progress: {0:.1f}% ({1}/{2}) | Copied: {3} | Skipped: {4} | Errors: {5}".format( + percent, total_processed, len(valid_model_paths), copied_count, skipped_count, error_count)) + + return copied_count, skipped_count, error_count + + +def main(): + log("ComfyUI models sync starting.") + log("Source directory: {0}".format(SOURCE_DIR)) + log("Model list: {0}".format(MODEL_LIST_PATH)) + + try: + # Get computer-specific configuration + dest_dir = get_computer_config() + + # Make dest_dir available globally for other functions + global DEST_DIR + DEST_DIR = dest_dir + + # Read model list + model_paths = read_model_list() + + # Analyze models + valid_model_paths, missing_files = analyze_models(model_paths) + + # Perform space checking (delta-aware) + space_info = calculate_required_space(valid_model_paths, dest_dir) + drive_info = get_drive_space_info(dest_dir) + + if drive_info and space_info: + # Apply 15% safety buffer like PowerShell script + buffer_multiplier = 1.15 + required_with_buffer = int(space_info['total_delta_bytes'] * buffer_multiplier) + required_with_buffer_gb = round(required_with_buffer / (1024**3), 2) + + log("Space check:") + log(" Required (with 15% buffer): {0} GB".format(required_with_buffer_gb)) + log(" Available: {0} GB".format(drive_info['free_gb'])) + + if drive_info['free_bytes'] < required_with_buffer: + shortfall_gb = round((required_with_buffer - drive_info['free_bytes']) / (1024**3), 2) + log("ERROR: Insufficient disk space!") + log(" Shortfall: {0} GB".format(shortfall_gb)) + raise RuntimeError("Insufficient disk space for model sync") + else: + log("Space check PASSED") + else: + log("WARNING: Could not perform space check - proceeding anyway") + + # Ensure destination exists + ensure_destination() + + # Cleanup extra files + cleanup_extra_files(valid_model_paths) + + # Sync models + copied_count, skipped_count, error_count = sync_models(valid_model_paths) + + # Final summary + log("Sync completed:") + log(" Models in list: {0}".format(len(model_paths))) + log(" Valid source files: {0}".format(len(valid_model_paths))) + log(" Missing source files: {0}".format(len(missing_files))) + log(" Files copied: {0}".format(copied_count)) + log(" Files skipped: {0}".format(skipped_count)) + log(" Files with errors: {0}".format(error_count)) + + if error_count > 0: + log("WARNING: Some files failed to copy") + sys.exit(1) + + except Exception as e: + log("ERROR: {0}".format(e)) + sys.exit(1) + + log("ComfyUI models sync finished successfully.") + + +if __name__ == "__main__": + try: + main() + except Exception as exc: # pragma: no cover - runtime safeguard + log("ERROR: {0}".format(exc)) + raise + + +def __main__(*args): # Deadline's ExecuteScript entry point + main() diff --git a/scripts/maintenance/ComfyUISync.py b/scripts/maintenance/ComfyUISync.py new file mode 100644 index 0000000..72d84f2 --- /dev/null +++ b/scripts/maintenance/ComfyUISync.py @@ -0,0 +1,176 @@ +from __future__ import absolute_import, print_function + +r""" +ComfyUI portable installation sync maintenance script. + +Mirrors a shared ComfyUI portable install to a local Worker path +using Robocopy /MIR, excluding user-specific input/ and output/ folders so +their contents are never erased. `.git` is included so workers can keep +the same Git metadata as the share. Only ComfyUI/input/example.png is copied separately +(as required for ComfyUI to work properly). + +Designed to be launched on Deadline Workers via the CommandLine or +DeadlineCommand plugin, +so all stdout/stderr (or Deadline logging if available) ends up in job logs. +""" + +import os +import subprocess +import sys +import time +from datetime import datetime + +SOURCE_DIR = os.environ.get("COMFY_SYNC_SOURCE", r"\\YOUR-SERVER\share\AI\ComfyUI_windows_portable") +DEST_DIR = os.environ.get("COMFY_SYNC_DEST", r"C:\AI\ComfyUI_windows_portable") +ROBOCOPY_EXECUTABLE = "robocopy" +ROBOCOPY_FLAGS = [ + "/MIR", # Mirror source to destination (adds + deletes) + "/FFT", # Assume FAT file times (two-second granularity) for cross-protocol copies + "/Z", # Restartable mode + "/R:3", # Retry failed copies 3 times + "/W:5", # Wait 5 seconds between retries + "/NFL", # No file list (keeps logs smaller) + "/NDL", # No directory list +] +EXCLUDED_FILE_PATTERNS = ["*.log", "*.tmp"] +PUBLISH_LOCK_FILE = SOURCE_DIR + ".sync_in_progress" +PUBLISH_LOCK_POLL_SECONDS = 10 +PUBLISH_LOCK_TIMEOUT_SECONDS = 60 * 60 +PUBLISH_LOCK_STALE_SECONDS = 6 * 60 * 60 +EXCLUDED_DIRS = [ + "__pycache__", + "node_modules", + "logs", + os.path.join(SOURCE_DIR, "ComfyUI", "input"), + os.path.join(SOURCE_DIR, "ComfyUI", "output"), +] + +try: + from Deadline.Scripting import ClientUtils # type: ignore +except Exception: # pragma: no cover - Deadline libs unavailable outside Worker + ClientUtils = None + + +def log(message): + """Log to Deadline if available, otherwise stdout.""" + timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S") + line = "[ComfyUISync] {0} {1}".format(timestamp, message) + if ClientUtils: + ClientUtils.LogText(line) + else: + print(line) + + +def ensure_destination(): + """Ensure destination directory exists.""" + if not os.path.isdir(DEST_DIR): + log("Destination directory does not exist, creating: {0}".format(DEST_DIR)) + os.makedirs(DEST_DIR) + + +def wait_for_publish_lock(): + """Avoid copying the share while it is being refreshed from the source machine.""" + start_time = time.time() + while os.path.exists(PUBLISH_LOCK_FILE): + try: + lock_age = time.time() - os.path.getmtime(PUBLISH_LOCK_FILE) + except OSError: + continue + + if lock_age >= PUBLISH_LOCK_STALE_SECONDS: + log( + "Ignoring stale publish lock older than {0} seconds: {1}".format( + PUBLISH_LOCK_STALE_SECONDS, + PUBLISH_LOCK_FILE, + ) + ) + return + + elapsed = time.time() - start_time + if elapsed >= PUBLISH_LOCK_TIMEOUT_SECONDS: + raise RuntimeError( + "Timed out waiting for publish lock to clear: {0}".format( + PUBLISH_LOCK_FILE + ) + ) + + log("Publish lock present, waiting before sync: {0}".format(PUBLISH_LOCK_FILE)) + time.sleep(PUBLISH_LOCK_POLL_SECONDS) + + +def mirror_comfyui(): + """Run robocopy mirror operation.""" + cmd = [ROBOCOPY_EXECUTABLE, SOURCE_DIR, DEST_DIR] + ROBOCOPY_FLAGS + if EXCLUDED_FILE_PATTERNS: + cmd.append("/XF") + cmd.extend(EXCLUDED_FILE_PATTERNS) + if EXCLUDED_DIRS: + cmd.append("/XD") + cmd.extend(EXCLUDED_DIRS) + log("Running command: {0}".format(" ".join(cmd))) + process = subprocess.Popen(cmd) + process.wait() + exit_code = process.returncode + + # Robocopy exit codes < 8 are success (0=No Change, 1=Copied, etc.) + if exit_code >= 8: + msg = "Robocopy failed with exit code {0}".format(exit_code) + log(msg) + raise RuntimeError(msg) + + log("Robocopy completed successfully with exit code {0}".format(exit_code)) + + +def ensure_example_png(): + """Copy ComfyUI/input/example.png from source to destination. Required for ComfyUI to work properly.""" + src_file = os.path.join(SOURCE_DIR, "ComfyUI", "input", "example.png") + dest_dir = os.path.join(DEST_DIR, "ComfyUI", "input") + + if not os.path.isfile(src_file): + log("Source example.png not found, skipping: {0}".format(src_file)) + return + + if not os.path.isdir(dest_dir): + log("Creating input directory: {0}".format(dest_dir)) + os.makedirs(dest_dir) + + cmd = [ + ROBOCOPY_EXECUTABLE, + os.path.join(SOURCE_DIR, "ComfyUI", "input"), + dest_dir, + "example.png", + "/NFL", "/NDL", + ] + log("Ensuring input/example.png exists: {0}".format(" ".join(cmd))) + process = subprocess.Popen(cmd) + process.wait() + exit_code = process.returncode + if exit_code >= 8: + log("Warning: Failed to copy example.png (exit code {0})".format(exit_code)) + else: + log("example.png ensured successfully.") + + +def main(): + log("ComfyUI portable installation sync starting.") + log("Source: {0}".format(SOURCE_DIR)) + log("Destination: {0}".format(DEST_DIR)) + + wait_for_publish_lock() + ensure_destination() + mirror_comfyui() + ensure_example_png() + + log("ComfyUI portable installation sync finished.") + + +if __name__ == "__main__": + try: + main() + except Exception as exc: # pragma: no cover - runtime safeguard + log("ERROR: {0}".format(exc)) + raise + + +def __main__(*args): # Deadline's ExecuteScript entry point + main() diff --git a/scripts/maintenance/README.md b/scripts/maintenance/README.md new file mode 100644 index 0000000..9780bf8 --- /dev/null +++ b/scripts/maintenance/README.md @@ -0,0 +1,39 @@ +# Deadline Maintenance Scripts + +These scripts are publishable templates for syncing ComfyUI installs and model files through Deadline maintenance jobs. + +## Scripts + +- `submit_comfy_sync.py` submits maintenance jobs using Deadline's `DeadlineCommand` plugin. +- `ComfyUISync.py` mirrors a shared ComfyUI portable install to local Worker storage. +- `ComfyModelsSync.py` syncs model files from a shared model list to local Worker storage. + +## Typical Setup + +1. Copy this folder to shared storage visible to all Workers. +2. Edit the default paths in the scripts, or set environment variables on Workers. +3. Submit maintenance jobs from a machine with `deadlinecommand` available. + +```powershell +python \\YOUR-SERVER\share\scripts\maintenance\submit_comfy_sync.py --type both +``` + +## Useful Overrides + +`ComfyUISync.py`: + +- `COMFY_SYNC_SOURCE`: shared ComfyUI portable source path. +- `COMFY_SYNC_DEST`: local Worker destination path. + +`ComfyModelsSync.py`: + +- `COMFY_MODELS_SOURCE`: shared model source path. +- `COMFY_MODEL_LIST`: text file containing model paths to sync. + +`submit_comfy_sync.py`: + +- `--type installation`: submit only ComfyUI install sync. +- `--type models`: submit only model sync. +- `--type both`: submit both jobs. +- `--allowlist`: comma-separated Worker allowlist. +- `--pool`, `--group`, `--region`: Deadline routing options. diff --git a/scripts/maintenance/submit_comfy_sync.py b/scripts/maintenance/submit_comfy_sync.py new file mode 100644 index 0000000..fab9a9a --- /dev/null +++ b/scripts/maintenance/submit_comfy_sync.py @@ -0,0 +1,348 @@ +#!/usr/bin/env python +from __future__ import absolute_import, print_function + +r""" +Utility script that submits ComfyUI sync jobs as Deadline maintenance jobs. + +Supports two types of sync: +1. ComfyUI portable installation sync (mirrors a shared ComfyUI portable install to local Workers) +2. ComfyUI models sync (syncs AI models from shared storage to local Worker drives, using modellist.txt) + +Defaults are publication-safe examples: +- Pool none / group none / no region label +- Maintenance flag enabled +- Empty machine allow list by default +- DeadlineCommand plugin executes sync scripts via deadlinecommand -ExecuteScript + +Usage examples: + # Submit both ComfyUI installation and models sync + python maintenance/submit_comfy_sync.py + + # Submit only ComfyUI installation sync + python maintenance/submit_comfy_sync.py --type installation + + # Submit only models sync + python maintenance/submit_comfy_sync.py --type models + + # Custom settings + python maintenance/submit_comfy_sync.py --pool mypool --allowlist "GPU01,GPU02" +""" + +import argparse +import io +import locale +import os +import subprocess +import sys +import tempfile + +# Default configuration +DEFAULT_POOL = "none" +DEFAULT_GROUP = "none" +DEFAULT_REGION = "" +DEFAULT_PRIORITY = 100 +DEFAULT_WORKER_DEADLINE_COMMAND = "deadlinecommand" +DEFAULT_ALLOWLIST = () +DEFAULT_ALLOWLIST_STRING = ",".join(DEFAULT_ALLOWLIST) + +# Script paths +COMFYUI_SYNC_SCRIPT = r"\\YOUR-SERVER\share\scripts\maintenance\ComfyUISync.py" +COMFY_MODELS_SYNC_SCRIPT = r"\\YOUR-SERVER\share\scripts\maintenance\ComfyModelsSync.py" + + +def normalize_allowlist(raw_value): + """Return a Deadline-ready comma-separated allowlist string and token list.""" + if not raw_value: + return "", [] + tokens = [] + for chunk in raw_value.replace("\n", ",").split(","): + candidate = chunk.strip() + if candidate: + tokens.append(candidate) + return ",".join(tokens), tokens + + +def parse_args(): + parser = argparse.ArgumentParser( + description="Submit ComfyUI sync maintenance jobs to Deadline." + ) + parser.add_argument( + "--type", + choices=["both", "installation", "models"], + default="both", + help="Type of sync job to submit (default: both)" + ) + parser.add_argument("--job-name", help="Base job name (will be suffixed with sync type)") + parser.add_argument( + "--comment", + help="Base job comment/description (will be suffixed with sync type)", + ) + parser.add_argument("--pool", default=DEFAULT_POOL, help="Deadline pool.") + parser.add_argument("--group", default=DEFAULT_GROUP, help="Deadline group.") + parser.add_argument("--region", default=DEFAULT_REGION, help="Deadline region label.") + parser.add_argument("--priority", type=int, default=DEFAULT_PRIORITY, help="Job priority (0-100).") + parser.add_argument( + "--machine-limit", + type=int, + help="Limit each job to N Workers (default: auto = size of allow list).", + ) + parser.add_argument( + "--allowlist", + default=DEFAULT_ALLOWLIST_STRING, + help="Comma or newline separated Worker names allowed to run the jobs.", + ) + parser.add_argument( + "--comfyui-script", + default=COMFYUI_SYNC_SCRIPT, + help="UNC path to ComfyUISync.py accessible by Workers.", + ) + parser.add_argument( + "--models-script", + default=COMFY_MODELS_SYNC_SCRIPT, + help="UNC path to ComfyModelsSync.py accessible by Workers.", + ) + parser.add_argument( + "--additional-arguments", + default="", + help="Extra arguments passed to sync scripts (optional).", + ) + parser.add_argument( + "--deadline-command", + default=None, + help="Path to deadlinecommand executable for submitting jobs. " + "Defaults to DEADLINE_PATH or PATH lookup.", + ) + parser.add_argument( + "--worker-deadline-command", + default=DEFAULT_WORKER_DEADLINE_COMMAND, + help="Executable that Workers should run inside the DeadlineCommand plugin.", + ) + parser.add_argument( + "--suspended", + action="store_true", + help="Submit jobs in suspended state (resume manually).", + ) + parser.add_argument( + "--user-name", + default="", + help="Override the Deadline job owner (leave blank to use current user).", + ) + parser.add_argument( + "--maintenance", + dest="maintenance", + action="store_true", + help="Mark as maintenance jobs (default).", + ) + parser.add_argument( + "--no-maintenance", + dest="maintenance", + action="store_false", + help="Submit without the MaintenanceJob flag.", + ) + parser.set_defaults(maintenance=True) + return parser.parse_args() + + +def find_deadline_command(explicit_path): + if explicit_path: + return explicit_path + + env_path = os.environ.get("DEADLINE_PATH") + candidates = [] + if env_path: + candidates.append(os.path.join(env_path, "deadlinecommand.exe")) + candidates.append(os.path.join(env_path, "deadlinecommand")) + + candidates.append("deadlinecommand") + + for candidate in candidates: + if candidate == "deadlinecommand": + return candidate + if os.path.isfile(candidate): + return candidate + + return "deadlinecommand" + + +def create_temp_file(lines, suffix): + fd, path = tempfile.mkstemp(suffix=suffix) + try: + with io.open(fd, "w", encoding="utf-8") as handle: + handle.write(u"\n".join(lines)) + handle.write(u"\n") + except Exception: + os.close(fd) + os.unlink(path) + raise + return path + + +def build_job_info(args, job_name, comment): + lines = [ + "Plugin=DeadlineCommand", + "Name={0}".format(job_name), + "Comment={0}".format(comment), + "Pool={0}".format(args.pool), + "Group={0}".format(args.group), + "Region={0}".format(args.region), + "Priority={0}".format(args.priority), + "MachineLimit={0}".format(args.machine_limit), + "MaintenanceJob={0}".format("true" if args.maintenance else "false"), + "OnJobComplete=Nothing", + ] + if args.allowlist: + lines.append("Whitelist={0}".format(args.allowlist)) + if args.user_name: + lines.append("UserName={0}".format(args.user_name)) + if args.suspended: + lines.append("SubmitSuspended=true") + return lines + + +def build_plugin_info(args, script_path): + argument_parts = ["-ExecuteScript", '"{0}"'.format(script_path)] + if args.additional_arguments: + argument_parts.append(args.additional_arguments) + + arguments = " ".join(argument_parts) + lines = [ + "Executable={0}".format(args.worker_deadline_command), + "Arguments={0}".format(arguments), + "ShellExecute=False", + "StartupDirectory=", + "SingleFramesOnly=False", + "HideWindow=False", + "IgnoreExitCode=False", + ] + return lines + + +def submit_job(deadline_command, job_info_path, plugin_info_path, job_name): + cmd = [deadline_command, "SubmitJob", job_info_path, plugin_info_path] + print("[submit_comfy_sync] Submitting job '{0}'...".format(job_name)) + print("[submit_comfy_sync] Running:", " ".join(cmd)) + + process = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE) + stdout_bytes, stderr_bytes = process.communicate() + encoding = locale.getpreferredencoding(False) or "utf-8" + + stdout_text = stdout_bytes.decode(encoding, "ignore") if isinstance(stdout_bytes, bytes) else stdout_bytes + stderr_text = stderr_bytes.decode(encoding, "ignore") if isinstance(stderr_bytes, bytes) else stderr_bytes + + if process.returncode != 0: + raise RuntimeError( + "deadlinecommand failed (code {0}):\nSTDOUT:\n{1}\nSTDERR:\n{2}".format( + process.returncode, stdout_text, stderr_text + ) + ) + + print(stdout_text.strip()) + return stdout_text + + +def submit_comfyui_sync(args, deadline_command): + """Submit ComfyUI installation sync job.""" + job_name = args.job_name or "ComfyUI Installation Sync" + comment = args.comment or "Mirrors ComfyUI portable installation from shared storage to local Workers" + + print("\n[submit_comfy_sync] Preparing ComfyUI installation sync job...") + + if not os.path.exists(args.comfyui_script): + print( + "[submit_comfy_sync] Warning: ComfyUI sync script not found right now ({0}). " + "Ensure Workers can reach this path.".format(args.comfyui_script) + ) + + job_info_lines = build_job_info(args, job_name, comment) + plugin_info_lines = build_plugin_info(args, args.comfyui_script) + + job_info_path = create_temp_file(job_info_lines, ".job") + plugin_info_path = create_temp_file(plugin_info_lines, ".plugin") + + try: + result = submit_job(deadline_command, job_info_path, plugin_info_path, job_name) + return result + finally: + for temp_path in (job_info_path, plugin_info_path): + try: + os.unlink(temp_path) + except OSError: + pass + + +def submit_models_sync(args, deadline_command): + """Submit ComfyUI models sync job.""" + job_name = (args.job_name + " - Models" if args.job_name else "ComfyUI Models Sync") + comment = (args.comment + " - Models" if args.comment else "Syncs ComfyUI AI models to local Worker storage based on modellist.txt") + + print("\n[submit_comfy_sync] Preparing ComfyUI models sync job...") + + if not os.path.exists(args.models_script): + print( + "[submit_comfy_sync] Warning: Models sync script not found right now ({0}). " + "Ensure Workers can reach this path.".format(args.models_script) + ) + + job_info_lines = build_job_info(args, job_name, comment) + plugin_info_lines = build_plugin_info(args, args.models_script) + + job_info_path = create_temp_file(job_info_lines, ".job") + plugin_info_path = create_temp_file(plugin_info_lines, ".plugin") + + try: + result = submit_job(deadline_command, job_info_path, plugin_info_path, job_name) + return result + finally: + for temp_path in (job_info_path, plugin_info_path): + try: + os.unlink(temp_path) + except OSError: + pass + + +def main(): + args = parse_args() + allowlist_str, allowlist_tokens = normalize_allowlist(args.allowlist) + args.allowlist = allowlist_str + allowlist_count = len(allowlist_tokens) + if args.machine_limit is None: + args.machine_limit = allowlist_count if allowlist_count else 0 + + print("[submit_comfy_sync] ComfyUI sync job submission starting...") + print("[submit_comfy_sync] Pool: {0}".format(args.pool)) + print("[submit_comfy_sync] Group: {0}".format(args.group)) + print("[submit_comfy_sync] Region: {0}".format(args.region)) + print("[submit_comfy_sync] Allowlist: {0} workers".format(allowlist_count)) + print("[submit_comfy_sync] Sync type: {0}".format(args.type)) + + deadline_command = find_deadline_command(args.deadline_command) + + jobs_submitted = [] + + try: + if args.type in ["both", "installation"]: + result = submit_comfyui_sync(args, deadline_command) + jobs_submitted.append(("ComfyUI Installation", result)) + + if args.type in ["both", "models"]: + result = submit_models_sync(args, deadline_command) + jobs_submitted.append(("ComfyUI Models", result)) + + print("\n[submit_comfy_sync] Job submission completed successfully!") + print("[submit_comfy_sync] Jobs submitted: {0}".format(len(jobs_submitted))) + for job_type, result in jobs_submitted: + print("[submit_comfy_sync] - {0}: {1}".format(job_type, result.strip() if result else "Unknown")) + + except Exception as e: + print("[submit_comfy_sync] ERROR: {0}".format(e)) + sys.exit(1) + + +if __name__ == "__main__": + try: + main() + except KeyboardInterrupt: + sys.exit(130) + except Exception as exc: + print("[submit_comfy_sync] ERROR:", exc) + sys.exit(1)