Compare commits

..
Author SHA1 Message Date
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 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 e011711600 feat: zip batch image support 2024-09-09 17:49:39 -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
bennykok 503dca8fb6 chore: add log 2024-08-30 12:16:41 -07:00
bennykok 73c149b4cb fix node meta 2024-08-30 12:16:41 -07:00
bennykok 65b5b0b8c7 fix: remove content length 2024-08-30 12:16:41 -07:00
bennykok 9d6ee85402 fix: upload file acl 2024-08-30 12:16:41 -07:00
bennykok cdaed8a571 fix: include upload time 2024-08-30 12:16:41 -07:00
bennykok 3129e89cce fix: log file error log 2024-08-30 12:16:41 -07:00
bennykok 7a693eabc8 fix: size 2024-08-30 12:16:41 -07:00
bennykok 8f677e520d chore: log more test for upload file debug 2024-08-30 12:16:41 -07:00
bennykok 4c8d32c5b0 fix 2024-08-30 12:16:41 -07:00
nick a99d2568e0 video and lora node fix 2024-08-28 13:08:15 -07:00
nick 649b61c580 default vid 2024-08-26 13:46:01 -07:00
nick edff5685f9 fix: random seed 2024-08-22 17:39:03 -07:00
bennykok 9fc0c2b4a2 chore: upload node data 2024-08-21 16:34:25 -07:00
bennykok d34e2e99b1 fix: external lora for new comfyui 2024-08-21 09:46:13 -07:00
bennykok f85043db07 fix: remove default value 2024-08-20 19:14:43 -07:00
bennykok 894d8e1503 Merge branch 'benny/async-upload-file' into public-main 2024-08-20 18:02:57 -07:00
bennykok 08d631d1eb feat: async file upload for the same node 2024-08-20 17:07:50 -07:00
karrix a1031487e1 add: all node support name and description 2024-08-20 20:15:29 +08:00
bennykok ca41207192 feat: max min int for all number inputs to enable negative number input 2024-08-19 13:27:46 -07:00
bennykok 507d5ef631 feat: add a init timeout of 10 seconds for retry logic 2024-08-18 17:31:48 -07:00
bennykok dd1d9df23f fix: resolve false possible error 2024-08-18 15:38:16 -07:00
bennykok 3a14e49ca5 fix: refresh workflows list 2024-08-17 16:04:14 -07:00
nick 8147c4bfb7 video node' 2024-08-15 12:50:29 -07:00
bennykok 10268825d9 feat: support new frontend! 2024-08-14 11:09:58 -07:00
bennykok f6ea252652 fix: log when random seed is applied 2024-08-10 10:35:48 -07:00
bennykok 98cd5ef79c fix: randomize noise RandomNoise, KSamplerAdvanced, SamplerCustom 2024-08-10 10:02:01 -07:00
Emmanuel Morales 4bce5cadfb fix(text): return correctly the text in external_text_list node 2024-08-10 09:44:37 -06:00
Nick Kao f362671041 Merge pull request #61 from BennyKok/node-error-no-throw
block on bad prompt
2024-08-08 10:01:33 -07:00
16 changed files with 1620 additions and 1135 deletions
+11 -1
View File
@@ -8,6 +8,16 @@ class ComfyUIDeployExternalBoolean:
{"multiline": False, "default": "input_bool"},
),
"default_value": ("BOOLEAN", {"default": False})
},
"optional": {
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -16,7 +26,7 @@ class ComfyUIDeployExternalBoolean:
FUNCTION = "run"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
print(f"Node '{input_id}' processing with switch set to {default_value}")
return [default_value]
+9 -1
View File
@@ -23,6 +23,14 @@ class ComfyUIDeployExternalCheckpoint:
},
"optional": {
"default_value": (folder_paths.get_filename_list("checkpoints"), ),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -33,7 +41,7 @@ class ComfyUIDeployExternalCheckpoint:
CATEGORY = "deploy"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
import requests
import os
import uuid
+9 -1
View File
@@ -15,6 +15,14 @@ class ComfyUIDeployExternalImage:
},
"optional": {
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -25,7 +33,7 @@ class ComfyUIDeployExternalImage:
CATEGORY = "image"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
image = default_value
try:
if input_id.startswith('http'):
+9 -1
View File
@@ -15,6 +15,14 @@ class ComfyUIDeployExternalImageAlpha:
},
"optional": {
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -25,7 +33,7 @@ class ComfyUIDeployExternalImageAlpha:
CATEGORY = "image"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
image = default_value
try:
if input_id.startswith('http'):
+31 -3
View File
@@ -21,6 +21,14 @@ class ComfyUIDeployExternalImageBatch:
},
"optional": {
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -31,14 +39,34 @@ class ComfyUIDeployExternalImageBatch:
CATEGORY = "image"
def run(self, input_id, images=None, default_value=None):
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'):
import requests
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'):
from io import BytesIO
print("Fetching image from url: ", img_input)
response = requests.get(img_input)
+26 -6
View File
@@ -25,10 +25,22 @@ 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": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
"lora_url": (
"STRING",
{"multiline": False, "default": ""},
),
},
}
@@ -39,12 +51,20 @@ class ComfyUIDeployExternalLora:
CATEGORY = "deploy"
def run(self, input_id, default_lora_name=None, lora_save_name=None):
def run(
self,
input_id,
default_lora_name=None,
lora_save_name=None,
display_name=None,
description=None,
lora_url=None,
):
import requests
import os
import uuid
if default_lora_name.startswith("http"):
if lora_url and lora_url.startswith("http"):
if lora_save_name:
existing_loras = folder_paths.get_filename_list("loras")
# Check if lora_save_name exists in the list
@@ -59,9 +79,9 @@ class ComfyUIDeployExternalLora:
folder_paths.folder_names_and_paths["loras"][0][0], lora_save_name
)
print(destination_path)
print("Downloading external lora - " + input_id + " to " + destination_path)
print("Downloading external lora - " + lora_url + " to " + destination_path)
response = requests.get(
input_id,
lora_url,
headers={"User-Agent": "Mozilla/5.0"},
allow_redirects=True,
)
+10 -2
View File
@@ -16,7 +16,15 @@ class ComfyUIDeployExternalNumber:
"optional": {
"default_value": (
"FLOAT",
{"multiline": True, "display": "number", "default": 0, "step": 0.01},
{"multiline": True, "display": "number", "default": 0, "min": -2147483647, "max": 2147483647, "step": 0.01},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -28,7 +36,7 @@ class ComfyUIDeployExternalNumber:
CATEGORY = "number"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
try:
float_value = float(input_id)
print("my number", float_value)
+10 -2
View File
@@ -16,7 +16,15 @@ class ComfyUIDeployExternalNumberInt:
"optional": {
"default_value": (
"INT",
{"multiline": True, "display": "number", "default": 0},
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 0},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -28,7 +36,7 @@ class ComfyUIDeployExternalNumberInt:
CATEGORY = "number"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
if not input_id or (isinstance(input_id, str) and not input_id.strip().isdigit()):
return [default_value]
return [int(input_id)]
+12 -4
View File
@@ -11,15 +11,23 @@ class ComfyUIDeployExternalNumberSlider:
"optional": {
"default_value": (
"FLOAT",
{"multiline": True, "display": "number", "default": 0.5, "step": 0.01},
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 0.5, "step": 0.01},
),
"min_value": (
"FLOAT",
{"multiline": True, "display": "number", "default": 0, "step": 0.01},
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 0, "step": 0.01},
),
"max_value": (
"FLOAT",
{"multiline": True, "display": "number", "default": 1, "step": 0.01},
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 1, "step": 0.01},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -31,7 +39,7 @@ class ComfyUIDeployExternalNumberSlider:
CATEGORY = "number"
def run(self, input_id, default_value=None, min_value=0, max_value=1):
def run(self, input_id, default_value=None, min_value=0, max_value=1, display_name=None, description=None):
try:
float_value = float(input_id)
if min_value <= float_value <= max_value:
+9 -1
View File
@@ -18,6 +18,14 @@ class ComfyUIDeployExternalText:
"STRING",
{"multiline": True, "default": ""},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -28,7 +36,7 @@ class ComfyUIDeployExternalText:
CATEGORY = "text"
def run(self, input_id, default_value=None):
def run(self, input_id, default_value=None, display_name=None, description=None):
return [default_value]
+12 -3
View File
@@ -17,6 +17,16 @@ class ComfyUIDeployExternalTextList:
"STRING",
{"multiline": True, "default": "[]"},
),
},
"optional": {
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -29,15 +39,14 @@ class ComfyUIDeployExternalTextList:
CATEGORY = "text"
def run(self, input_id, text=None):
def run(self, input_id, text=None, display_name=None, description=None):
text_list = []
try:
text_list = json.loads(text) # Assuming text is a JSON array string
except Exception as e:
print(f"Error processing images: {e}")
pass
return [text_list]
return ([text_list],)
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalTextList": ComfyUIDeployExternalTextList}
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyUIDeployExternalTextList": "External Text List (ComfyUI Deploy)"}
+14 -5
View File
@@ -764,7 +764,15 @@ class ComfyUIDeployExternalVideo:
"optional": {
"meta_batch": ("VHS_BatchManager",),
"vae": ("VAE",),
"default_value": (sorted(files),),
"default_video": (sorted(files),),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
},
"hidden": {
"unique_id": "UNIQUE_ID"
@@ -796,8 +804,6 @@ class ComfyUIDeployExternalVideo:
meta_batch = kwargs.get("meta_batch")
unique_id = kwargs.get("unique_id")
video = kwargs.get("default_value")
video_path = folder_paths.get_annotated_filepath(video.strip('"'))
input_dir = folder_paths.get_input_directory()
if input_id.startswith("http"):
@@ -827,8 +833,11 @@ class ComfyUIDeployExternalVideo:
leave=True,
):
out_file.write(chunk)
print("video path: ", video_path)
else:
video = kwargs.get("default_video", None)
if video is None:
raise "No default video given and no external video provided"
video_path = folder_paths.get_annotated_filepath(video.strip('"'))
return load_video_cv(
video=video_path,
+247 -53
View File
@@ -1,4 +1,5 @@
from io import BytesIO
from pprint import pprint
from aiohttp import web
import os
import requests
@@ -17,12 +18,13 @@ from urllib.parse import quote
import threading
import hashlib
import aiohttp
from aiohttp import ClientSession, web
import aiofiles
from typing import Dict, List, Union, Any, Optional
from PIL import Image
import copy
import struct
from aiohttp import ClientError
from aiohttp import web, ClientSession, ClientError, ClientTimeout
import atexit
# Global session
@@ -50,28 +52,70 @@ def exit_handler():
atexit.register(exit_handler)
max_retries = int(os.environ.get('MAX_RETRIES', '3'))
max_retries = int(os.environ.get('MAX_RETRIES', '5'))
retry_delay_multiplier = float(os.environ.get('RETRY_DELAY_MULTIPLIER', '2'))
print(f"max_retries: {max_retries}, retry_delay_multiplier: {retry_delay_multiplier}")
async def async_request_with_retry(method, url, **kwargs):
import time
async def async_request_with_retry(method, url, disable_timeout=False, token=None, **kwargs):
global client_session
await ensure_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:
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
logger.warning(f"Request failed (attempt {attempt + 1}/{max_retries}): {e}")
await asyncio.sleep(retry_delay)
retry_delay *= retry_delay_multiplier # Exponential backoff
raise
await asyncio.sleep(retry_delay)
retry_delay *= retry_delay_multiplier
total_time = time.time() - start_time
raise Exception(f"Request failed after {max_retries} attempts and {total_time:.2f} seconds")
from logging import basicConfig, getLogger
@@ -217,13 +261,30 @@ def apply_random_seed_to_workflow(workflow_api):
workflow_api (dict): The workflow API dictionary to modify.
"""
for key in workflow_api:
if 'inputs' in workflow_api[key] and 'seed' in workflow_api[key]['inputs']:
if isinstance(workflow_api[key]['inputs']['seed'], list):
continue
if workflow_api[key]['class_type'] == "PromptExpansion":
workflow_api[key]['inputs']['seed'] = randomSeed(8);
continue
workflow_api[key]['inputs']['seed'] = randomSeed();
if 'inputs' in workflow_api[key]:
if 'seed' in workflow_api[key]['inputs']:
if isinstance(workflow_api[key]['inputs']['seed'], list):
continue
if workflow_api[key]['class_type'] == "PromptExpansion":
workflow_api[key]['inputs']['seed'] = randomSeed(8)
logger.info(f"Applied random seed {workflow_api[key]['inputs']['seed']} to PromptExpansion")
continue
workflow_api[key]['inputs']['seed'] = randomSeed()
logger.info(f"Applied random seed {workflow_api[key]['inputs']['seed']} to {workflow_api[key]['class_type']}")
if 'noise_seed' in workflow_api[key]['inputs']:
if workflow_api[key]['class_type'] == "RandomNoise":
workflow_api[key]['inputs']['noise_seed'] = randomSeed()
logger.info(f"Applied random noise_seed {workflow_api[key]['inputs']['noise_seed']} to RandomNoise")
continue
if workflow_api[key]['class_type'] == "KSamplerAdvanced":
workflow_api[key]['inputs']['noise_seed'] = randomSeed()
logger.info(f"Applied random noise_seed {workflow_api[key]['inputs']['noise_seed']} to KSamplerAdvanced")
continue
if workflow_api[key]['class_type'] == "SamplerCustom":
workflow_api[key]['inputs']['noise_seed'] = randomSeed()
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):
# Loop through each of the inputs and replace them
@@ -258,7 +319,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"]["default_lora_name"] = new_value
value["inputs"]["lora_url"] = new_value
if value["class_type"] == "ComfyUIDeployExternalSlider":
value["inputs"]["default_value"] = new_value
@@ -305,6 +366,14 @@ 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
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]
data = await request.json()
# In older version, we use workflow_api, but this has inputs already swapped in nextjs frontend, which is tricky
@@ -320,13 +389,14 @@ async def comfy_deploy_run(request):
prompt = {
"prompt": workflow_api,
"client_id": "comfy_deploy_instance", #api.client_id
"prompt_id": prompt_id
"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
workflow_api=workflow_api,
token=token
)
try:
@@ -349,7 +419,7 @@ async def comfy_deploy_run(request):
status = 200
if "node_errors" in res and res["node_errors"] is not None:
if "node_errors" in res and res["node_errors"] is not None and len(res["node_errors"]) > 0:
# Even tho there are node_errors it can still be run
status = 400
await update_run_with_output(prompt_id, {
@@ -364,7 +434,7 @@ async def comfy_deploy_run(request):
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
workflow_api = data.get("workflow_api_raw")
# The prompt id generated from comfy deploy, can be None
@@ -384,7 +454,8 @@ async def stream_prompt(data):
prompt_metadata[prompt_id] = SimplePrompt(
status_endpoint=data.get('status_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)
@@ -410,7 +481,7 @@ async def stream_prompt(data):
status = 200
if "node_errors" in res and res["node_errors"] is not None:
if "node_errors" in res and res["node_errors"] is not None and len(res["node_errors"]) > 0:
# Even tho there are node_errors it can still be run
status = 400
await update_run_with_output(prompt_id, {
@@ -433,6 +504,14 @@ 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()
@@ -444,7 +523,7 @@ async def stream_response(request):
log('info', 'Streaming prompt')
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(.encode('utf-8'))
await response.drain() # Ensure the buffer is flushed
@@ -573,9 +652,10 @@ async def upload_file_endpoint(request):
with open(file_path, 'rb') as f:
headers = {
"Content-Type": file_type,
# "x-amz-acl": "public-read",
"Content-Length": str(file_size)
# "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({
@@ -821,7 +901,51 @@ 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:
# Create a new ClientSession for each request
async with ClientSession() as client_session:
# Forward the request
client_req = await client_session.request(
method=request.method,
url=target_url,
headers={k: v for k, v in request.headers.items() if k.lower() not in ('host', 'content-length')},
data=await request.read(),
allow_redirects=False,
)
# Read the entire response content
content = await client_req.read()
# Try to decode the content as JSON
try:
json_data = json.loads(content)
# If successful, return a JSON response
return web.json_response(json_data, status=client_req.status)
except json.JSONDecodeError:
# If it's not valid JSON, return the content as-is
return web.Response(body=content, status=client_req.status, headers=client_req.headers)
except ClientError as e:
print(f"Client error occurred while proxying request: {str(e)}")
return web.Response(status=502, text=f"Bad Gateway: {str(e)}")
except Exception as e:
print(f"Error occurred while proxying request: {str(e)}")
return web.Response(status=500, text=f"Internal Server Error: {str(e)}")
prompt_server = server.PromptServer.instance
@@ -847,6 +971,7 @@ 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
@@ -914,17 +1039,22 @@ 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'))
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'))
# update_run_with_output(prompt_id, data.get('output'))
@@ -939,6 +1069,7 @@ 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
@@ -962,7 +1093,27 @@ 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, 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):
@@ -993,7 +1144,8 @@ async def update_run(prompt_id: str, status: Status):
try:
# requests.post(status_endpoint, json=body)
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):
try:
@@ -1023,7 +1175,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)
except Exception as log_error:
logger.info(f"Error reading log file: {log_error}")
@@ -1044,7 +1196,7 @@ async def update_run(prompt_id: str, status: Status):
})
async def upload_file(prompt_id, filename, subfolder=None, content_type="image/png", type="output"):
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
:return: None
@@ -1074,33 +1226,47 @@ 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)
target_url = f"{file_upload_endpoint}?file_name={filename}&run_id={prompt_id}&type={content_type}"
async with aiofiles.open(file, 'rb') as f:
data = await f.read()
size = str(len(data))
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
result = requests.get(target_url)
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 = result.json()
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}")
start_time = time.time() # Start timing here
with open(file, 'rb') as f:
data = f.read()
start_time = time.time() # Start timing here
headers = {
# "x-amz-acl": "public-read",
"Content-Type": content_type,
"Content-Length": str(len(data)),
# "Content-Length": size,
}
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}")
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:
@@ -1193,8 +1359,10 @@ async def update_file_status(prompt_id: str, data, uploading, have_error=False,
async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, default_content_type: str):
items = data.get(key, [])
# upload_tasks = []
for item in items:
# # Skipping temp files
# Skipping temp files
if item.get("type") == "temp":
continue
@@ -1207,29 +1375,50 @@ async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, d
elif file_extension == '.webp':
file_type = 'image/webp'
# upload_tasks.append(upload_file(
# prompt_id,
# item.get("filename"),
# subfolder=item.get("subfolder"),
# type=item.get("type"),
# content_type=file_type,
# item=item
# ))
await upload_file(
prompt_id,
item.get("filename"),
subfolder=item.get("subfolder"),
type=item.get("type"),
content_type=file_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):
async def upload_in_background(prompt_id: str, data, node_id=None, have_upload=True, node_meta=None):
try:
await handle_upload(prompt_id, data, 'images', "content_type", "image/png")
await handle_upload(prompt_id, data, 'files', "content_type", "image/png")
# This will also be mp4
await handle_upload(prompt_id, data, 'gifs', "format", "image/gif")
await handle_upload(prompt_id, data, 'mesh', "format", "application/octet-stream")
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):
async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
if prompt_id not in prompt_metadata:
return
@@ -1240,9 +1429,13 @@ async def update_run_with_output(prompt_id, data, node_id=None):
body = {
"run_id": prompt_id,
"output_data": data
"output_data": data,
"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:
print("CD_BYPASS_UPLOAD is enabled, skipping the upload of the output:", node_id)
return
@@ -1254,15 +1447,16 @@ async def update_run_with_output(prompt_id, data, node_id=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))
# 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)
except Exception as e:
await handle_error(prompt_id, data, e)
# requests.post(status_endpoint, json=body)
if status_endpoint is not None:
await async_request_with_retry('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)
await send('outputs_uploaded', {
"prompt_id": prompt_id
+2
View File
@@ -29,6 +29,8 @@ 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()
+1
View File
@@ -2,4 +2,5 @@ aiofiles
pydantic
opencv-python
imageio-ffmpeg
brotli
# logfire
+1208 -1052
View File
File diff suppressed because it is too large Load Diff