Compare commits

..
Author SHA1 Message Date
nick f5940ac899 better 2024-08-20 19:43:04 -07:00
nick 67703abb8a syntax 2024-08-20 19:41:23 -07:00
nick dfee31f0ed fix: noise seed 2024-08-20 19:36:29 -07:00
18 changed files with 181 additions and 557 deletions
+2 -2
View File
@@ -12,11 +12,11 @@ class ComfyUIDeployExternalBoolean:
"optional": {
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+2 -2
View File
@@ -25,11 +25,11 @@ class ComfyUIDeployExternalCheckpoint:
"default_value": (folder_paths.get_filename_list("checkpoints"), ),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
-108
View File
@@ -1,108 +0,0 @@
from PIL import Image, ImageOps
import numpy as np
import torch
import folder_paths
class AnyType(str):
def __ne__(self, __value: object) -> bool:
return False
WILDCARD = AnyType("*")
class ComfyUIDeployExternalFaceModel:
@classmethod
def INPUT_TYPES(s):
return {
"required": {
"input_id": (
"STRING",
{"multiline": False, "default": "input_reactor_face_model"},
),
},
"optional": {
"default_face_model_name": (
"STRING",
{"multiline": False, "default": ""},
),
"face_model_save_name": ( # if `default_face_model_name` is a link to download a file, we will attempt to save it with this name
"STRING",
{"multiline": False, "default": ""},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
"face_model_url": (
"STRING",
{"multiline": False, "default": ""},
),
},
}
RETURN_TYPES = (WILDCARD,)
RETURN_NAMES = ("path",)
FUNCTION = "run"
CATEGORY = "deploy"
def run(
self,
input_id,
default_face_model_name=None,
face_model_save_name=None,
display_name=None,
description=None,
face_model_url=None,
):
import requests
import os
import uuid
if face_model_url and face_model_url.startswith("http"):
if face_model_save_name:
existing_face_models = folder_paths.get_filename_list("reactor/faces")
# Check if face_model_save_name exists in the list
if face_model_save_name in existing_face_models:
print(f"using face model: {face_model_save_name}")
return (face_model_save_name,)
else:
face_model_save_name = str(uuid.uuid4()) + ".safetensors"
print(face_model_save_name)
print(folder_paths.folder_names_and_paths["reactor/faces"][0][0])
destination_path = os.path.join(
folder_paths.folder_names_and_paths["reactor/faces"][0][0],
face_model_save_name,
)
print(destination_path)
print(
"Downloading external face model - "
+ face_model_url
+ " to "
+ destination_path
)
response = requests.get(
face_model_url,
headers={"User-Agent": "Mozilla/5.0"},
allow_redirects=True,
)
with open(destination_path, "wb") as out_file:
out_file.write(response.content)
return (face_model_save_name,)
else:
print(f"using face model: {default_face_model_name}")
return (default_face_model_name,)
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalFaceModel": ComfyUIDeployExternalFaceModel}
NODE_DISPLAY_NAME_MAPPINGS = {
"ComfyUIDeployExternalFaceModel": "External Face Model (ComfyUI Deploy)"
}
+2 -2
View File
@@ -17,11 +17,11 @@ class ComfyUIDeployExternalImage:
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+2 -2
View File
@@ -17,11 +17,11 @@ class ComfyUIDeployExternalImageAlpha:
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+4 -24
View File
@@ -23,11 +23,11 @@ class ComfyUIDeployExternalImageBatch:
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
@@ -39,34 +39,14 @@ class ComfyUIDeployExternalImageBatch:
CATEGORY = "image"
def process_image(self, image):
image = ImageOps.exif_transpose(image)
image = image.convert("RGB")
image = np.array(image).astype(np.float32) / 255.0
image_tensor = torch.from_numpy(image)[None,]
return image_tensor
def run(self, input_id, images=None, default_value=None, display_name=None, description=None):
import requests
import zipfile
import io
processed_images = []
try:
images_list = json.loads(images) # Assuming images is a JSON array string
print(images_list)
for img_input in images_list:
if img_input.startswith('http') and img_input.endswith('.zip'):
print("Fetching zip file from url: ", img_input)
response = requests.get(img_input)
zip_file = zipfile.ZipFile(io.BytesIO(response.content))
for file_name in zip_file.namelist():
if file_name.lower().endswith(('.png', '.jpg', '.jpeg')):
with zip_file.open(file_name) as file:
image = Image.open(file)
image = self.process_image(image)
processed_images.append(image)
elif img_input.startswith('http'):
if img_input.startswith('http'):
import requests
from io import BytesIO
print("Fetching image from url: ", img_input)
response = requests.get(img_input)
+8 -20
View File
@@ -25,21 +25,17 @@ class ComfyUIDeployExternalLora:
},
"optional": {
"default_lora_name": (folder_paths.get_filename_list("loras"),),
"lora_save_name": ( # if `default_lora_name` is a link to download a file, we will attempt to save it with this name
"lora_save_name": ( # if `default_lora_name` is a link to download a file, we will attempt to save it with this name
"STRING",
{"multiline": False, "default": ""},
),
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
"lora_url": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
},
}
@@ -51,20 +47,12 @@ class ComfyUIDeployExternalLora:
CATEGORY = "deploy"
def run(
self,
input_id,
default_lora_name=None,
lora_save_name=None,
display_name=None,
description=None,
lora_url=None,
):
def run(self, input_id, default_lora_name=None, lora_save_name=None, display_name=None, description=None):
import requests
import os
import uuid
if lora_url and lora_url.startswith("http"):
if default_lora_name.startswith("http"):
if lora_save_name:
existing_loras = folder_paths.get_filename_list("loras")
# Check if lora_save_name exists in the list
@@ -79,9 +67,9 @@ class ComfyUIDeployExternalLora:
folder_paths.folder_names_and_paths["loras"][0][0], lora_save_name
)
print(destination_path)
print("Downloading external lora - " + lora_url + " to " + destination_path)
print("Downloading external lora - " + input_id + " to " + destination_path)
response = requests.get(
lora_url,
input_id,
headers={"User-Agent": "Mozilla/5.0"},
allow_redirects=True,
)
+2 -2
View File
@@ -20,11 +20,11 @@ class ComfyUIDeployExternalNumber:
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+2 -2
View File
@@ -20,11 +20,11 @@ class ComfyUIDeployExternalNumberInt:
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+2 -2
View File
@@ -23,11 +23,11 @@ class ComfyUIDeployExternalNumberSlider:
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+2 -2
View File
@@ -20,11 +20,11 @@ class ComfyUIDeployExternalText:
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
-46
View File
@@ -1,46 +0,0 @@
class AnyType(str):
def __ne__(self, __value: object) -> bool:
return False
WILDCARD = AnyType("*")
class ComfyUIDeployExternalTextAny:
@classmethod
def INPUT_TYPES(s):
return {
"required": {
"input_id": (
"STRING",
{"multiline": False, "default": "input_text"},
),
},
"optional": {
"default_value": (
"STRING",
{"multiline": True, "default": ""},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
RETURN_TYPES = (WILDCARD,)
RETURN_NAMES = ("text",)
FUNCTION = "run"
CATEGORY = "text"
def run(self, input_id, default_value=None, display_name=None, description=None):
return [default_value]
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalTextAny": ComfyUIDeployExternalTextAny}
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyUIDeployExternalTextAny": "External Text Any (ComfyUI Deploy)"}
+2 -2
View File
@@ -21,11 +21,11 @@ class ComfyUIDeployExternalTextList:
"optional": {
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
}
}
+4 -4
View File
@@ -764,14 +764,14 @@ class ComfyUIDeployExternalVideo:
"optional": {
"meta_batch": ("VHS_BatchManager",),
"vae": ("VAE",),
"default_video": (sorted(files),),
"default_value": (sorted(files),),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
{"multiline": False, "default": "Name of the node (optional)"},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "default": "Description of the node (optional)"},
),
},
"hidden": {
@@ -834,7 +834,7 @@ class ComfyUIDeployExternalVideo:
):
out_file.write(chunk)
else:
video = kwargs.get("default_video", None)
video = kwargs.get("default_value", "")
if video is None:
raise "No default video given and no external video provided"
video_path = folder_paths.get_annotated_filepath(video.strip('"'))
+59 -244
View File
@@ -1,5 +1,4 @@
from io import BytesIO
from pprint import pprint
from aiohttp import web
import os
import requests
@@ -24,7 +23,7 @@ from typing import Dict, List, Union, Any, Optional
from PIL import Image
import copy
import struct
from aiohttp import web, ClientSession, ClientError, ClientTimeout, ClientResponseError
from aiohttp import web, ClientSession, ClientError, ClientTimeout
import atexit
# Global session
@@ -34,7 +33,7 @@ client_session = None
# global client_session
# if client_session is None:
# client_session = aiohttp.ClientSession()
async def ensure_client_session():
global client_session
if client_session is None:
@@ -44,7 +43,7 @@ async def cleanup():
global client_session
if client_session:
await client_session.close()
def exit_handler():
print("Exiting the application. Initiating cleanup...")
loop = asyncio.get_event_loop()
@@ -57,65 +56,39 @@ retry_delay_multiplier = float(os.environ.get('RETRY_DELAY_MULTIPLIER', '2'))
print(f"max_retries: {max_retries}, retry_delay_multiplier: {retry_delay_multiplier}")
import time
async def async_request_with_retry(method, url, disable_timeout=False, token=None, **kwargs):
async def async_request_with_retry(method, url, disable_timeout=False, **kwargs):
global client_session
await ensure_client_session()
# async with aiohttp.ClientSession() as client_session:
retry_delay = 1 # Start with 1 second delay
initial_timeout = 5 # 5 seconds timeout for the initial connection
start_time = time.time()
for attempt in range(max_retries):
try:
# Set a timeout for the initial connection
if not disable_timeout:
timeout = ClientTimeout(total=None, connect=initial_timeout)
kwargs['timeout'] = timeout
if token is not None:
if 'headers' not in kwargs:
kwargs['headers'] = {}
kwargs['headers']['Authorization'] = f"Bearer {token}"
request_start = time.time()
async with client_session.request(method, url, **kwargs) as response:
request_end = time.time()
# logger.info(f"Request attempt {attempt + 1} took {request_end - request_start:.2f} seconds")
if response.status != 200:
error_body = await response.text()
# logger.error(f"Request failed with status {response.status} and body {error_body}")
# raise Exception(f"Request failed with status {response.status}")
response.raise_for_status()
if method.upper() == 'GET':
await response.read()
total_time = time.time() - start_time
# logger.info(f"Request succeeded after {total_time:.2f} seconds (attempt {attempt + 1}/{max_retries})")
return response
except asyncio.TimeoutError:
logger.warning(f"Request timed out after {initial_timeout} seconds (attempt {attempt + 1}/{max_retries})")
except ClientError as e:
end_time = time.time()
logger.error(f"Request failed (attempt {attempt + 1}/{max_retries}): {e}")
logger.error(f"Time taken for failed attempt: {end_time - request_start:.2f} seconds")
logger.error(f"Total time elapsed: {end_time - start_time:.2f} seconds")
# Log the response body for ClientError as well
if hasattr(e, 'response') and e.response is not None:
error_body = await e.response.text()
logger.error(f"Error response body: {error_body}")
if attempt == max_retries - 1:
logger.error(f"Request failed after {max_retries} attempts: {e}")
raise
await asyncio.sleep(retry_delay)
retry_delay *= retry_delay_multiplier
# raise
logger.warning(f"Request failed (attempt {attempt + 1}/{max_retries}): {e}")
total_time = time.time() - start_time
raise Exception(f"Request failed after {max_retries} attempts and {total_time:.2f} seconds")
# Wait before retrying
await asyncio.sleep(retry_delay)
retry_delay *= retry_delay_multiplier # Exponential backoff
# If all retries fail, raise an exception
raise Exception(f"Request failed after {max_retries} attempts")
from logging import basicConfig, getLogger
@@ -143,7 +116,7 @@ def log(level, message, **kwargs):
getattr(logger, level)(message, **kwargs)
else:
getattr(logger, level)(f"{message} {kwargs}")
# For a span, you might need to create a context manager
from contextlib import contextmanager
@@ -286,7 +259,7 @@ def apply_random_seed_to_workflow(workflow_api):
logger.info(f"Applied random noise_seed {workflow_api[key]['inputs']['noise_seed']} to SamplerCustom")
continue
def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str | None = None):
# Loop through each of the inputs and replace them
for key, value in workflow_api.items():
if 'inputs' in value:
@@ -309,7 +282,7 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
value['inputs']["input_id"] = new_value
# Fix for external text default value
if (value["class_type"] == "ComfyUIDeployExternalText" or value["class_type"] == "ComfyUIDeployExternalTextAny"):
if (value["class_type"] == "ComfyUIDeployExternalText"):
value['inputs']["default_value"] = new_value
if (value["class_type"] == "ComfyUIDeployExternalCheckpoint"):
@@ -319,7 +292,7 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
value['inputs']["images"] = new_value
if value["class_type"] == "ComfyUIDeployExternalLora":
value["inputs"]["lora_url"] = new_value
value["inputs"]["default_lora_name"] = new_value
if value["class_type"] == "ComfyUIDeployExternalSlider":
value["inputs"]["default_value"] = new_value
@@ -327,13 +300,9 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
if value["class_type"] == "ComfyUIDeployExternalBoolean":
value["inputs"]["default_value"] = new_value
if value["class_type"] == "ComfyUIDeployExternalFaceModel":
value["inputs"]["face_model_url"] = new_value
def send_prompt(sid: str, inputs: StreamingPrompt):
# workflow_api = inputs.workflow_api
workflow_api = copy.deepcopy(inputs.workflow_api)
workflow = copy.deepcopy(inputs.workflow)
# Random seed
apply_random_seed_to_workflow(workflow_api)
@@ -349,8 +318,7 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
prompt = {
"prompt": workflow_api,
"client_id": sid, #"comfy_deploy_instance", #api.client_id
"prompt_id": prompt_id,
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
"prompt_id": prompt_id
}
try:
@@ -371,25 +339,13 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
@server.PromptServer.instance.routes.post("/comfyui-deploy/run")
async def comfy_deploy_run(request):
# Extract the bearer token from the Authorization header
data = await request.json()
if "cd_token" in data:
token = data["cd_token"]
else:
auth_header = request.headers.get('Authorization')
token = None
if auth_header:
parts = auth_header.split()
if len(parts) == 2 and parts[0].lower() == 'bearer':
token = parts[1]
# 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")
# The prompt id generated from comfy deploy, can be None
prompt_id = data.get("prompt_id")
inputs = data.get("inputs")
workflow = data.get("workflow")
# Now it handles directly in here
apply_random_seed_to_workflow(workflow_api)
@@ -398,15 +354,13 @@ async def comfy_deploy_run(request):
prompt = {
"prompt": workflow_api,
"client_id": "comfy_deploy_instance", #api.client_id
"prompt_id": prompt_id,
"extra_data": {"extra_pnginfo": {"workflow": workflow}}
"prompt_id": prompt_id
}
prompt_metadata[prompt_id] = SimplePrompt(
status_endpoint=data.get('status_endpoint'),
file_upload_endpoint=data.get('file_upload_endpoint'),
workflow_api=workflow_api,
token=token
workflow_api=workflow_api
)
try:
@@ -444,13 +398,12 @@ async def comfy_deploy_run(request):
return web.json_response(res, status=status)
async def stream_prompt(data, token):
async def stream_prompt(data):
# 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")
# The prompt id generated from comfy deploy, can be None
prompt_id = data.get("prompt_id")
inputs = data.get("inputs")
workflow = data.get("workflow")
# Now it handles directly in here
apply_random_seed_to_workflow(workflow_api)
@@ -459,15 +412,13 @@ async def stream_prompt(data, token):
prompt = {
"prompt": workflow_api,
"client_id": "comfy_deploy_instance", #api.client_id
"prompt_id": prompt_id,
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
"prompt_id": prompt_id
}
prompt_metadata[prompt_id] = SimplePrompt(
status_endpoint=data.get('status_endpoint'),
file_upload_endpoint=data.get('file_upload_endpoint'),
workflow_api=workflow_api,
token=token
workflow_api=workflow_api
)
# log('info', "Begin prompt", prompt=prompt)
@@ -516,14 +467,6 @@ comfy_message_queues: Dict[str, asyncio.Queue] = {}
async def stream_response(request):
response = web.StreamResponse(status=200, reason='OK', headers={'Content-Type': 'text/event-stream'})
await response.prepare(request)
# Extract the bearer token from the Authorization header
auth_header = request.headers.get('Authorization')
token = None
if auth_header:
parts = auth_header.split()
if len(parts) == 2 and parts[0].lower() == 'bearer':
token = parts[1]
pending = True
data = await request.json()
@@ -535,7 +478,7 @@ async def stream_response(request):
log('info', 'Streaming prompt')
try:
result = await stream_prompt(data=data, token=token)
result = await stream_prompt(data=data)
await response.write(f"event: event_update\ndata: {json.dumps(result)}\n\n".encode('utf-8'))
# await response.write(.encode('utf-8'))
await response.drain() # Ensure the buffer is flushed
@@ -664,10 +607,9 @@ async def upload_file_endpoint(request):
with open(file_path, 'rb') as f:
headers = {
"Content-Type": file_type,
# "Content-Length": str(file_size)
# "x-amz-acl": "public-read",
"Content-Length": str(file_size)
}
if content.get('include_acl') is True:
headers["x-amz-acl"] = "public-read"
upload_response = await async_request_with_retry('PUT', upload_url, data=f, headers=headers)
if upload_response.status == 200:
return web.json_response({
@@ -794,7 +736,6 @@ async def websocket_handler(request):
inputs={},
status_endpoint=status_endpoint,
file_upload_endpoint=request.rel_url.query.get('file_upload_endpoint', None),
workflow=workflow["workflow"],
)
await update_realtime_run_status(realtime_id, status_endpoint, Status.RUNNING)
@@ -914,19 +855,19 @@ async def send(event, data, sid=None):
except Exception as e:
logger.info(f"Exception: {e}")
traceback.print_exc()
@server.PromptServer.instance.routes.get('/comfydeploy/{tail:.*}')
@server.PromptServer.instance.routes.post('/comfydeploy/{tail:.*}')
async def proxy_to_comfydeploy(request):
# Get the base URL
base_url = f'https://www.comfydeploy.com/{request.match_info["tail"]}'
# Get all query parameters
query_params = request.query_string
# Construct the full target URL with query parameters
target_url = f"{base_url}?{query_params}" if query_params else base_url
# print(f"Proxying request to: {target_url}")
try:
@@ -984,7 +925,6 @@ async def send_json_override(self, event, data, sid=None):
"data": data
})
asyncio.create_task(update_run_ws_event(prompt_id, event, data))
# event_emitter.emit("send_json", {
# "event": event,
# "data": data
@@ -1052,22 +992,17 @@ async def send_json_override(self, event, data, sid=None):
# await update_run_with_output(prompt_id, data)
if event == 'executed' and 'node' in data and 'output' in data:
node_meta = None
if prompt_id in prompt_metadata:
node = data.get('node')
class_type = prompt_metadata[prompt_id].workflow_api[node]['class_type']
logger.info(f"Executed {class_type} {data}")
node_meta = {
"node_id": node,
"node_class": class_type,
}
if class_type == "PreviewImage":
logger.info("Skipping preview image")
return
else:
logger.info(f"Executed {data}")
await update_run_with_output(prompt_id, data.get('output'), node_id=data.get('node'), node_meta=node_meta)
await update_run_with_output(prompt_id, data.get('output'), node_id=data.get('node'))
# await update_run_with_output(prompt_id, data.get('output'), node_id=data.get('node'))
# update_run_with_output(prompt_id, data.get('output'))
@@ -1082,7 +1017,6 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
return
status_endpoint = prompt_metadata[prompt_id].status_endpoint
token = prompt_metadata[prompt_id].token
if (status_endpoint is None):
return
@@ -1106,27 +1040,7 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
})
# requests.post(status_endpoint, json=body)
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
async def update_run_ws_event(prompt_id: str, event: str, data: dict):
if prompt_id not in prompt_metadata:
return
# print("update_run_ws_event", prompt_id, event, data)
status_endpoint = prompt_metadata[prompt_id].status_endpoint
if status_endpoint is None:
return
token = prompt_metadata[prompt_id].token
body = {
"run_id": prompt_id,
"ws_event": {
"event": event,
"data": data,
},
}
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
await async_request_with_retry('POST', status_endpoint, json=body)
async def update_run(prompt_id: str, status: Status):
@@ -1157,8 +1071,7 @@ async def update_run(prompt_id: str, status: Status):
try:
# requests.post(status_endpoint, json=body)
if (status_endpoint is not None):
token = prompt_metadata[prompt_id].token
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
await async_request_with_retry('POST', status_endpoint, json=body)
if (status_endpoint is not None) and cd_enable_run_log and (status == Status.SUCCESS or status == Status.FAILED):
try:
@@ -1188,7 +1101,7 @@ async def update_run(prompt_id: str, status: Status):
]
}
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
await async_request_with_retry('POST', status_endpoint, json=body)
# requests.post(status_endpoint, json=body)
except Exception as log_error:
logger.info(f"Error reading log file: {log_error}")
@@ -1209,48 +1122,7 @@ async def update_run(prompt_id: str, status: Status):
})
async def file_sender(file_object, chunk_size):
while True:
chunk = await file_object.read(chunk_size)
if not chunk:
break
yield chunk
chunk_size = 1024 * 1024 # 1MB chunks, adjust as needed
async def upload_with_retry(session, url, headers, data, max_retries=3, initial_delay=1):
start_time = time.time() # Start timing here
for attempt in range(max_retries):
try:
async with session.put(url, headers=headers, data=data) as response:
upload_duration = time.time() - start_time
logger.info(f"Upload attempt {attempt + 1} completed in {upload_duration:.2f} seconds")
logger.info(f"Upload response status: {response.status}")
response.raise_for_status() # This will raise an exception for 4xx and 5xx status codes
response_text = await response.text()
logger.info(f"Response body: {response_text[:1000]}...")
logger.info("Upload successful")
return response # Successful upload, exit the retry loop
except (ClientError, ClientResponseError) as e:
logger.error(f"Upload attempt {attempt + 1} failed: {str(e)}")
if attempt < max_retries - 1: # If it's not the last attempt
delay = initial_delay * (2 ** attempt) # Exponential backoff
logger.info(f"Retrying in {delay} seconds...")
await asyncio.sleep(delay)
else:
logger.error("Max retries reached. Upload failed.")
raise # Re-raise the last exception if all retries are exhausted
except Exception as e:
logger.error(f"Unexpected error during upload: {str(e)}")
logger.error(traceback.format_exc())
raise # Re-raise unexpected exceptions immediately
async def upload_file(prompt_id, filename, subfolder=None, content_type="image/png", type="output", item=None):
async def upload_file(prompt_id, filename, subfolder=None, content_type="image/png", type="output"):
"""
Uploads file to S3 bucket using S3 client object
:return: None
@@ -1280,7 +1152,7 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
logger.info(f"Uploading file {file}")
file_upload_endpoint = prompt_metadata[prompt_id].file_upload_endpoint
token = prompt_metadata[prompt_id].token
filename = quote(filename)
prompt_id = quote(prompt_id)
content_type = quote(content_type)
@@ -1288,52 +1160,25 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
target_url = f"{file_upload_endpoint}?file_name={filename}&run_id={prompt_id}&type={content_type}&version=v2"
start_time = time.time() # Start timing here
logger.info(f"Target URL: {target_url}")
result = await async_request_with_retry("GET", target_url, disable_timeout=True, token=token)
result = await async_request_with_retry("GET", target_url, disable_timeout=True)
end_time = time.time() # End timing after the request is complete
logger.info("Time taken for getting file upload endpoint: {:.2f} seconds".format(end_time - start_time))
ok = await result.json()
logger.info(f"Result: {ok}")
start_time = time.time() # Start timing here
async with aiofiles.open(file, 'rb') as f:
data = await f.read()
size = str(len(data))
# logger.info(f"Image size: {size}")
start_time = time.time() # Start timing here
headers = {
# "x-amz-acl": "public-read",
"Content-Type": content_type,
"Content-Length": size,
"Content-Length": str(len(data)),
}
logger.info(headers)
if ok.get('include_acl') is True:
headers["x-amz-acl"] = "public-read"
# response = requests.put(ok.get("url"), headers=headers, data=data)
# response = await async_request_with_retry('PUT', ok.get("url"), headers=headers, data=data)
# logger.info(f"Upload file response status: {response.status}, status text: {response.reason}")
async with aiohttp.ClientSession() as session:
try:
response = await upload_with_retry(session, ok.get("url"), headers, data)
# Process successful response...
except Exception as e:
# Handle final failure...
logger.error(f"Upload ultimately failed: {str(e)}")
response = await async_request_with_retry('PUT', ok.get("url"), headers=headers, data=data)
logger.info(f"Upload file response status: {response.status}, status text: {response.reason}")
end_time = time.time() # End timing after the request is complete
logger.info("Upload time: {:.2f} seconds".format(end_time - start_time))
if item is not None:
file_download_url = ok.get("download_url")
if file_download_url is not None:
item["url"] = file_download_url
item["upload_duration"] = end_time - start_time
if ok.get("is_public") is not None:
item["is_public"] = ok.get("is_public")
def have_pending_upload(prompt_id):
if prompt_id in prompt_metadata and len(prompt_metadata[prompt_id].uploading_nodes) > 0:
@@ -1447,55 +1292,30 @@ async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, d
item.get("filename"),
subfolder=item.get("subfolder"),
type=item.get("type"),
content_type=file_type,
item=item
content_type=file_type
))
# await upload_file(
# prompt_id,
# item.get("filename"),
# subfolder=item.get("subfolder"),
# type=item.get("type"),
# content_type=file_type,
# item=item
# )
# Execute all upload tasks concurrently
await asyncio.gather(*upload_tasks)
# Upload files in the background
async def upload_in_background(prompt_id: str, data, node_id=None, have_upload=True, node_meta=None):
async def upload_in_background(prompt_id: str, data, node_id=None, have_upload=True):
try:
# await handle_upload(prompt_id, data, 'images', "content_type", "image/png")
# await handle_upload(prompt_id, data, 'files', "content_type", "image/png")
# await handle_upload(prompt_id, data, 'gifs', "format", "image/gif")
# await handle_upload(prompt_id, data, 'mesh', "format", "application/octet-stream")
upload_tasks = [
handle_upload(prompt_id, data, "images", "content_type", "image/png"),
handle_upload(prompt_id, data, "files", "content_type", "image/png"),
handle_upload(prompt_id, data, "gifs", "format", "image/gif"),
handle_upload(
prompt_id, data, "mesh", "format", "application/octet-stream"
),
handle_upload(prompt_id, data, 'images', "content_type", "image/png"),
handle_upload(prompt_id, data, 'files', "content_type", "image/png"),
handle_upload(prompt_id, data, 'gifs', "format", "image/gif"),
handle_upload(prompt_id, data, 'mesh', "format", "application/octet-stream")
]
await asyncio.gather(*upload_tasks)
status_endpoint = prompt_metadata[prompt_id].status_endpoint
token = prompt_metadata[prompt_id].token
if have_upload:
if status_endpoint is not None:
body = {
"run_id": prompt_id,
"output_data": data,
"node_meta": node_meta,
}
# pprint(body)
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
await update_file_status(prompt_id, data, False, node_id=node_id)
except Exception as e:
await handle_error(prompt_id, data, e)
async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
async def update_run_with_output(prompt_id, data, node_id=None):
if prompt_id not in prompt_metadata:
return
@@ -1506,13 +1326,9 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
body = {
"run_id": prompt_id,
"output_data": data,
"node_meta": node_meta,
"output_data": data
}
pprint(body)
have_upload_media = False
if data is not None:
have_upload_media = 'images' in data or 'files' in data or 'gifs' in data or 'mesh' in data
have_upload_media = 'images' in data or 'files' in data or 'gifs' in data or 'mesh' in data
if bypass_upload and have_upload_media:
print("CD_BYPASS_UPLOAD is enabled, skipping the upload of the output:", node_id)
return
@@ -1524,16 +1340,15 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
if have_upload_media:
await update_file_status(prompt_id, data, True, node_id=node_id)
# asyncio.create_task(upload_in_background(prompt_id, data, node_id=node_id, have_upload=have_upload_media, node_meta=node_meta))
await upload_in_background(prompt_id, data, node_id=node_id, have_upload=have_upload_media, node_meta=node_meta)
asyncio.create_task(upload_in_background(prompt_id, data, node_id=node_id, have_upload=have_upload_media))
# await upload_in_background(prompt_id, data, node_id=node_id, have_upload=have_upload)
except Exception as e:
await handle_error(prompt_id, data, e)
# requests.post(status_endpoint, json=body)
elif status_endpoint is not None:
token = prompt_metadata[prompt_id].token
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
if status_endpoint is not None:
await async_request_with_retry('POST', status_endpoint, json=body)
await send('outputs_uploaded', {
"prompt_id": prompt_id
@@ -1592,4 +1407,4 @@ if cd_enable_log:
@server.PromptServer.instance.routes.get("/comfyui-deploy/filename_list_cache")
async def get_filename_list_cache(_):
from folder_paths import filename_list_cache
return web.json_response({'filename_list': filename_list_cache})
return web.json_response({'filename_list': filename_list_cache})
-3
View File
@@ -24,14 +24,11 @@ class StreamingPrompt(BaseModel):
running_prompt_ids: set[str] = set()
status_endpoint: Optional[str]
file_upload_endpoint: Optional[str]
workflow: Any
class SimplePrompt(BaseModel):
status_endpoint: Optional[str]
file_upload_endpoint: Optional[str]
token: Optional[str]
workflow_api: dict
status: Status = Status.NOT_STARTED
progress: set = set()
+88 -89
View File
@@ -192,13 +192,11 @@ const ext = {
registerCustomNodes() {
/** @type {LGraphNode}*/
class ComfyDeploy extends LGraphNode {
class ComfyDeploy {
color = LGraphCanvas.node_colors.yellow.color;
bgcolor = LGraphCanvas.node_colors.yellow.bgcolor;
groupcolor = LGraphCanvas.node_colors.yellow.groupcolor;
constructor() {
super();
this.color = LGraphCanvas.node_colors.yellow.color;
this.bgcolor = LGraphCanvas.node_colors.yellow.bgcolor;
this.groupcolor = LGraphCanvas.node_colors.yellow.groupcolor;
if (!this.properties) {
this.properties = {};
this.properties.workflow_name = "";
@@ -206,75 +204,65 @@ const ext = {
this.properties.version = "";
}
this.addWidget(
"text",
ComfyWidgets.STRING(
this,
"workflow_name",
this.properties.workflow_name,
(v) => {
this.properties.workflow_name = v;
},
{ multiline: false }
[
"",
{
default: this.properties.workflow_name,
multiline: false,
},
],
app,
);
this.addWidget(
"text",
ComfyWidgets.STRING(
this,
"workflow_id",
this.properties.workflow_id,
(v) => {
this.properties.workflow_id = v;
},
{ multiline: false }
[
"",
{
default: this.properties.workflow_id,
multiline: false,
},
],
app,
);
this.addWidget(
"text",
ComfyWidgets.STRING(
this,
"version",
this.properties.version,
(v) => {
this.properties.version = v;
},
{ multiline: false }
["", { default: this.properties.version, multiline: false }],
app,
);
// this.widgets.forEach((w) => {
// // w.computeSize = () => [200,10]
// w.computedHeight = 2;
// })
this.widgets_start_y = 10;
this.setSize(this.computeSize());
// const config = { };
// console.log(this);
this.serialize_widgets = true;
this.isVirtualNode = true;
}
onExecute() {
// This method is called when the node is executed
// You can add any necessary logic here
}
onSerialize(o) {
// This method is called when the node is being serialized
// Ensure all necessary data is saved
if (!o.properties) {
o.properties = {};
}
o.properties.workflow_name = this.properties.workflow_name;
o.properties.workflow_id = this.properties.workflow_id;
o.properties.version = this.properties.version;
}
onConfigure(o) {
// This method is called when the node is being configured (e.g., when loading a saved graph)
// Ensure all necessary data is restored
if (o.properties) {
this.properties = { ...this.properties, ...o.properties };
this.widgets[0].value = this.properties.workflow_name || "";
this.widgets[1].value = this.properties.workflow_id || "";
this.widgets[2].value = this.properties.version || "1";
}
}
}
// Register the node type
LiteGraph.registerNodeType("ComfyDeploy", Object.assign(ComfyDeploy, {
title: "Comfy Deploy",
title_mode: LiteGraph.NORMAL_TITLE,
collapsable: true,
}));
// Load default visibility
LiteGraph.registerNodeType(
"ComfyDeploy",
Object.assign(ComfyDeploy, {
title_mode: LiteGraph.NORMAL_TITLE,
title: "Comfy Deploy",
collapsable: true,
}),
);
ComfyDeploy.category = "deploy";
},
@@ -443,10 +431,10 @@ function createDynamicUIHtml(data) {
<h3 style="font-size: 14px; font-weight: semibold; margin-bottom: 8px;">Missing Nodes</h3>
<p style="font-size: 12px;">These nodes are not found with any matching custom_nodes in the ComfyUI Manager Database</p>
${data.missing_nodes
.map((node) => {
return `<p style="font-size: 14px; color: #d69e2e;">${node}</p>`;
})
.join("")}
.map((node) => {
return `<p style="font-size: 14px; color: #d69e2e;">${node}</p>`;
})
.join("")}
</div>
`;
}
@@ -454,14 +442,17 @@ function createDynamicUIHtml(data) {
Object.values(data.custom_nodes).forEach((node) => {
html += `
<div style="border-bottom: 1px solid #e2e8f0; padding-top: 16px;">
<a href="${node.url
}" target="_blank" style="font-size: 18px; font-weight: semibold; color: white; text-decoration: none;">${node.name
}</a>
<a href="${
node.url
}" target="_blank" style="font-size: 18px; font-weight: semibold; color: white; text-decoration: none;">${
node.name
}</a>
<p style="font-size: 14px; color: #4b5563;">${node.hash}</p>
${node.warning
? `<p style="font-size: 14px; color: #d69e2e;">${node.warning}</p>`
: ""
}
${
node.warning
? `<p style="font-size: 14px; color: #d69e2e;">${node.warning}</p>`
: ""
}
</div>
`;
});
@@ -475,8 +466,9 @@ function createDynamicUIHtml(data) {
Object.entries(data.models).forEach(([section, items]) => {
html += `
<div style="border-bottom: 1px solid #e2e8f0; padding-top: 8px; padding-bottom: 8px;">
<h3 style="font-size: 18px; font-weight: semibold; margin-bottom: 8px;">${section.charAt(0).toUpperCase() + section.slice(1)
}</h3>`;
<h3 style="font-size: 18px; font-weight: semibold; margin-bottom: 8px;">${
section.charAt(0).toUpperCase() + section.slice(1)
}</h3>`;
items.forEach((item) => {
html += `<p style="font-size: 14px; color: ${textColor};">${item.name}</p>`;
});
@@ -492,8 +484,9 @@ function createDynamicUIHtml(data) {
Object.entries(data.files).forEach(([section, items]) => {
html += `
<div style="border-bottom: 1px solid #e2e8f0; padding-top: 8px; padding-bottom: 8px;">
<h3 style="font-size: 18px; font-weight: semibold; margin-bottom: 8px;">${section.charAt(0).toUpperCase() + section.slice(1)
}</h3>`;
<h3 style="font-size: 18px; font-weight: semibold; margin-bottom: 8px;">${
section.charAt(0).toUpperCase() + section.slice(1)
}</h3>`;
items.forEach((item) => {
html += `<p style="font-size: 14px; color: ${textColor};">${item.name}</p>`;
});
@@ -1013,12 +1006,14 @@ export class LoadingDialog extends ComfyDialog {
showLoading(title, message) {
this.show(`
<div style="width: 400px; display: flex; gap: 18px; flex-direction: column; overflow: unset">
<h3 style="margin: 0px; display: flex; align-items: center; justify-content: center; gap: 12px;">${title} ${this.loadingIcon
}</h3>
${message
? `<label style="max-width: 100%; white-space: pre-wrap; word-wrap: break-word;">${message}</label>`
: ""
}
<h3 style="margin: 0px; display: flex; align-items: center; justify-content: center; gap: 12px;">${title} ${
this.loadingIcon
}</h3>
${
message
? `<label style="max-width: 100%; white-space: pre-wrap; word-wrap: break-word;">${message}</label>`
: ""
}
</div>
`);
}
@@ -1284,17 +1279,21 @@ export class ConfigDialog extends ComfyDialog {
</label>
<label style="color: white; width: 100%;">
Endpoint:
<input id="endpoint" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;" type="text" value="${data.endpoint
}">
<input id="endpoint" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;" type="text" value="${
data.endpoint
}">
</label>
<div style="color: white;">
API Key: User / Org <button style="font-size: 18px;">${data.displayName ?? ""
}</button>
<input id="apiKey" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;" type="password" value="${data.apiKey
}">
API Key: User / Org <button style="font-size: 18px;">${
data.displayName ?? ""
}</button>
<input id="apiKey" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;" type="password" value="${
data.apiKey
}">
<button id="loginButton" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;">
${data.apiKey ? "Re-login with ComfyDeploy" : "Login with ComfyDeploy"
}
${
data.apiKey ? "Re-login with ComfyDeploy" : "Login with ComfyDeploy"
}
</button>
</div>
</div>
-1
View File
@@ -6,5 +6,4 @@ export const customInputNodes: Record<string, string> = {
ComfyUIDeployExternalNumberInt: "integer",
ComfyUIDeployExternalLora: "string - (public lora download url)",
ComfyUIDeployExternalCheckpoint: "string - (public checkpoints download url)",
ComfyUIDeployExternalFaceModel: "string - (public face model download url)",
};