Compare commits
29
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a4e3bc2571 | ||
|
|
bf8e2a17b9 | ||
|
|
60c93b65f4 | ||
|
|
5d2ade4fd4 | ||
|
|
1c0a4d5950 | ||
|
|
1a8e089c7d | ||
|
|
557799cdb1 | ||
|
|
d5f52131db | ||
|
|
755399a255 | ||
|
|
ee3717b4b3 | ||
|
|
c13828af56 | ||
|
|
ffda951cfa | ||
|
|
ac0d8ba725 | ||
|
|
31c38ddb9b | ||
|
|
2f2c63fa11 | ||
|
|
7675c5ba90 | ||
|
|
7d65d23f84 | ||
|
|
a91edf6371 | ||
|
|
967709dd07 | ||
|
|
07a827d56b | ||
|
|
7f6eea361e | ||
|
|
14829a2b12 | ||
|
|
03c4e2d85d | ||
|
|
bad8b1104f | ||
|
|
ac8779dc54 | ||
|
|
2a5223f2f0 | ||
|
|
f02aef4cb3 | ||
|
|
22458a1cd6 | ||
|
|
0a58eba554 |
@@ -2,6 +2,10 @@
|
||||
|
||||
Open source comfyui deployment platform, a `vercel` for generative workflow infra. (serverless hosted gpu with vertical intergation with comfyui)
|
||||
|
||||
Check out our latest lcoal demo -> https://github.com/comfy-deploy/comfyui-api-comfydeploy
|
||||
|
||||
Full backend and frontend is here -> https://github.com/comfy-deploy/comfydeploy
|
||||
|
||||
> [!NOTE]
|
||||
> Im looking for creative hacker to join ComfyDeploy's core team! DM me on [twitter](https://x.com/BennyKokMusic)
|
||||
|
||||
|
||||
@@ -0,0 +1,137 @@
|
||||
import folder_paths
|
||||
|
||||
|
||||
class AnyType(str):
|
||||
def __ne__(self, __value: object) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
WILDCARD = AnyType("*")
|
||||
|
||||
|
||||
class ComfyUIDeployExternalFile:
|
||||
@classmethod
|
||||
def INPUT_TYPES(s):
|
||||
return {
|
||||
"required": {
|
||||
"input_id": (
|
||||
"STRING",
|
||||
{"multiline": False, "default": "input_file"},
|
||||
),
|
||||
},
|
||||
"optional": {
|
||||
"display_name": (
|
||||
"STRING",
|
||||
{"multiline": False, "default": ""},
|
||||
),
|
||||
"description": (
|
||||
"STRING",
|
||||
{"multiline": True, "default": ""},
|
||||
),
|
||||
"file_url": (
|
||||
"STRING",
|
||||
{"multiline": False, "default": ""},
|
||||
),
|
||||
},
|
||||
}
|
||||
|
||||
RETURN_TYPES = (WILDCARD,)
|
||||
RETURN_NAMES = ("path",)
|
||||
FUNCTION = "run"
|
||||
CATEGORY = "🔗ComfyDeploy"
|
||||
|
||||
def run(
|
||||
self,
|
||||
input_id,
|
||||
display_name=None,
|
||||
description=None,
|
||||
file_url=None,
|
||||
):
|
||||
import requests
|
||||
import os
|
||||
import uuid
|
||||
from urllib.parse import urlparse
|
||||
|
||||
if file_url:
|
||||
if file_url.startswith("http"):
|
||||
# Use cache directory for saving files
|
||||
cache_dir = folder_paths.get_temp_directory()
|
||||
if not os.path.exists(cache_dir):
|
||||
os.makedirs(cache_dir)
|
||||
|
||||
# Always generate random filename to avoid conflicts
|
||||
parsed_url = urlparse(file_url)
|
||||
original_filename = os.path.basename(parsed_url.path)
|
||||
|
||||
# Extract file extension from original filename if available
|
||||
file_extension = ""
|
||||
if original_filename and "." in original_filename:
|
||||
file_extension = os.path.splitext(original_filename)[1]
|
||||
else:
|
||||
# Try to determine extension from content-type if no extension found
|
||||
file_extension = ".bin"
|
||||
|
||||
# Generate random filename with preserved extension
|
||||
filename = str(uuid.uuid4()) + file_extension
|
||||
|
||||
destination_path = os.path.join(cache_dir, filename)
|
||||
print(f"Cache directory: {cache_dir}")
|
||||
print(f"Destination path: {destination_path}")
|
||||
print(
|
||||
"Downloading external file - "
|
||||
+ file_url
|
||||
+ " to "
|
||||
+ destination_path
|
||||
)
|
||||
|
||||
headers = {"User-Agent": "Mozilla/5.0"}
|
||||
|
||||
try:
|
||||
response = requests.get(
|
||||
file_url,
|
||||
headers=headers,
|
||||
allow_redirects=True,
|
||||
timeout=30, # Add timeout to prevent hanging
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
with open(destination_path, "wb") as out_file:
|
||||
out_file.write(response.content)
|
||||
print(f"External file downloaded: {file_url} to {destination_path}")
|
||||
return (destination_path,)
|
||||
|
||||
except requests.exceptions.HTTPError as e:
|
||||
error_msg = f"HTTP Error {e.response.status_code}: {e.response.reason} for URL: {file_url}"
|
||||
print(f"⚠️ Download failed - {error_msg}")
|
||||
if e.response.status_code == 404:
|
||||
print(
|
||||
"💡 This URL might have expired or the file may have been deleted"
|
||||
)
|
||||
# Return empty string instead of crashing
|
||||
return ("",)
|
||||
|
||||
except requests.exceptions.RequestException as e:
|
||||
error_msg = (
|
||||
f"Network error downloading file from {file_url}: {str(e)}"
|
||||
)
|
||||
print(f"⚠️ Download failed - {error_msg}")
|
||||
return ("",)
|
||||
|
||||
except Exception as e:
|
||||
error_msg = (
|
||||
f"Unexpected error downloading file from {file_url}: {str(e)}"
|
||||
)
|
||||
print(f"⚠️ Download failed - {error_msg}")
|
||||
return ("",)
|
||||
else:
|
||||
print(f"External file loading: {file_url}")
|
||||
return (file_url,)
|
||||
else:
|
||||
print(f"No file URL provided")
|
||||
return ("",)
|
||||
|
||||
|
||||
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalFile": ComfyUIDeployExternalFile}
|
||||
NODE_DISPLAY_NAME_MAPPINGS = {
|
||||
"ComfyUIDeployExternalFile": "External File (ComfyUI Deploy)"
|
||||
}
|
||||
@@ -0,0 +1,78 @@
|
||||
# In file: comfyui-deploy/comfy-nodes/output_exr.py
|
||||
|
||||
import os
|
||||
import numpy as np
|
||||
import folder_paths
|
||||
|
||||
# Try to set up OpenCV for EXR writing.
|
||||
try:
|
||||
os.environ["OPENCV_IO_ENABLE_OPENEXR"] = "1"
|
||||
import cv2
|
||||
OPENCV_AVAILABLE = True
|
||||
except ImportError:
|
||||
print("Warning: OpenCV not found for ComfyDeployOutputEXR. Please add opencv-python-headless to requirements.txt")
|
||||
OPENCV_AVAILABLE = False
|
||||
|
||||
# ALIGNED: Renamed class to match project conventions
|
||||
class ComfyDeployOutputEXR:
|
||||
def __init__(self):
|
||||
self.output_dir = folder_paths.get_output_directory()
|
||||
self.type = "output"
|
||||
|
||||
@classmethod
|
||||
def INPUT_TYPES(s):
|
||||
return {
|
||||
"required": {
|
||||
"images": ("IMAGE", ),
|
||||
"filename_prefix": ("STRING", {"default": "ComfyDeploy_EXR"})
|
||||
},
|
||||
# ADDED: Optional output_id for consistency with other ComfyDeploy nodes
|
||||
"optional": {
|
||||
"output_id": ("STRING", {"multiline": False, "default": "output_exr"}),
|
||||
},
|
||||
}
|
||||
|
||||
RETURN_TYPES = ()
|
||||
# ALIGNED: Changed function name to 'run'
|
||||
FUNCTION = "run"
|
||||
OUTPUT_NODE = True
|
||||
# ALIGNED: Matched the category name
|
||||
CATEGORY = "🔗ComfyDeploy"
|
||||
DESCRIPTION = "Saves the input images as EXR (HDR) files."
|
||||
|
||||
def run(self, images, filename_prefix="ComfyDeploy_EXR", output_id="output_exr"):
|
||||
if not OPENCV_AVAILABLE:
|
||||
raise ImportError("OpenCV is required to save EXR files. Please ensure opencv-python-headless is in requirements.txt.")
|
||||
|
||||
full_output_folder, filename, counter, subfolder, filename_prefix = (
|
||||
folder_paths.get_save_image_path(
|
||||
filename_prefix, self.output_dir, images[0].shape[1], images[0].shape[0]
|
||||
)
|
||||
)
|
||||
results = list()
|
||||
|
||||
for image in images:
|
||||
image_np = image.cpu().numpy()
|
||||
|
||||
if image_np.dtype != np.float32:
|
||||
image_np = image_np.astype(np.float32)
|
||||
|
||||
file = f"{filename}_{counter:05}.exr"
|
||||
file_path = os.path.join(full_output_folder, file)
|
||||
|
||||
image_np_bgr = cv2.cvtColor(image_np, cv2.COLOR_RGB2BGR)
|
||||
cv2.imwrite(file_path, image_np_bgr)
|
||||
|
||||
results.append({
|
||||
"filename": file,
|
||||
"subfolder": subfolder,
|
||||
"type": self.type,
|
||||
"output_id": output_id, # ADDED
|
||||
})
|
||||
counter += 1
|
||||
|
||||
return {"ui": {"images": results}}
|
||||
|
||||
# ALIGNED: Mappings are defined at the bottom of the node file in this project
|
||||
NODE_CLASS_MAPPINGS = {"ComfyDeployOutputEXR": ComfyDeployOutputEXR}
|
||||
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyDeployOutputEXR": "EXR Output (ComfyDeploy)"}
|
||||
@@ -0,0 +1,168 @@
|
||||
import folder_paths
|
||||
import os
|
||||
import shutil
|
||||
import uuid
|
||||
|
||||
|
||||
class AnyType(str):
|
||||
def __ne__(self, __value: object) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
WILDCARD = AnyType("*")
|
||||
|
||||
|
||||
class ComfyDeployOutputFile:
|
||||
def __init__(self):
|
||||
self.output_dir = folder_paths.get_output_directory()
|
||||
self.type = "output"
|
||||
|
||||
@classmethod
|
||||
def INPUT_TYPES(s):
|
||||
return {
|
||||
"required": {
|
||||
"file_path": (
|
||||
"STRING",
|
||||
{
|
||||
"forceInput": True,
|
||||
"tooltip": "Path to the file to output and upload.",
|
||||
},
|
||||
),
|
||||
},
|
||||
"optional": {
|
||||
"output_id": (
|
||||
"STRING",
|
||||
{"multiline": False, "default": "output_file"},
|
||||
),
|
||||
},
|
||||
}
|
||||
|
||||
RETURN_TYPES = ()
|
||||
FUNCTION = "run"
|
||||
OUTPUT_NODE = True
|
||||
CATEGORY = "🔗ComfyDeploy"
|
||||
DESCRIPTION = "Outputs any file by path for upload to ComfyDeploy."
|
||||
|
||||
def run(self, file_path, output_id="output_file"):
|
||||
if not file_path or not os.path.exists(file_path):
|
||||
print(f"⚠️ File not found: {file_path}")
|
||||
return {"ui": {"files": []}}
|
||||
|
||||
# Security checks - ensure file is within safe ComfyUI paths
|
||||
try:
|
||||
# Get absolute paths for comparison
|
||||
file_abs_path = os.path.abspath(file_path)
|
||||
base_path = folder_paths.base_path
|
||||
temp_dir = folder_paths.get_temp_directory()
|
||||
|
||||
# Check if file is within ComfyUI base path or temp directory
|
||||
if not (
|
||||
file_abs_path.startswith(os.path.abspath(base_path))
|
||||
or file_abs_path.startswith(os.path.abspath(temp_dir))
|
||||
):
|
||||
print(f"⚠️ Security: File outside allowed ComfyUI paths: {file_path}")
|
||||
return {"ui": {"files": []}}
|
||||
|
||||
# Check for path traversal attempts (but allow absolute paths within ComfyUI)
|
||||
if ".." in file_path:
|
||||
print(f"⚠️ Security: Path traversal attempt detected: {file_path}")
|
||||
return {"ui": {"files": []}}
|
||||
|
||||
except Exception as e:
|
||||
print(f"⚠️ Security check failed: {str(e)}")
|
||||
return {"ui": {"files": []}}
|
||||
|
||||
# Get the original filename and extension
|
||||
original_filename = os.path.basename(file_path)
|
||||
file_extension = os.path.splitext(original_filename)[1]
|
||||
|
||||
# Additional filename security check
|
||||
if ".." in original_filename:
|
||||
print(f"⚠️ Security: Insecure filename: {original_filename}")
|
||||
return {"ui": {"files": []}}
|
||||
|
||||
results = []
|
||||
|
||||
# Check if file is in output folder, if not, symlink it there
|
||||
try:
|
||||
if file_path.startswith(self.output_dir):
|
||||
# File is already in output directory - use as is
|
||||
relative_path = os.path.relpath(file_path, self.output_dir)
|
||||
path_parts = relative_path.split(os.sep)
|
||||
|
||||
if len(path_parts) > 1:
|
||||
subfolder = os.sep.join(path_parts[:-1])
|
||||
else:
|
||||
subfolder = ""
|
||||
|
||||
filename = path_parts[-1]
|
||||
file_type = self.type
|
||||
else:
|
||||
# File is not in output folder - symlink it to output/temp
|
||||
print(
|
||||
f"File is not in output folder, symlinking to output/temp: {file_path}"
|
||||
)
|
||||
output_temp_dir = os.path.join(self.output_dir, "temp")
|
||||
if not os.path.exists(output_temp_dir):
|
||||
os.makedirs(output_temp_dir)
|
||||
|
||||
# Use the existing filename but with UUID prefix to avoid conflicts
|
||||
file_ext = os.path.splitext(original_filename)[1]
|
||||
temp_filename = f"{uuid.uuid4()}{file_ext}"
|
||||
temp_path = os.path.join(output_temp_dir, temp_filename)
|
||||
|
||||
# Create symlink to file in output/temp directory where upload system expects it
|
||||
try:
|
||||
# Remove existing symlink if it exists
|
||||
if os.path.exists(temp_path):
|
||||
os.remove(temp_path)
|
||||
os.symlink(file_path, temp_path)
|
||||
print(f"File symlinked to output/temp: {temp_path} -> {file_path}")
|
||||
except OSError as e:
|
||||
# Fall back to copying if symlink fails
|
||||
print(f"Symlink failed ({e}), falling back to copy")
|
||||
shutil.copy2(file_path, temp_path)
|
||||
print(f"File copied to output/temp: {temp_path}")
|
||||
|
||||
# Use output/temp directory structure for upload
|
||||
subfolder = "temp"
|
||||
filename = temp_filename
|
||||
file_type = self.type
|
||||
|
||||
results.append(
|
||||
{
|
||||
"filename": filename,
|
||||
"subfolder": subfolder,
|
||||
"type": file_type,
|
||||
"output_id": output_id,
|
||||
}
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
print(f"⚠️ Error processing file path: {str(e)}")
|
||||
return {"ui": {"files": []}}
|
||||
|
||||
# Determine the appropriate UI key based on file type
|
||||
file_ext = file_extension.lower()
|
||||
if file_ext in [".png", ".jpg", ".jpeg", ".webp", ".gif", ".bmp", ".tiff"]:
|
||||
ui_key = "images"
|
||||
elif file_ext in [".mp3", ".wav", ".flac", ".aac", ".ogg"]:
|
||||
ui_key = "audio"
|
||||
elif file_ext in [".txt", ".json", ".md", ".csv"]:
|
||||
ui_key = "text_file"
|
||||
elif file_ext in [".exr", ".hdr"]:
|
||||
ui_key = "images" # EXR files are still images
|
||||
elif file_ext in [".zip", ".psb", ".psd"]:
|
||||
ui_key = "files" # Archives and Photoshop project files
|
||||
else:
|
||||
ui_key = "files" # Generic files
|
||||
|
||||
return {"ui": {ui_key: results}}
|
||||
|
||||
|
||||
NODE_CLASS_MAPPINGS = {
|
||||
"ComfyDeployOutputFile": ComfyDeployOutputFile,
|
||||
}
|
||||
NODE_DISPLAY_NAME_MAPPINGS = {
|
||||
"ComfyDeployOutputFile": "File Output (ComfyDeploy)",
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import numpy as np
|
||||
from PIL import Image
|
||||
from PIL.PngImagePlugin import PngInfo
|
||||
import folder_paths
|
||||
from comfy.cli_args import args
|
||||
|
||||
|
||||
class ComfyDeployOutputImage:
|
||||
@@ -64,12 +65,14 @@ class ComfyDeployOutputImage:
|
||||
for batch_number, image in enumerate(images):
|
||||
i = 255.0 * image.cpu().numpy()
|
||||
img = Image.fromarray(np.clip(i, 0, 255).astype(np.uint8))
|
||||
metadata = PngInfo()
|
||||
if prompt is not None:
|
||||
metadata.add_text("prompt", json.dumps(prompt))
|
||||
if extra_pnginfo is not None:
|
||||
for x in extra_pnginfo:
|
||||
metadata.add_text(x, json.dumps(extra_pnginfo[x]))
|
||||
metadata = None
|
||||
if not args.disable_metadata:
|
||||
metadata = PngInfo()
|
||||
if prompt is not None:
|
||||
metadata.add_text("prompt", json.dumps(prompt))
|
||||
if extra_pnginfo is not None:
|
||||
for x in extra_pnginfo:
|
||||
metadata.add_text(x, json.dumps(extra_pnginfo[x]))
|
||||
|
||||
filename_with_batch_num = filename.replace("%batch_num%", str(batch_number))
|
||||
file = f"{filename_with_batch_num}_{counter:05}_.{file_type}"
|
||||
|
||||
+699
-94
@@ -265,7 +265,7 @@ def clear_current_prompt(sid):
|
||||
streaming_prompt_metadata[sid].running_prompt_ids.clear()
|
||||
|
||||
|
||||
def post_prompt(json_data):
|
||||
async def post_prompt(json_data):
|
||||
prompt_server = server.PromptServer.instance
|
||||
json_data = prompt_server.trigger_on_prompt(json_data)
|
||||
|
||||
@@ -281,7 +281,48 @@ def post_prompt(json_data):
|
||||
|
||||
if "prompt" in json_data:
|
||||
prompt = json_data["prompt"]
|
||||
valid = execution.validate_prompt(prompt)
|
||||
prompt_id = json_data.get("prompt_id") or str(uuid.uuid4())
|
||||
|
||||
partial_execution_targets = None
|
||||
if "partial_execution_targets" in json_data:
|
||||
partial_execution_targets = json_data["partial_execution_targets"]
|
||||
|
||||
# Handle different validate_prompt signatures (newest to oldest)
|
||||
valid = None
|
||||
last_error = None
|
||||
|
||||
# v0.3.48 (3 args)
|
||||
try:
|
||||
valid = await execution.validate_prompt(
|
||||
prompt_id, prompt, partial_execution_targets
|
||||
)
|
||||
except TypeError as e:
|
||||
last_error = e
|
||||
logger.debug(
|
||||
f"validate_prompt with 3 params not supported, trying with 2. Debug: {last_error}"
|
||||
)
|
||||
|
||||
# v0.3.45 - 0.3.47 (2 args)
|
||||
if valid is None:
|
||||
try:
|
||||
valid = await execution.validate_prompt(prompt_id, prompt)
|
||||
except TypeError as e:
|
||||
last_error = e
|
||||
logger.debug(
|
||||
f"validate_prompt with 2 params not supported, trying legacy signature. Debug: {last_error}"
|
||||
)
|
||||
|
||||
# v0.3.44 or older (1 arg)
|
||||
if valid is None:
|
||||
try:
|
||||
valid = execution.validate_prompt(prompt)
|
||||
except TypeError as e:
|
||||
last_error = e
|
||||
logger.error(
|
||||
f"validate_prompt failed with all signatures. Last error: {last_error}"
|
||||
)
|
||||
raise
|
||||
|
||||
extra_data = {}
|
||||
if "extra_data" in json_data:
|
||||
extra_data = json_data["extra_data"]
|
||||
@@ -292,8 +333,6 @@ def post_prompt(json_data):
|
||||
if "client_id" in json_data:
|
||||
extra_data["client_id"] = json_data["client_id"]
|
||||
if valid[0]:
|
||||
# if the prompt id is provided
|
||||
prompt_id = json_data.get("prompt_id") or str(uuid.uuid4())
|
||||
outputs_to_execute = valid[2]
|
||||
prompt_server.prompt_queue.put(
|
||||
(number, prompt_id, prompt, extra_data, outputs_to_execute)
|
||||
@@ -477,6 +516,9 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
|
||||
if value["class_type"] == "ComfyUIDeployExternalEXR":
|
||||
value["inputs"]["exr_file"] = new_value
|
||||
|
||||
if value["class_type"] == "ComfyUIDeployExternalFile":
|
||||
value["inputs"]["file_url"] = new_value
|
||||
|
||||
if value["class_type"] == "ComfyUIDeployExternalSeed":
|
||||
logger.info(
|
||||
f"Applied random seed {new_value} to {value['class_type']}"
|
||||
@@ -500,15 +542,15 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
|
||||
|
||||
prompt_id = str(uuid.uuid4())
|
||||
|
||||
prompt = {
|
||||
"prompt": workflow_api,
|
||||
"client_id": sid, # "comfy_deploy_instance", #api.client_id
|
||||
"prompt_id": prompt_id,
|
||||
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
|
||||
}
|
||||
# prompt = {
|
||||
# "prompt": workflow_api,
|
||||
# "client_id": sid, # "comfy_deploy_instance", #api.client_id
|
||||
# "prompt_id": prompt_id,
|
||||
# "extra_data": {"extra_pnginfo": {"workflow": workflow}},
|
||||
# }
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
# res = post_prompt(prompt)
|
||||
inputs.running_prompt_ids.add(prompt_id)
|
||||
prompt_metadata[prompt_id] = SimplePrompt(
|
||||
status_endpoint=inputs.status_endpoint,
|
||||
@@ -519,7 +561,7 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
|
||||
except Exception as e:
|
||||
error_type = type(e).__name__
|
||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||
stack_trace = traceback.format_exc().strip()
|
||||
# stack_trace = traceback.format_exc().strip()
|
||||
logger.info(f"error: {error_type}, {e}")
|
||||
logger.info(f"stack trace: {stack_trace_short}")
|
||||
|
||||
@@ -595,7 +637,7 @@ async def comfy_deploy_run(request):
|
||||
)
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
res = await post_prompt(prompt)
|
||||
except Exception as e:
|
||||
error_type = type(e).__name__
|
||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||
@@ -633,6 +675,14 @@ async def comfy_deploy_run(request):
|
||||
return web.json_response(res, status=status)
|
||||
|
||||
|
||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/interrupt")
|
||||
async def interrupt_prompt(request):
|
||||
data = await request.json()
|
||||
prompt_id = data.get("prompt_id")
|
||||
await update_run(prompt_id, Status.CANCELLED)
|
||||
return web.json_response({"message": "Prompt interrupted"}, status=200)
|
||||
|
||||
|
||||
async def stream_prompt(data, token):
|
||||
# In older version, we use workflow_api, but this has inputs already swapped in nextjs frontend, which is tricky
|
||||
workflow_api = data.get("workflow_api_raw")
|
||||
@@ -664,7 +714,7 @@ async def stream_prompt(data, token):
|
||||
# log('info', "Begin prompt", prompt=prompt)
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
res = await post_prompt(prompt)
|
||||
except Exception as e:
|
||||
error_type = type(e).__name__
|
||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||
@@ -849,6 +899,10 @@ async def upload_file_endpoint(request):
|
||||
file_type = "image/png"
|
||||
elif file_extension == ".webp":
|
||||
file_type = "image/webp"
|
||||
elif file_extension == ".zip":
|
||||
file_type = "application/zip"
|
||||
elif file_extension in [".psd", ".psb"]:
|
||||
file_type = "image/vnd.adobe.photoshop"
|
||||
else:
|
||||
file_type = (
|
||||
"application/octet-stream" # Default to binary file type if unknown
|
||||
@@ -1266,22 +1320,11 @@ def handle_execute(class_type, last_node_id, prompt_id, server, unique_id):
|
||||
|
||||
try:
|
||||
origin_execute = execution.execute
|
||||
is_async = asyncio.iscoroutinefunction(origin_execute)
|
||||
|
||||
def swizzle_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
):
|
||||
unique_id = current_item
|
||||
class_type = dynprompt.get_node(unique_id)["class_type"]
|
||||
last_node_id = server.last_node_id
|
||||
result = origin_execute(
|
||||
if is_async:
|
||||
|
||||
async def swizzle_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
@@ -1291,12 +1334,61 @@ try:
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
)
|
||||
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||
return result
|
||||
pending_async_nodes,
|
||||
):
|
||||
unique_id = current_item
|
||||
class_type = dynprompt.get_node(unique_id)["class_type"]
|
||||
last_node_id = server.last_node_id
|
||||
|
||||
result = await origin_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
pending_async_nodes,
|
||||
)
|
||||
|
||||
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||
return result
|
||||
else:
|
||||
|
||||
def swizzle_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
):
|
||||
unique_id = current_item
|
||||
class_type = dynprompt.get_node(unique_id)["class_type"]
|
||||
last_node_id = server.last_node_id
|
||||
|
||||
result = origin_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
)
|
||||
|
||||
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||
return result
|
||||
|
||||
execution.execute = swizzle_execute
|
||||
except Exception as e:
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@@ -1440,10 +1532,6 @@ async def send_json_override(self, event, data, sid=None):
|
||||
logger.info(format_table(headers, table_data))
|
||||
# print("========================\n")
|
||||
|
||||
timeline = format_execution_timeline(NODE_EXECUTION_TIMES)
|
||||
logger.info(f"\nNode Execution Timeline:\n{timeline}")
|
||||
# Clear the execution times for the next run
|
||||
|
||||
# the last executing event is none, then the workflow is finished
|
||||
if event == "executing" and data.get("node") is None:
|
||||
mark_prompt_done(prompt_id=prompt_id)
|
||||
@@ -1466,18 +1554,26 @@ async def send_json_override(self, event, data, sid=None):
|
||||
)
|
||||
|
||||
if event == "executing" and data.get("node") is not None:
|
||||
node = data.get("node")
|
||||
raw_node = data.get("node")
|
||||
node = str(raw_node)
|
||||
|
||||
if prompt_id in prompt_metadata:
|
||||
# if 'progress' not in prompt_metadata[prompt_id]:
|
||||
# prompt_metadata[prompt_id]["progress"] = set()
|
||||
wf_api = prompt_metadata[prompt_id].workflow_api
|
||||
|
||||
# Normalize dotted display ids like "23.0.0.1" to base "23"
|
||||
if node not in wf_api and "." in node:
|
||||
base = node.split(".")[0]
|
||||
if base in wf_api:
|
||||
node = base
|
||||
|
||||
# If still unknown, skip safely
|
||||
if node not in wf_api:
|
||||
logger.info(f"Skipping unknown node id in 'executing': {raw_node}")
|
||||
return
|
||||
|
||||
prompt_metadata[prompt_id].progress.add(node)
|
||||
calculated_progress = len(prompt_metadata[prompt_id].progress) / len(
|
||||
prompt_metadata[prompt_id].workflow_api
|
||||
)
|
||||
calculated_progress = len(prompt_metadata[prompt_id].progress) / len(wf_api)
|
||||
calculated_progress = round(calculated_progress, 2)
|
||||
# logger.info("calculated_progress", calculated_progress)
|
||||
|
||||
if (
|
||||
prompt_metadata[prompt_id].last_updated_node is not None
|
||||
@@ -1485,7 +1581,8 @@ async def send_json_override(self, event, data, sid=None):
|
||||
):
|
||||
return
|
||||
prompt_metadata[prompt_id].last_updated_node = node
|
||||
class_type = prompt_metadata[prompt_id].workflow_api[node]["class_type"]
|
||||
|
||||
class_type = wf_api[node]["class_type"]
|
||||
logger.info(f"At: {round(calculated_progress * 100)}% - {class_type}")
|
||||
asyncio.create_task(
|
||||
send(
|
||||
@@ -2017,6 +2114,10 @@ async def handle_upload(
|
||||
file_type = "image/png"
|
||||
elif file_extension == ".webp":
|
||||
file_type = "image/webp"
|
||||
elif file_extension == ".zip":
|
||||
file_type = "application/zip"
|
||||
elif file_extension in [".psd", ".psb"]:
|
||||
file_type = "image/vnd.adobe.photoshop"
|
||||
|
||||
upload_tasks.append(
|
||||
upload_file(
|
||||
@@ -2059,6 +2160,7 @@ async def upload_in_background(
|
||||
("model_file", "format", "application/octet-stream"),
|
||||
("result", "format", "application/octet-stream"),
|
||||
("text_file", "format", "text/plain"),
|
||||
("audio", "format", "audio/mpeg"),
|
||||
]:
|
||||
items = data.get(file_type, [])
|
||||
|
||||
@@ -2146,6 +2248,7 @@ async def update_run_with_output(
|
||||
or "model_file" in data
|
||||
or "result" in data
|
||||
or "text_file" in data
|
||||
or "audio" in data
|
||||
)
|
||||
if bypass_upload and have_upload_media:
|
||||
print(
|
||||
@@ -2480,6 +2583,10 @@ class UploadQueue:
|
||||
content_type = "image/webp"
|
||||
elif file_extension == ".gif":
|
||||
content_type = "image/gif"
|
||||
elif file_extension == ".zip":
|
||||
content_type = "application/zip"
|
||||
elif file_extension in [".psd", ".psb"]:
|
||||
content_type = "image/vnd.adobe.photoshop"
|
||||
else:
|
||||
content_type = file_info.get("content_type", "application/octet-stream")
|
||||
|
||||
@@ -2754,59 +2861,57 @@ class UploadQueue:
|
||||
logger.error(f"Upload failed: {str(e)}")
|
||||
logger.error(traceback.format_exc())
|
||||
finally:
|
||||
# Remove this upload from tracking
|
||||
if prompt_id in self.pending_uploads:
|
||||
self.pending_uploads[prompt_id].discard(upload_id)
|
||||
# Remove from node tracking if applicable
|
||||
if (
|
||||
node_id
|
||||
and prompt_id in self.node_uploads
|
||||
and node_id in self.node_uploads[prompt_id]
|
||||
):
|
||||
self.node_uploads[prompt_id][node_id].discard(upload_id)
|
||||
async with self.lock: # Acquire lock to protect shared dict access
|
||||
if prompt_id in self.pending_uploads:
|
||||
self.pending_uploads[prompt_id].discard(upload_id)
|
||||
|
||||
# If this was the last upload for this node, clean up node data
|
||||
if not self.node_uploads[prompt_id][node_id]:
|
||||
del self.node_uploads[prompt_id][node_id]
|
||||
if prompt_id in self.node_output_data:
|
||||
if node_id in self.node_output_data[prompt_id]:
|
||||
if self.node_output_data[prompt_id][node_id][
|
||||
"data"
|
||||
]:
|
||||
# Send final node data to API before cleanup
|
||||
if prompt_metadata[
|
||||
prompt_id
|
||||
].status_endpoint:
|
||||
body = {
|
||||
"run_id": prompt_id,
|
||||
"output_data": self.node_output_data[
|
||||
if (
|
||||
node_id
|
||||
and prompt_id in self.node_uploads
|
||||
and node_id in self.node_uploads[prompt_id]
|
||||
):
|
||||
self.node_uploads[prompt_id][node_id].discard(upload_id)
|
||||
|
||||
if not self.node_uploads[prompt_id][node_id]:
|
||||
del self.node_uploads[prompt_id][node_id]
|
||||
|
||||
if (
|
||||
prompt_id in self.node_output_data
|
||||
and node_id in self.node_output_data[prompt_id]
|
||||
):
|
||||
node_data = self.node_output_data[prompt_id][
|
||||
node_id
|
||||
]
|
||||
if node_data["data"]:
|
||||
body = {
|
||||
"run_id": prompt_id,
|
||||
"output_data": node_data["data"],
|
||||
"node_meta": {"node_id": node_id},
|
||||
}
|
||||
try:
|
||||
await async_request_with_retry(
|
||||
"POST",
|
||||
prompt_metadata[
|
||||
prompt_id
|
||||
][node_id]["data"],
|
||||
"node_meta": {"node_id": node_id},
|
||||
}
|
||||
try:
|
||||
await async_request_with_retry(
|
||||
"POST",
|
||||
prompt_metadata[
|
||||
prompt_id
|
||||
].status_endpoint,
|
||||
token=prompt_metadata[
|
||||
prompt_id
|
||||
].token,
|
||||
json=body,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"Failed to send final node data: {str(e)}"
|
||||
)
|
||||
].status_endpoint,
|
||||
token=prompt_metadata[
|
||||
prompt_id
|
||||
].token,
|
||||
json=body,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"Failed to send final node data: {str(e)}"
|
||||
)
|
||||
|
||||
# Safe to delete now (re-check not strictly needed with lock, but harmless)
|
||||
del self.node_output_data[prompt_id][node_id]
|
||||
|
||||
# Send status update
|
||||
await self.update_queue_status(prompt_id)
|
||||
|
||||
# If no more pending uploads for this prompt and it's done, update status
|
||||
if not self.pending_uploads[prompt_id] and is_prompt_done(
|
||||
prompt_id
|
||||
if (
|
||||
prompt_id in self.pending_uploads
|
||||
and not self.pending_uploads[prompt_id]
|
||||
and is_prompt_done(prompt_id)
|
||||
):
|
||||
# Clean up all data for this prompt
|
||||
if prompt_id in self.node_uploads:
|
||||
@@ -2819,9 +2924,12 @@ class UploadQueue:
|
||||
loop.create_task(update_run(prompt_id, Status.SUCCESS))
|
||||
loop.create_task(send("success", {"prompt_id": prompt_id}))
|
||||
|
||||
# Mark task as done
|
||||
# Mark task as done (outside lock to avoid holding it unnecessarily)
|
||||
self.queue.task_done()
|
||||
|
||||
# Send status update (also outside lock)
|
||||
await self.update_queue_status(prompt_id)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error in upload worker: {str(e)}")
|
||||
logger.error(traceback.format_exc())
|
||||
@@ -2919,6 +3027,7 @@ async def create_workflow_proxy(request):
|
||||
name = data.get("name")
|
||||
workflow_json = data.get("workflow_json")
|
||||
workflow_api = data.get("workflow_api")
|
||||
machine_id = data.get("machine_id")
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
|
||||
auth_header = request.headers.get("Authorization")
|
||||
@@ -2938,6 +3047,7 @@ async def create_workflow_proxy(request):
|
||||
"name": name,
|
||||
"workflow_json": json.dumps(workflow_json),
|
||||
"workflow_api": json.dumps(workflow_api),
|
||||
"machine_id": machine_id,
|
||||
}
|
||||
|
||||
try:
|
||||
@@ -3060,3 +3170,498 @@ async def get_workflow_proxy(request):
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# fetch workflow versions (infinite scroll support)
|
||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/workflow/versions")
|
||||
async def get_workflow_versions_proxy(request):
|
||||
workflow_id = request.rel_url.query.get("workflow_id")
|
||||
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
|
||||
search = request.rel_url.query.get("search", "")
|
||||
limit = request.rel_url.query.get("limit", "20")
|
||||
offset = request.rel_url.query.get("offset", "0")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
# Build target URL with query params
|
||||
params = {"limit": limit, "offset": offset}
|
||||
if search:
|
||||
params["search"] = search
|
||||
|
||||
query = urlencode(params)
|
||||
target_url = f"{api_url}/api/workflow/{workflow_id}/versions"
|
||||
if query:
|
||||
target_url += f"?{query}"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.get(
|
||||
target_url, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# fetch a specific workflow version json
|
||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/workflow/version")
|
||||
async def get_workflow_version_proxy(request):
|
||||
workflow_id = request.rel_url.query.get("workflow_id")
|
||||
version = request.rel_url.query.get("version")
|
||||
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/workflow/{workflow_id}/version/{version}"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.get(
|
||||
target_url, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# for getting a machine by id
|
||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/machine")
|
||||
async def get_machine_proxy(request):
|
||||
machine_id = request.rel_url.query.get("machine_id")
|
||||
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/machine/{machine_id}"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.get(
|
||||
target_url, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# for fetching docker steps from current snapshot
|
||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/snapshot-to-docker")
|
||||
async def snapshot_to_docker_proxy(request):
|
||||
data = await request.json()
|
||||
snapshot = data.get("snapshot")
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/snapshot-to-docker"
|
||||
|
||||
request_body = snapshot
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url, json=request_body, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# update a serverless machine with machine id
|
||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/machine/update")
|
||||
async def update_machine_proxy(request):
|
||||
data = await request.json()
|
||||
machine_id = data.get("machine_id")
|
||||
comfyui_version = data.get("comfyui_version", None)
|
||||
docker_steps = data.get("docker_steps")
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/machine/serverless/{machine_id}"
|
||||
|
||||
request_body = {"docker_command_steps": docker_steps}
|
||||
|
||||
if comfyui_version:
|
||||
request_body["comfyui_version"] = comfyui_version
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.patch(
|
||||
target_url, json=request_body, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/machine/create")
|
||||
async def create_machine_proxy(request):
|
||||
data = await request.json()
|
||||
name = data.get("name")
|
||||
docker_command_steps = data.get("docker_command_steps")
|
||||
comfyui_version = data.get("comfyui_version")
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/machine/serverless"
|
||||
|
||||
request_body = {
|
||||
"name": name,
|
||||
"docker_command_steps": docker_command_steps,
|
||||
"comfyui_version": comfyui_version,
|
||||
"gpu": "A10G",
|
||||
}
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url, json=request_body, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# get latest comfyui version
|
||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/comfyui-version")
|
||||
async def get_comfyui_version_proxy(request):
|
||||
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/latest-hashes"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.get(
|
||||
target_url, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Proxy: generate single-part upload URL
|
||||
@server.PromptServer.instance.routes.post(
|
||||
"/comfyui-deploy/volume/file/generate-upload-url"
|
||||
)
|
||||
async def proxy_generate_upload_url(request):
|
||||
data = await request.json()
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/volume/file/generate-upload-url"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url,
|
||||
json={
|
||||
"filename": data.get("filename"),
|
||||
"contentType": data.get("contentType"),
|
||||
"size": data.get("size"),
|
||||
},
|
||||
headers={"Authorization": auth_header},
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Proxy: initiate multipart upload
|
||||
@server.PromptServer.instance.routes.post(
|
||||
"/comfyui-deploy/volume/file/initiate-multipart-upload"
|
||||
)
|
||||
async def proxy_initiate_multipart_upload(request):
|
||||
data = await request.json()
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/volume/file/initiate-multipart-upload"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url,
|
||||
json={
|
||||
"filename": data.get("filename"),
|
||||
"contentType": data.get("contentType"),
|
||||
"size": data.get("size"),
|
||||
},
|
||||
headers={"Authorization": auth_header},
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Proxy: generate part upload URL
|
||||
@server.PromptServer.instance.routes.post(
|
||||
"/comfyui-deploy/volume/file/generate-part-upload-url"
|
||||
)
|
||||
async def proxy_generate_part_upload_url(request):
|
||||
data = await request.json()
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/volume/file/generate-part-upload-url"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url,
|
||||
json={
|
||||
"uploadId": data.get("uploadId"),
|
||||
"key": data.get("key"),
|
||||
"partNumber": data.get("partNumber"),
|
||||
},
|
||||
headers={"Authorization": auth_header},
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Proxy: complete multipart upload
|
||||
@server.PromptServer.instance.routes.post(
|
||||
"/comfyui-deploy/volume/file/complete-multipart-upload"
|
||||
)
|
||||
async def proxy_complete_multipart_upload(request):
|
||||
data = await request.json()
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/volume/file/complete-multipart-upload"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url,
|
||||
json={
|
||||
"uploadId": data.get("uploadId"),
|
||||
"key": data.get("key"),
|
||||
"parts": data.get("parts"),
|
||||
},
|
||||
headers={"Authorization": auth_header},
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Proxy: abort multipart upload
|
||||
@server.PromptServer.instance.routes.post(
|
||||
"/comfyui-deploy/volume/file/abort-multipart-upload"
|
||||
)
|
||||
async def proxy_abort_multipart_upload(request):
|
||||
data = await request.json()
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/volume/file/abort-multipart-upload"
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url,
|
||||
json={
|
||||
"uploadId": data.get("uploadId"),
|
||||
"key": data.get("key"),
|
||||
},
|
||||
headers={"Authorization": auth_header},
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Proxy: add model (unified endpoint)
|
||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/volume/model")
|
||||
async def proxy_add_model(request):
|
||||
data = await request.json()
|
||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||
auth_header = request.headers.get("Authorization")
|
||||
|
||||
if not auth_header:
|
||||
return web.json_response(
|
||||
{"error": "Authorization header is required"}, status=401
|
||||
)
|
||||
|
||||
target_url = f"{api_url}/api/volume/model"
|
||||
|
||||
# pass body through but remove api_url key
|
||||
forward_body = dict(data)
|
||||
if "api_url" in forward_body:
|
||||
forward_body.pop("api_url")
|
||||
|
||||
try:
|
||||
await ensure_client_session()
|
||||
async with client_session.post(
|
||||
target_url, json=forward_body, headers={"Authorization": auth_header}
|
||||
) as response:
|
||||
json_data = await response.json()
|
||||
return web.json_response(json_data, status=response.status)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# FS: stat file (size)
|
||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/fs/stat")
|
||||
async def fs_stat(request):
|
||||
try:
|
||||
import os
|
||||
|
||||
file_path = request.rel_url.query.get("path")
|
||||
if not file_path:
|
||||
return web.json_response({"error": "path is required"}, status=400)
|
||||
|
||||
# Basic safeguard: ensure it's a ComfyUI models path
|
||||
if "/models/" not in file_path:
|
||||
return web.json_response({"error": "invalid path"}, status=400)
|
||||
|
||||
st = os.stat(file_path)
|
||||
return web.json_response({"size": st.st_size})
|
||||
except FileNotFoundError:
|
||||
return web.json_response({"error": "not found"}, status=404)
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
|
||||
# Upload a multipart part directly from machine filesystem to S3 upload URL
|
||||
@server.PromptServer.instance.routes.post(
|
||||
"/comfyui-deploy/volume/file/upload-part-from-path"
|
||||
)
|
||||
async def upload_part_from_path(request):
|
||||
try:
|
||||
import os
|
||||
import io
|
||||
|
||||
data = await request.json()
|
||||
file_path = data.get("filePath")
|
||||
upload_url = data.get("uploadUrl")
|
||||
start = int(data.get("start", 0))
|
||||
end = int(data.get("end", 0))
|
||||
if not file_path or not upload_url:
|
||||
return web.json_response(
|
||||
{"error": "filePath and uploadUrl are required"}, status=400
|
||||
)
|
||||
if "/models/" not in file_path:
|
||||
return web.json_response({"error": "invalid filePath"}, status=400)
|
||||
if end <= start:
|
||||
return web.json_response({"error": "invalid byte range"}, status=400)
|
||||
|
||||
size = end - start
|
||||
|
||||
await ensure_client_session()
|
||||
|
||||
# Important: S3 pre-signed part uploads do not support chunked transfer
|
||||
# Buffer the exact part into memory to provide a Content-Length header
|
||||
buffer = bytearray()
|
||||
chunk_size = 4 * 1024 * 1024
|
||||
with open(file_path, "rb") as f:
|
||||
f.seek(start)
|
||||
remaining = size
|
||||
while remaining > 0:
|
||||
to_read = chunk_size if remaining >= chunk_size else remaining
|
||||
chunk = f.read(to_read)
|
||||
if not chunk:
|
||||
break
|
||||
buffer.extend(chunk)
|
||||
remaining -= len(chunk)
|
||||
|
||||
if len(buffer) != size:
|
||||
return web.json_response(
|
||||
{
|
||||
"error": "read size mismatch",
|
||||
"expected": size,
|
||||
"actual": len(buffer),
|
||||
},
|
||||
status=500,
|
||||
)
|
||||
|
||||
headers = {
|
||||
"Content-Length": str(size),
|
||||
"Content-Type": "application/octet-stream",
|
||||
}
|
||||
|
||||
async with client_session.put(
|
||||
upload_url, data=bytes(buffer), headers=headers
|
||||
) as resp:
|
||||
text = await resp.text()
|
||||
if resp.status < 200 or resp.status >= 300:
|
||||
return web.json_response(
|
||||
{"error": f"upload failed: {resp.status}", "body": text},
|
||||
status=resp.status,
|
||||
)
|
||||
etag = resp.headers.get("ETag") or resp.headers.get("etag") or ""
|
||||
etag = etag.replace('"', "")
|
||||
return web.json_response({"eTag": etag, "bytesSent": size})
|
||||
except Exception as e:
|
||||
return web.json_response({"error": str(e)}, status=500)
|
||||
|
||||
@@ -18,6 +18,7 @@ class Status(Enum):
|
||||
SUCCESS = "success"
|
||||
FAILED = "failed"
|
||||
UPLOADING = "uploading"
|
||||
CANCELLED = "cancelled"
|
||||
|
||||
|
||||
class StreamingPrompt(BaseModel):
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
[project]
|
||||
name = "comfyui-deploy"
|
||||
description = "Open source comfyui deployment platform, a vercel for generative workflow infra."
|
||||
version = "2.1.0"
|
||||
version = "2.3.9"
|
||||
license = { file = "LICENSE" }
|
||||
dependencies = ["aiofiles", "pydantic", "opencv-python", "imageio-ffmpeg", "tabulate", "brotli"]
|
||||
|
||||
|
||||
+1063
-348
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,68 @@
|
||||
// Snapshot Utilities
|
||||
// Centralized snapshot fetching with ComfyUI version fallback
|
||||
|
||||
/**
|
||||
* Fetches the current snapshot with ComfyUI version fallback
|
||||
* If the snapshot response has null comfyui field, it will fetch the latest ComfyUI version
|
||||
* and update the snapshot with the comfyui_hash
|
||||
*
|
||||
* @param {Function} getDataFn - Function that returns { apiKey, apiUrl } for ComfyUI version API calls
|
||||
* @returns {Promise<Object>} - The snapshot data with comfyui field populated
|
||||
*/
|
||||
export async function fetchSnapshot(getDataFn = null) {
|
||||
try {
|
||||
// Fetch the current snapshot
|
||||
const response = await fetch("/snapshot/get_current");
|
||||
if (!response.ok) {
|
||||
throw new Error(`Snapshot fetch failed: ${response.status}`);
|
||||
}
|
||||
|
||||
const snapshot = await response.json();
|
||||
|
||||
// Check if comfyui field is null and we have getDataFn for fallback
|
||||
if (snapshot.comfyui === null && getDataFn) {
|
||||
console.log(
|
||||
"ComfyUI version is null in snapshot, fetching latest version..."
|
||||
);
|
||||
|
||||
try {
|
||||
const data = getDataFn();
|
||||
if (data && data.apiKey) {
|
||||
const comfyuiVersionResponse = await fetch(
|
||||
`/comfyui-deploy/comfyui-version?api_url=${encodeURIComponent(
|
||||
data.apiUrl || "https://api.comfydeploy.com"
|
||||
)}`,
|
||||
{
|
||||
headers: {
|
||||
Authorization: `Bearer ${data.apiKey}`,
|
||||
},
|
||||
}
|
||||
);
|
||||
|
||||
if (comfyuiVersionResponse.ok) {
|
||||
const versionData = await comfyuiVersionResponse.json();
|
||||
if (versionData.comfyui_hash) {
|
||||
console.log(
|
||||
`Using ComfyUI hash from API: ${versionData.comfyui_hash}`
|
||||
);
|
||||
snapshot.comfyui = versionData.comfyui_hash;
|
||||
}
|
||||
} else {
|
||||
console.warn(
|
||||
"Failed to fetch ComfyUI version from API:",
|
||||
comfyuiVersionResponse.status
|
||||
);
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
console.warn("Error fetching ComfyUI version fallback:", error);
|
||||
// Continue with original snapshot even if fallback fails
|
||||
}
|
||||
}
|
||||
|
||||
return snapshot;
|
||||
} catch (error) {
|
||||
console.error("Error fetching snapshot:", error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
+957
-3
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user