Compare commits

...
Author SHA1 Message Date
bennykok e2fcf67aec fix: graph load 2024-09-25 12:59:00 -07:00
bennykok 69f63f4869 Merge branch 'jeff/fix-workflow-in-extra-data' into workspace-v3 2024-09-24 19:58:16 -07:00
bennykok 50860cd500 test 2024-09-24 19:45:53 -07:00
bennykok 2eb02fc92e fi 2024-09-24 19:36:57 -07:00
EdwinWong 5c6defbe62 fix: add workflow data to extra data 2024-09-24 15:35:48 -07:00
bennykok d1c54b2b6d fix: state 2024-09-23 19:01:47 -07:00
bennykok 3a6c3b1ae9 feat: add native run proxy 2024-09-23 15:31:13 -07:00
bennykok aea456cba9 fix face loader extenal load 2024-09-21 10:51:51 -07:00
bennykok 8c5e5c4277 feat: add ComfyUIDeployExternalTextAny 2024-09-21 10:39:34 -07:00
bennykok 02430ee62d remove some logs 2024-09-20 18:10:04 -07:00
Fawaz Kadem 764a8fee82 Add new external deploy node for face models (#66) 2024-09-18 17:00:51 -07:00
bennykok 61acffd355 fix 2024-09-18 08:20:35 -07:00
bennykok aa47f3523f fix 2024-09-17 23:36:24 -07:00
bennykok 7ed4284a6f fix 2024-09-17 23:25:19 -07:00
bennykok a403daa314 fix 2024-09-17 23:09:42 -07:00
bennykok ba9b187dcc fix 2024-09-17 22:59:27 -07:00
bennykok 1243fa4e58 fix 2024-09-17 22:55:08 -07:00
bennykok 0d1537963c fix 2024-09-17 21:48:42 -07:00
bennykok 0083b38dcc chore: log image size 2024-09-17 20:44:44 -07:00
bennykok b8dded1535 Revert "fix: roll back to unique session per request"
This reverts commit 5a78ca97bd.
2024-09-17 20:26:39 -07:00
bennykok 4927d81e73 chore: accept cd_token 2024-09-17 18:57:15 -07:00
bennykok fb6bb2357a Reapply "fix: back to sequential file upload"
This reverts commit 1f5a88b888.
2024-09-17 14:28:56 -07:00
bennykok 086d642360 Merge branch 'benny/log-sync' into public-main 2024-09-17 14:27:59 -07:00
bennykok 212daa838c Revert "feat: experiment with await + asyncio.gather for multi file in same node"
This reverts commit c08b68c41f.
2024-09-17 14:25:13 -07:00
bennykok c08b68c41f feat: experiment with await + asyncio.gather for multi file in same node 2024-09-17 12:56:42 -07:00
bennykok 5a78ca97bd fix: roll back to unique session per request 2024-09-16 23:57:40 -07:00
bennykok 1f5a88b888 Revert "fix: back to sequential file upload"
This reverts commit 3d099f88ea.
2024-09-16 23:55:16 -07:00
bennykok 946571e32e fix: await 2024-09-16 18:54:05 -07:00
bennykok e692beb009 feat: realtime log sync 2024-09-16 15:34:20 -07:00
bennykok 3d099f88ea fix: back to sequential file upload 2024-09-16 13:55:02 -07:00
karrix 65f7576748 fix: non type error when upload output 2024-09-16 12:45:53 -07:00
bennykok 2d72cd8175 fix: batch zip image input 2024-09-14 21:49:17 -07:00
bennykok 5554c95f44 Merge branch 'benny/auth_token' into public-main 2024-09-12 14:14:16 -07:00
bennykok c1003f7e31 Merge branch 'benny/zip-batch-image' into public-main 2024-09-12 14:14:08 -07:00
EdwinWong 71d60a5dd1 fix: comfydeploy node backward compatible in every comfyui 2024-09-10 01:03:50 -07:00
bennykok 4df9d38e56 feat: embed file public status into image output 2024-09-03 23:07:48 -07:00
bennykok 9cd626e1f6 feat: send token for cd update api 2024-09-03 21:58:39 -07:00
7 changed files with 495 additions and 98 deletions
+108
View File
@@ -0,0 +1,108 @@
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)"
}
+6 -6
View File
@@ -56,12 +56,7 @@ class ComfyUIDeployExternalImageBatch:
images_list = json.loads(images) # Assuming images is a JSON array string images_list = json.loads(images) # Assuming images is a JSON array string
print(images_list) print(images_list)
for img_input in images_list: for img_input in images_list:
if img_input.startswith('http'): if img_input.startswith('http') and img_input.endswith('.zip'):
from io import BytesIO
print("Fetching image from url: ", img_input)
response = requests.get(img_input)
image = Image.open(BytesIO(response.content))
elif img_input.startswith('http') and img_input.endswith('.zip'):
print("Fetching zip file from url: ", img_input) print("Fetching zip file from url: ", img_input)
response = requests.get(img_input) response = requests.get(img_input)
zip_file = zipfile.ZipFile(io.BytesIO(response.content)) zip_file = zipfile.ZipFile(io.BytesIO(response.content))
@@ -71,6 +66,11 @@ class ComfyUIDeployExternalImageBatch:
image = Image.open(file) image = Image.open(file)
image = self.process_image(image) image = self.process_image(image)
processed_images.append(image) processed_images.append(image)
elif img_input.startswith('http'):
from io import BytesIO
print("Fetching image from url: ", img_input)
response = requests.get(img_input)
image = Image.open(BytesIO(response.content))
elif img_input.startswith('data:image/png;base64,') or img_input.startswith('data:image/jpeg;base64,') or img_input.startswith('data:image/jpg;base64,'): elif img_input.startswith('data:image/png;base64,') or img_input.startswith('data:image/jpeg;base64,') or img_input.startswith('data:image/jpg;base64,'):
import base64 import base64
from io import BytesIO from io import BytesIO
+46
View File
@@ -0,0 +1,46 @@
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)"}
+213 -47
View File
@@ -24,7 +24,7 @@ from typing import Dict, List, Union, Any, Optional
from PIL import Image from PIL import Image
import copy import copy
import struct import struct
from aiohttp import web, ClientSession, ClientError, ClientTimeout from aiohttp import web, ClientSession, ClientError, ClientTimeout, ClientResponseError
import atexit import atexit
# Global session # Global session
@@ -59,7 +59,7 @@ print(f"max_retries: {max_retries}, retry_delay_multiplier: {retry_delay_multipl
import time import time
async def async_request_with_retry(method, url, disable_timeout=False, **kwargs): async def async_request_with_retry(method, url, disable_timeout=False, token=None, **kwargs):
global client_session global client_session
await ensure_client_session() await ensure_client_session()
retry_delay = 1 # Start with 1 second delay retry_delay = 1 # Start with 1 second delay
@@ -72,14 +72,19 @@ async def async_request_with_retry(method, url, disable_timeout=False, **kwargs)
timeout = ClientTimeout(total=None, connect=initial_timeout) timeout = ClientTimeout(total=None, connect=initial_timeout)
kwargs['timeout'] = 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() request_start = time.time()
async with client_session.request(method, url, **kwargs) as response: async with client_session.request(method, url, **kwargs) as response:
request_end = time.time() request_end = time.time()
logger.info(f"Request attempt {attempt + 1} took {request_end - request_start:.2f} seconds") # logger.info(f"Request attempt {attempt + 1} took {request_end - request_start:.2f} seconds")
if response.status != 200: if response.status != 200:
error_body = await response.text() error_body = await response.text()
logger.error(f"Request failed with status {response.status} and body {error_body}") # logger.error(f"Request failed with status {response.status} and body {error_body}")
# raise Exception(f"Request failed with status {response.status}") # raise Exception(f"Request failed with status {response.status}")
response.raise_for_status() response.raise_for_status()
@@ -87,7 +92,7 @@ async def async_request_with_retry(method, url, disable_timeout=False, **kwargs)
await response.read() await response.read()
total_time = time.time() - start_time total_time = time.time() - start_time
logger.info(f"Request succeeded after {total_time:.2f} seconds (attempt {attempt + 1}/{max_retries})") # logger.info(f"Request succeeded after {total_time:.2f} seconds (attempt {attempt + 1}/{max_retries})")
return response return response
except asyncio.TimeoutError: except asyncio.TimeoutError:
logger.warning(f"Request timed out after {initial_timeout} seconds (attempt {attempt + 1}/{max_retries})") logger.warning(f"Request timed out after {initial_timeout} seconds (attempt {attempt + 1}/{max_retries})")
@@ -188,6 +193,9 @@ bypass_upload = os.environ.get('CD_BYPASS_UPLOAD', 'false').lower() == 'true'
logger.info(f"CD_BYPASS_UPLOAD {bypass_upload}") logger.info(f"CD_BYPASS_UPLOAD {bypass_upload}")
create_native_run_endpoint = None
status_endpoint = None
file_upload_endpoint = None
def clear_current_prompt(sid): def clear_current_prompt(sid):
prompt_server = server.PromptServer.instance prompt_server = server.PromptServer.instance
@@ -304,7 +312,7 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
value['inputs']["input_id"] = new_value value['inputs']["input_id"] = new_value
# Fix for external text default value # Fix for external text default value
if (value["class_type"] == "ComfyUIDeployExternalText"): if (value["class_type"] == "ComfyUIDeployExternalText" or value["class_type"] == "ComfyUIDeployExternalTextAny"):
value['inputs']["default_value"] = new_value value['inputs']["default_value"] = new_value
if (value["class_type"] == "ComfyUIDeployExternalCheckpoint"): if (value["class_type"] == "ComfyUIDeployExternalCheckpoint"):
@@ -322,9 +330,13 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
if value["class_type"] == "ComfyUIDeployExternalBoolean": if value["class_type"] == "ComfyUIDeployExternalBoolean":
value["inputs"]["default_value"] = new_value 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): def send_prompt(sid: str, inputs: StreamingPrompt):
# workflow_api = inputs.workflow_api # workflow_api = inputs.workflow_api
workflow_api = copy.deepcopy(inputs.workflow_api) workflow_api = copy.deepcopy(inputs.workflow_api)
workflow = copy.deepcopy(inputs.workflow)
# Random seed # Random seed
apply_random_seed_to_workflow(workflow_api) apply_random_seed_to_workflow(workflow_api)
@@ -340,7 +352,8 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
prompt = { prompt = {
"prompt": workflow_api, "prompt": workflow_api,
"client_id": sid, #"comfy_deploy_instance", #api.client_id "client_id": sid, #"comfy_deploy_instance", #api.client_id
"prompt_id": prompt_id "prompt_id": prompt_id,
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
} }
try: try:
@@ -359,15 +372,55 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
logger.info(f"error: {error_type}, {e}") logger.info(f"error: {error_type}, {e}")
logger.info(f"stack trace: {stack_trace_short}") logger.info(f"stack trace: {stack_trace_short}")
# # Add custom logic here
# if 'prompt_id' in response:
# prompt_id = response['prompt_id']
# if prompt_id in prompt_metadata:
# metadata = prompt_metadata[prompt_id]
# # Add additional information to the response
# response['status_endpoint'] = metadata.status_endpoint
# response['file_upload_endpoint'] = metadata.file_upload_endpoint
return response
@server.PromptServer.instance.routes.post("/comfyui-deploy/run") @server.PromptServer.instance.routes.post("/comfyui-deploy/run")
async def comfy_deploy_run(request): async def comfy_deploy_run(request):
# Extract the bearer token from the Authorization header
data = await request.json() data = await request.json()
client_id = data.get("client_id")
# We proxy the request to Comfy Deploy, this is a native run
if "is_native_run" in data:
async with aiohttp.ClientSession() as session:
pprint(data)
# headers = request.headers.copy()
# headers['Content-Type'] = 'application/json'
async with session.post(data.get("native_run_api_endpoint"), json=data, headers={
'Content-Type': 'application/json',
'Authorization': request.headers.get('Authorization')
}) as response:
data = await response.json()
print(data)
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 # 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") workflow_api = data.get("workflow_api_raw")
# The prompt id generated from comfy deploy, can be None # The prompt id generated from comfy deploy, can be None
prompt_id = data.get("prompt_id") prompt_id = data.get("prompt_id")
inputs = data.get("inputs") inputs = data.get("inputs")
workflow = data.get("workflow")
# Now it handles directly in here # Now it handles directly in here
apply_random_seed_to_workflow(workflow_api) apply_random_seed_to_workflow(workflow_api)
@@ -375,14 +428,16 @@ async def comfy_deploy_run(request):
prompt = { prompt = {
"prompt": workflow_api, "prompt": workflow_api,
"client_id": "comfy_deploy_instance", #api.client_id "client_id": "comfy_deploy_instance" if client_id is None else client_id,
"prompt_id": prompt_id "prompt_id": prompt_id,
"extra_data": {"extra_pnginfo": {"workflow": workflow}}
} }
prompt_metadata[prompt_id] = SimplePrompt( prompt_metadata[prompt_id] = SimplePrompt(
status_endpoint=data.get('status_endpoint'), status_endpoint=data.get('status_endpoint'),
file_upload_endpoint=data.get('file_upload_endpoint'), file_upload_endpoint=data.get('file_upload_endpoint'),
workflow_api=workflow_api workflow_api=workflow_api,
token=token
) )
try: try:
@@ -420,12 +475,13 @@ async def comfy_deploy_run(request):
return web.json_response(res, status=status) return web.json_response(res, status=status)
async def stream_prompt(data): 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 # 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") workflow_api = data.get("workflow_api_raw")
# The prompt id generated from comfy deploy, can be None # The prompt id generated from comfy deploy, can be None
prompt_id = data.get("prompt_id") prompt_id = data.get("prompt_id")
inputs = data.get("inputs") inputs = data.get("inputs")
workflow = data.get("workflow")
# Now it handles directly in here # Now it handles directly in here
apply_random_seed_to_workflow(workflow_api) apply_random_seed_to_workflow(workflow_api)
@@ -434,13 +490,15 @@ async def stream_prompt(data):
prompt = { prompt = {
"prompt": workflow_api, "prompt": workflow_api,
"client_id": "comfy_deploy_instance", #api.client_id "client_id": "comfy_deploy_instance", #api.client_id
"prompt_id": prompt_id "prompt_id": prompt_id,
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
} }
prompt_metadata[prompt_id] = SimplePrompt( prompt_metadata[prompt_id] = SimplePrompt(
status_endpoint=data.get('status_endpoint'), status_endpoint=data.get('status_endpoint'),
file_upload_endpoint=data.get('file_upload_endpoint'), file_upload_endpoint=data.get('file_upload_endpoint'),
workflow_api=workflow_api workflow_api=workflow_api,
token=token
) )
# log('info', "Begin prompt", prompt=prompt) # log('info', "Begin prompt", prompt=prompt)
@@ -489,6 +547,14 @@ comfy_message_queues: Dict[str, asyncio.Queue] = {}
async def stream_response(request): async def stream_response(request):
response = web.StreamResponse(status=200, reason='OK', headers={'Content-Type': 'text/event-stream'}) response = web.StreamResponse(status=200, reason='OK', headers={'Content-Type': 'text/event-stream'})
await response.prepare(request) 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 pending = True
data = await request.json() data = await request.json()
@@ -500,7 +566,7 @@ async def stream_response(request):
log('info', 'Streaming prompt') log('info', 'Streaming prompt')
try: try:
result = await stream_prompt(data=data) result = await stream_prompt(data=data, token=token)
await response.write(f"event: event_update\ndata: {json.dumps(result)}\n\n".encode('utf-8')) await response.write(f"event: event_update\ndata: {json.dumps(result)}\n\n".encode('utf-8'))
# await response.write(.encode('utf-8')) # await response.write(.encode('utf-8'))
await response.drain() # Ensure the buffer is flushed await response.drain() # Ensure the buffer is flushed
@@ -759,6 +825,7 @@ async def websocket_handler(request):
inputs={}, inputs={},
status_endpoint=status_endpoint, status_endpoint=status_endpoint,
file_upload_endpoint=request.rel_url.query.get('file_upload_endpoint', None), 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) await update_realtime_run_status(realtime_id, status_endpoint, Status.RUNNING)
@@ -948,6 +1015,7 @@ async def send_json_override(self, event, data, sid=None):
"data": data "data": data
}) })
asyncio.create_task(update_run_ws_event(prompt_id, event, data))
# event_emitter.emit("send_json", { # event_emitter.emit("send_json", {
# "event": event, # "event": event,
# "data": data # "data": data
@@ -1045,6 +1113,7 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
return return
status_endpoint = prompt_metadata[prompt_id].status_endpoint status_endpoint = prompt_metadata[prompt_id].status_endpoint
token = prompt_metadata[prompt_id].token
if (status_endpoint is None): if (status_endpoint is None):
return return
@@ -1068,7 +1137,27 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
}) })
# requests.post(status_endpoint, json=body) # requests.post(status_endpoint, json=body)
await async_request_with_retry('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)
async def update_run(prompt_id: str, status: Status): async def update_run(prompt_id: str, status: Status):
@@ -1099,7 +1188,8 @@ async def update_run(prompt_id: str, status: Status):
try: try:
# requests.post(status_endpoint, json=body) # requests.post(status_endpoint, json=body)
if (status_endpoint is not None): if (status_endpoint is not None):
await async_request_with_retry('POST', status_endpoint, json=body) token = prompt_metadata[prompt_id].token
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
if (status_endpoint is not None) and cd_enable_run_log and (status == Status.SUCCESS or status == Status.FAILED): if (status_endpoint is not None) and cd_enable_run_log and (status == Status.SUCCESS or status == Status.FAILED):
try: try:
@@ -1129,7 +1219,7 @@ async def update_run(prompt_id: str, status: Status):
] ]
} }
await async_request_with_retry('POST', status_endpoint, json=body) await async_request_with_retry('POST', status_endpoint, token=token, json=body)
# requests.post(status_endpoint, json=body) # requests.post(status_endpoint, json=body)
except Exception as log_error: except Exception as log_error:
logger.info(f"Error reading log file: {log_error}") logger.info(f"Error reading log file: {log_error}")
@@ -1150,6 +1240,47 @@ 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", item=None):
""" """
Uploads file to S3 bucket using S3 client object Uploads file to S3 bucket using S3 client object
@@ -1180,45 +1311,60 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
logger.info(f"Uploading file {file}") logger.info(f"Uploading file {file}")
file_upload_endpoint = prompt_metadata[prompt_id].file_upload_endpoint file_upload_endpoint = prompt_metadata[prompt_id].file_upload_endpoint
token = prompt_metadata[prompt_id].token
filename = quote(filename) filename = quote(filename)
prompt_id = quote(prompt_id) prompt_id = quote(prompt_id)
content_type = quote(content_type) content_type = quote(content_type)
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)
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}")
async with aiofiles.open(file, 'rb') as f: async with aiofiles.open(file, 'rb') as f:
data = await f.read() data = await f.read()
size = str(len(data)) size = str(len(data))
target_url = f"{file_upload_endpoint}?file_name={filename}&run_id={prompt_id}&type={content_type}&version=v2" # logger.info(f"Image size: {size}")
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)
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 start_time = time.time() # Start timing here
headers = { headers = {
"Content-Type": content_type, "Content-Type": content_type,
# "Content-Length": size, "Content-Length": size,
} }
logger.info(headers)
if ok.get('include_acl') is True: if ok.get('include_acl') is True:
headers["x-amz-acl"] = "public-read" headers["x-amz-acl"] = "public-read"
# response = requests.put(ok.get("url"), headers=headers, data=data) # response = requests.put(ok.get("url"), headers=headers, data=data)
response = await async_request_with_retry('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}") # 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)}")
end_time = time.time() # End timing after the request is complete end_time = time.time() # End timing after the request is complete
logger.info("Upload time: {:.2f} seconds".format(end_time - start_time)) logger.info("Upload time: {:.2f} seconds".format(end_time - start_time))
if item is not None: if item is not None:
file_download_url = ok.get("download_url") file_download_url = ok.get("download_url")
if file_download_url is not None: if file_download_url is not None:
item["url"] = file_download_url item["url"] = file_download_url
item["upload_duration"] = end_time - start_time 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): def have_pending_upload(prompt_id):
if prompt_id in prompt_metadata and len(prompt_metadata[prompt_id].uploading_nodes) > 0: if prompt_id in prompt_metadata and len(prompt_metadata[prompt_id].uploading_nodes) > 0:
@@ -1335,6 +1481,14 @@ async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, d
content_type=file_type, content_type=file_type,
item=item item=item
)) ))
# 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 # Execute all upload tasks concurrently
await asyncio.gather(*upload_tasks) await asyncio.gather(*upload_tasks)
@@ -1342,16 +1496,23 @@ async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, d
# Upload files in the background # 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, node_meta=None):
try: 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 = [ upload_tasks = [
handle_upload(prompt_id, data, 'images', "content_type", "image/png"), 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, "files", "content_type", "image/png"),
handle_upload(prompt_id, data, 'gifs', "format", "image/gif"), handle_upload(prompt_id, data, "gifs", "format", "image/gif"),
handle_upload(prompt_id, data, 'mesh', "format", "application/octet-stream") handle_upload(
prompt_id, data, "mesh", "format", "application/octet-stream"
),
] ]
await asyncio.gather(*upload_tasks) await asyncio.gather(*upload_tasks)
status_endpoint = prompt_metadata[prompt_id].status_endpoint status_endpoint = prompt_metadata[prompt_id].status_endpoint
token = prompt_metadata[prompt_id].token
if have_upload: if have_upload:
if status_endpoint is not None: if status_endpoint is not None:
body = { body = {
@@ -1360,7 +1521,7 @@ async def upload_in_background(prompt_id: str, data, node_id=None, have_upload=T
"node_meta": node_meta, "node_meta": node_meta,
} }
# pprint(body) # pprint(body)
await async_request_with_retry('POST', status_endpoint, json=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) await update_file_status(prompt_id, data, False, node_id=node_id)
except Exception as e: except Exception as e:
await handle_error(prompt_id, data, e) await handle_error(prompt_id, data, e)
@@ -1379,7 +1540,10 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
"output_data": data, "output_data": data,
"node_meta": node_meta, "node_meta": node_meta,
} }
have_upload_media = 'images' in data or 'files' in data or 'gifs' in data or 'mesh' in 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
if bypass_upload and have_upload_media: if bypass_upload and have_upload_media:
print("CD_BYPASS_UPLOAD is enabled, skipping the upload of the output:", node_id) print("CD_BYPASS_UPLOAD is enabled, skipping the upload of the output:", node_id)
return return
@@ -1391,14 +1555,16 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
if have_upload_media: if have_upload_media:
await update_file_status(prompt_id, data, True, node_id=node_id) 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)) # 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)
# await upload_in_background(prompt_id, data, node_id=node_id, have_upload=have_upload) # await upload_in_background(prompt_id, data, node_id=node_id, have_upload=have_upload)
except Exception as e: except Exception as e:
await handle_error(prompt_id, data, e) await handle_error(prompt_id, data, e)
# requests.post(status_endpoint, json=body) # requests.post(status_endpoint, json=body)
elif status_endpoint is not None: elif status_endpoint is not None:
await async_request_with_retry('POST', status_endpoint, json=body) token = prompt_metadata[prompt_id].token
await async_request_with_retry('POST', status_endpoint, token=token, json=body)
await send('outputs_uploaded', { await send('outputs_uploaded', {
"prompt_id": prompt_id "prompt_id": prompt_id
+3
View File
@@ -24,11 +24,14 @@ class StreamingPrompt(BaseModel):
running_prompt_ids: set[str] = set() running_prompt_ids: set[str] = set()
status_endpoint: Optional[str] status_endpoint: Optional[str]
file_upload_endpoint: Optional[str] file_upload_endpoint: Optional[str]
workflow: Any
class SimplePrompt(BaseModel): class SimplePrompt(BaseModel):
status_endpoint: Optional[str] status_endpoint: Optional[str]
file_upload_endpoint: Optional[str] file_upload_endpoint: Optional[str]
token: Optional[str]
workflow_api: dict workflow_api: dict
status: Status = Status.NOT_STARTED status: Status = Status.NOT_STARTED
progress: set = set() progress: set = set()
+118 -45
View File
@@ -83,6 +83,23 @@ function dispatchAPIEventData(data) {
} }
} }
const context = {
selectedWorkflowInfo: null,
};
// let selectedWorkflowInfo = {
// workflow_id: "05da8f2b-63af-4c0c-86dd-08d01ec512b7",
// machine_id: "45ac5f85-b7b6-436f-8d97-2383b25485f3",
// native_run_api_endpoint: "http://localhost:3011/api/run",
// };
function getSelectedWorkflowInfo() {
return context.selectedWorkflowInfo;
}
function setSelectedWorkflowInfo(info) {
context.selectedWorkflowInfo = info;
}
/** @typedef {import('../../../web/types/comfy.js').ComfyExtension} ComfyExtension*/ /** @typedef {import('../../../web/types/comfy.js').ComfyExtension} ComfyExtension*/
/** @type {ComfyExtension} */ /** @type {ComfyExtension} */
const ext = { const ext = {
@@ -103,10 +120,10 @@ const ext = {
sendEventToCD("cd_plugin_onInit"); sendEventToCD("cd_plugin_onInit");
app.queuePrompt = ((originalFunction) => async () => { // app.queuePrompt = ((originalFunction) => async () => {
// const prompt = await app.graphToPrompt(); // // const prompt = await app.graphToPrompt();
sendEventToCD("cd_plugin_onQueuePromptTrigger"); // sendEventToCD("cd_plugin_onQueuePromptTrigger");
})(app.queuePrompt); // })(app.queuePrompt);
// // Intercept the onkeydown event // // Intercept the onkeydown event
// window.addEventListener( // window.addEventListener(
@@ -192,11 +209,13 @@ const ext = {
registerCustomNodes() { registerCustomNodes() {
/** @type {LGraphNode}*/ /** @type {LGraphNode}*/
class ComfyDeploy { class ComfyDeploy extends LGraphNode {
color = LGraphCanvas.node_colors.yellow.color;
bgcolor = LGraphCanvas.node_colors.yellow.bgcolor;
groupcolor = LGraphCanvas.node_colors.yellow.groupcolor;
constructor() { 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) { if (!this.properties) {
this.properties = {}; this.properties = {};
this.properties.workflow_name = ""; this.properties.workflow_name = "";
@@ -204,62 +223,75 @@ const ext = {
this.properties.version = ""; this.properties.version = "";
} }
ComfyWidgets.STRING( this.addWidget(
this, "text",
"workflow_name", "workflow_name",
[ this.properties.workflow_name,
"", (v) => {
{ this.properties.workflow_name = v;
default: this.properties.workflow_name, },
multiline: false, { multiline: false },
},
],
app,
); );
ComfyWidgets.STRING( this.addWidget(
this, "text",
"workflow_id", "workflow_id",
[ this.properties.workflow_id,
"", (v) => {
{ this.properties.workflow_id = v;
default: this.properties.workflow_id, },
multiline: false, { multiline: false },
},
],
app,
); );
ComfyWidgets.STRING( this.addWidget(
this, "text",
"version", "version",
["", { default: this.properties.version, multiline: false }], this.properties.version,
app, (v) => {
this.properties.version = v;
},
{ multiline: false },
); );
// this.widgets.forEach((w) => {
// // w.computeSize = () => [200,10]
// w.computedHeight = 2;
// })
this.widgets_start_y = 10; this.widgets_start_y = 10;
this.setSize(this.computeSize());
// const config = { };
// console.log(this);
this.serialize_widgets = true; this.serialize_widgets = true;
this.isVirtualNode = 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";
}
}
} }
// Load default visibility // Register the node type
LiteGraph.registerNodeType( LiteGraph.registerNodeType(
"ComfyDeploy", "ComfyDeploy",
Object.assign(ComfyDeploy, { Object.assign(ComfyDeploy, {
title_mode: LiteGraph.NORMAL_TITLE,
title: "Comfy Deploy", title: "Comfy Deploy",
title_mode: LiteGraph.NORMAL_TITLE,
collapsable: true, collapsable: true,
}), }),
); );
@@ -281,6 +313,14 @@ const ext = {
// This part of the code would depend on how the ComfyUI expects to receive and process the workflow data // This part of the code would depend on how the ComfyUI expects to receive and process the workflow data
// For demonstration, let's assume there's a loadWorkflow method in the ComfyUI API // For demonstration, let's assume there's a loadWorkflow method in the ComfyUI API
if (comfyUIWorkflow && app && app.loadGraphData) { if (comfyUIWorkflow && app && app.loadGraphData) {
try {
await window["app"].ui.settings.setSettingValueAsync(
"Comfy.Validation.Workflows",
false,
);
} catch (error) {
console.warning("Error setting validation to false, is fine to ignore this", error);
}
console.log("loadGraphData"); console.log("loadGraphData");
app.loadGraphData(comfyUIWorkflow); app.loadGraphData(comfyUIWorkflow);
} }
@@ -368,6 +408,8 @@ const ext = {
} }
animate(); animate();
} else if (message.type === "workflow_info") {
setSelectedWorkflowInfo(message.data);
} }
// else if (message.type === "refresh") { // else if (message.type === "refresh") {
// sendEventToCD("cd_plugin_onRefresh"); // sendEventToCD("cd_plugin_onRefresh");
@@ -1482,3 +1524,34 @@ async function loadWorkflowApi(versionId) {
// Show an error message to the user // Show an error message to the user
} }
} }
const orginal_fetch_api = api.fetchApi;
api.fetchApi = async (route, options) => {
console.log("Fetch API called with args:", route, options);
const info = getSelectedWorkflowInfo();
if (info && route.startsWith("/prompt")) {
const body = JSON.parse(options.body);
const data = {
client_id: body.client_id,
workflow_api_json: body.prompt,
workflow: body?.extra_data?.extra_pnginfo?.workflow,
is_native_run: true,
machine_id: info.machine_id,
workflow_id: info.workflow_id,
native_run_api_endpoint: info.native_run_api_endpoint,
};
return await fetch("/comfyui-deploy/run", {
method: "POST",
headers: {
Authorization: `Bearer ${info.cd_token}`,
"Content-Type": "application/json",
},
body: JSON.stringify(data),
});
}
return await orginal_fetch_api.call(api, route, options);
};
+1
View File
@@ -6,4 +6,5 @@ export const customInputNodes: Record<string, string> = {
ComfyUIDeployExternalNumberInt: "integer", ComfyUIDeployExternalNumberInt: "integer",
ComfyUIDeployExternalLora: "string - (public lora download url)", ComfyUIDeployExternalLora: "string - (public lora download url)",
ComfyUIDeployExternalCheckpoint: "string - (public checkpoints download url)", ComfyUIDeployExternalCheckpoint: "string - (public checkpoints download url)",
ComfyUIDeployExternalFaceModel: "string - (public face model download url)",
}; };