Compare commits
40
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
177010b0d4 | ||
|
|
797180b5c7 | ||
|
|
d00ca375a2 | ||
|
|
be5d5d2b54 | ||
|
|
d592a6ba12 | ||
|
|
35fed9aa4d | ||
|
|
3b6a753472 | ||
|
|
7d2c521645 | ||
|
|
f363b7e871 | ||
|
|
1b25cfdd6c | ||
|
|
5da56b5507 | ||
|
|
03d12e4099 | ||
|
|
e66712425d | ||
|
|
81f315e14d | ||
|
|
7189f13263 | ||
|
|
e73392ba8b | ||
|
|
1bfbd91708 | ||
|
|
a640e1eb79 | ||
|
|
011d36edce | ||
|
|
3df549c25c | ||
|
|
619a9728c0 | ||
|
|
410d03cd2b | ||
|
|
32c6d1215b | ||
|
|
9e79c434a9 | ||
|
|
19511e55ba | ||
|
|
2d59fd2b1b | ||
|
|
542b72bde5 | ||
|
|
7b653201ae | ||
|
|
1c9c32e9e4 | ||
|
|
97096a9035 | ||
|
|
e87bb63c6f | ||
|
|
a643fa0999 | ||
|
|
cc31840d41 | ||
|
|
25e62af24c | ||
|
|
9d0ded7ecc | ||
|
|
ec620dbc53 | ||
|
|
45d37879c2 | ||
|
|
ddbf6848a7 | ||
|
|
4ce2c98ae9 | ||
|
|
6e068590a0 |
@@ -16,7 +16,7 @@ class ComfyUIDeployExternalCheckpoint:
|
|||||||
),
|
),
|
||||||
},
|
},
|
||||||
"optional": {
|
"optional": {
|
||||||
"default_checkpoint_name": (folder_paths.get_filename_list("checkpoints"), ),
|
"default_value": (folder_paths.get_filename_list("checkpoints"), ),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -27,12 +27,12 @@ class ComfyUIDeployExternalCheckpoint:
|
|||||||
|
|
||||||
CATEGORY = "deploy"
|
CATEGORY = "deploy"
|
||||||
|
|
||||||
def run(self, input_id, default_checkpoint_name=None):
|
def run(self, input_id, default_value=None):
|
||||||
import requests
|
import requests
|
||||||
import os
|
import os
|
||||||
import uuid
|
import uuid
|
||||||
|
|
||||||
if input_id and input_id.startswith('http'):
|
if default_value.startswith('http'):
|
||||||
unique_filename = str(uuid.uuid4()) + ".safetensors"
|
unique_filename = str(uuid.uuid4()) + ".safetensors"
|
||||||
print(unique_filename)
|
print(unique_filename)
|
||||||
print(folder_paths.folder_names_and_paths["checkpoints"][0][0])
|
print(folder_paths.folder_names_and_paths["checkpoints"][0][0])
|
||||||
@@ -59,7 +59,7 @@ class ComfyUIDeployExternalCheckpoint:
|
|||||||
out_file.write(chunk)
|
out_file.write(chunk)
|
||||||
return (unique_filename,)
|
return (unique_filename,)
|
||||||
else:
|
else:
|
||||||
return (default_checkpoints_name,)
|
return (default_value,)
|
||||||
|
|
||||||
|
|
||||||
NODE_CLASS_MAPPINGS = {
|
NODE_CLASS_MAPPINGS = {
|
||||||
|
|||||||
@@ -0,0 +1,85 @@
|
|||||||
|
import folder_paths
|
||||||
|
from PIL import Image, ImageOps
|
||||||
|
import numpy as np
|
||||||
|
import torch
|
||||||
|
import json
|
||||||
|
import comfy
|
||||||
|
|
||||||
|
class ComfyUIDeployExternalImageBatch:
|
||||||
|
@classmethod
|
||||||
|
def INPUT_TYPES(s):
|
||||||
|
return {
|
||||||
|
"required": {
|
||||||
|
"input_id": (
|
||||||
|
"STRING",
|
||||||
|
{"multiline": False, "default": "input_images"},
|
||||||
|
),
|
||||||
|
"images": (
|
||||||
|
"STRING",
|
||||||
|
{"multiline": False, "default": "[]"},
|
||||||
|
),
|
||||||
|
},
|
||||||
|
"optional": {
|
||||||
|
"default_value": ("IMAGE",),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
RETURN_TYPES = ("IMAGE",)
|
||||||
|
RETURN_NAMES = ("image",)
|
||||||
|
|
||||||
|
FUNCTION = "run"
|
||||||
|
|
||||||
|
CATEGORY = "image"
|
||||||
|
|
||||||
|
def run(self, input_id, images=None, default_value=None):
|
||||||
|
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
|
||||||
|
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,'):
|
||||||
|
import base64
|
||||||
|
from io import BytesIO
|
||||||
|
print("Decoding base64 image")
|
||||||
|
base64_image = img_input[img_input.find(",")+1:]
|
||||||
|
decoded_image = base64.b64decode(base64_image)
|
||||||
|
image = Image.open(BytesIO(decoded_image))
|
||||||
|
else:
|
||||||
|
raise ValueError("Invalid image url or base64 data provided.")
|
||||||
|
|
||||||
|
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,]
|
||||||
|
processed_images.append(image_tensor)
|
||||||
|
except Exception as e:
|
||||||
|
print(f"Error processing images: {e}")
|
||||||
|
pass
|
||||||
|
|
||||||
|
if default_value is not None and len(images_list) == 0:
|
||||||
|
processed_images.append(default_value) # Assuming default_value is a pre-processed image tensor
|
||||||
|
|
||||||
|
# Resize images if necessary and concatenate from MakeImageBatch in ImpactPack
|
||||||
|
if processed_images:
|
||||||
|
base_shape = processed_images[0].shape[1:] # Get the shape of the first image for comparison
|
||||||
|
batch_tensor = processed_images[0]
|
||||||
|
for i in range(1, len(processed_images)):
|
||||||
|
if processed_images[i].shape[1:] != base_shape:
|
||||||
|
# Resize to match the first image's dimensions
|
||||||
|
processed_images[i] = comfy.utils.common_upscale(processed_images[i].movedim(-1, 1), base_shape[1], base_shape[0], "lanczos", "center").movedim(1, -1)
|
||||||
|
|
||||||
|
batch_tensor = torch.cat((batch_tensor, processed_images[i]), dim=0)
|
||||||
|
# Concatenate using torch.cat
|
||||||
|
else:
|
||||||
|
batch_tensor = None # or handle the empty case as needed
|
||||||
|
return (batch_tensor, )
|
||||||
|
|
||||||
|
|
||||||
|
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalImageBatch": ComfyUIDeployExternalImageBatch}
|
||||||
|
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyUIDeployExternalImageBatch": "External Image Batch (ComfyUI Deploy)"}
|
||||||
@@ -32,7 +32,12 @@ class ComfyUIDeployExternalLora:
|
|||||||
import os
|
import os
|
||||||
import uuid
|
import uuid
|
||||||
|
|
||||||
if input_id and input_id.startswith('http'):
|
print('external lora using')
|
||||||
|
print("input id: ", input_id)
|
||||||
|
print("default lora : ", default_lora_name)
|
||||||
|
|
||||||
|
if input_id:
|
||||||
|
if input_id.startswith('http'):
|
||||||
unique_filename = str(uuid.uuid4()) + ".safetensors"
|
unique_filename = str(uuid.uuid4()) + ".safetensors"
|
||||||
print(unique_filename)
|
print(unique_filename)
|
||||||
print(folder_paths.folder_names_and_paths["loras"][0][0])
|
print(folder_paths.folder_names_and_paths["loras"][0][0])
|
||||||
@@ -44,6 +49,8 @@ class ComfyUIDeployExternalLora:
|
|||||||
out_file.write(response.content)
|
out_file.write(response.content)
|
||||||
return (unique_filename,)
|
return (unique_filename,)
|
||||||
else:
|
else:
|
||||||
|
return (input_id,)
|
||||||
|
|
||||||
return (default_lora_name,)
|
return (default_lora_name,)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ class ComfyUIDeployExternalNumber:
|
|||||||
"optional": {
|
"optional": {
|
||||||
"default_value": (
|
"default_value": (
|
||||||
"FLOAT",
|
"FLOAT",
|
||||||
{"multiline": True, "display": "number", "default": 0},
|
{"multiline": True, "display": "number", "default": 0, "step": 0.01},
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -29,9 +29,12 @@ class ComfyUIDeployExternalNumber:
|
|||||||
CATEGORY = "number"
|
CATEGORY = "number"
|
||||||
|
|
||||||
def run(self, input_id, default_value=None):
|
def run(self, input_id, default_value=None):
|
||||||
if not input_id or not input_id.strip().isdigit():
|
try:
|
||||||
|
float_value = float(input_id)
|
||||||
|
print("my number", float_value)
|
||||||
|
return [float_value]
|
||||||
|
except ValueError:
|
||||||
return [default_value]
|
return [default_value]
|
||||||
return [int(input_id)]
|
|
||||||
|
|
||||||
|
|
||||||
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalNumber": ComfyUIDeployExternalNumber}
|
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalNumber": ComfyUIDeployExternalNumber}
|
||||||
|
|||||||
@@ -0,0 +1,66 @@
|
|||||||
|
import folder_paths
|
||||||
|
from PIL import Image, ImageOps
|
||||||
|
import numpy as np
|
||||||
|
import torch
|
||||||
|
from server import PromptServer, BinaryEventTypes
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from globals import streaming_prompt_metadata, max_output_id_length
|
||||||
|
|
||||||
|
class ComfyDeployWebscoketImageInput:
|
||||||
|
@classmethod
|
||||||
|
def INPUT_TYPES(s):
|
||||||
|
return {
|
||||||
|
"required": {
|
||||||
|
"input_id": (
|
||||||
|
"STRING",
|
||||||
|
{"multiline": False, "default": "input_id"},
|
||||||
|
),
|
||||||
|
"seed": ("INT", {"default": 0, "min": 0, "max": 0xffffffffffffffff}),
|
||||||
|
},
|
||||||
|
"optional": {
|
||||||
|
"default_value": ("IMAGE", ),
|
||||||
|
"client_id": (
|
||||||
|
"STRING",
|
||||||
|
{"multiline": False, "default": ""},
|
||||||
|
),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
OUTPUT_NODE = True
|
||||||
|
|
||||||
|
RETURN_TYPES = ("IMAGE", )
|
||||||
|
RETURN_NAMES = ("images",)
|
||||||
|
|
||||||
|
FUNCTION = "run"
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def VALIDATE_INPUTS(s, input_id):
|
||||||
|
try:
|
||||||
|
if len(input_id.encode('ascii')) > max_output_id_length:
|
||||||
|
raise ValueError(f"input_id size is greater than {max_output_id_length} bytes")
|
||||||
|
except UnicodeEncodeError:
|
||||||
|
raise ValueError("input_id is not ASCII encodable")
|
||||||
|
|
||||||
|
return True
|
||||||
|
|
||||||
|
def run(self, input_id, seed, default_value=None ,client_id=None):
|
||||||
|
# print(streaming_prompt_metadata[client_id].inputs)
|
||||||
|
if client_id in streaming_prompt_metadata and input_id in streaming_prompt_metadata[client_id].inputs:
|
||||||
|
if isinstance(streaming_prompt_metadata[client_id].inputs[input_id], Image.Image):
|
||||||
|
print("Returning image from websocket input")
|
||||||
|
|
||||||
|
image = streaming_prompt_metadata[client_id].inputs[input_id]
|
||||||
|
|
||||||
|
image = ImageOps.exif_transpose(image)
|
||||||
|
image = image.convert("RGB")
|
||||||
|
image = np.array(image).astype(np.float32) / 255.0
|
||||||
|
image = torch.from_numpy(image)[None,]
|
||||||
|
|
||||||
|
return [image]
|
||||||
|
|
||||||
|
print("Returning default value")
|
||||||
|
return [default_value]
|
||||||
|
|
||||||
|
NODE_CLASS_MAPPINGS = {"ComfyDeployWebscoketImageInput": ComfyDeployWebscoketImageInput}
|
||||||
|
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyDeployWebscoketImageInput": "Image Websocket Input (ComfyDeploy)"}
|
||||||
@@ -0,0 +1,71 @@
|
|||||||
|
import folder_paths
|
||||||
|
from PIL import Image, ImageOps
|
||||||
|
import numpy as np
|
||||||
|
import torch
|
||||||
|
from server import PromptServer, BinaryEventTypes
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from globals import send_image, max_output_id_length
|
||||||
|
|
||||||
|
class ComfyDeployWebscoketImageOutput:
|
||||||
|
@classmethod
|
||||||
|
def INPUT_TYPES(s):
|
||||||
|
return {
|
||||||
|
"required": {
|
||||||
|
"output_id": (
|
||||||
|
"STRING",
|
||||||
|
{"multiline": False, "default": "output_id"},
|
||||||
|
),
|
||||||
|
"images": ("IMAGE", ),
|
||||||
|
"file_type": (["WEBP", "PNG", "JPEG"], ),
|
||||||
|
"quality": ("INT", {"default": 80, "min": 1, "max": 100, "step": 1}),
|
||||||
|
},
|
||||||
|
"optional": {
|
||||||
|
"client_id": (
|
||||||
|
"STRING",
|
||||||
|
{"multiline": False, "default": ""},
|
||||||
|
),
|
||||||
|
}
|
||||||
|
# "hidden": {"client_id": "CLIENT_ID"},
|
||||||
|
}
|
||||||
|
|
||||||
|
OUTPUT_NODE = True
|
||||||
|
|
||||||
|
RETURN_TYPES = ()
|
||||||
|
RETURN_NAMES = ("text",)
|
||||||
|
|
||||||
|
FUNCTION = "run"
|
||||||
|
|
||||||
|
CATEGORY = "output"
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def VALIDATE_INPUTS(s, output_id):
|
||||||
|
try:
|
||||||
|
if len(output_id.encode('ascii')) > max_output_id_length:
|
||||||
|
raise ValueError(f"output_id size is greater than {max_output_id_length} bytes")
|
||||||
|
except UnicodeEncodeError:
|
||||||
|
raise ValueError("output_id is not ASCII encodable")
|
||||||
|
|
||||||
|
return True
|
||||||
|
|
||||||
|
def run(self, output_id, images, file_type, quality, client_id):
|
||||||
|
prompt_server = PromptServer.instance
|
||||||
|
loop = prompt_server.loop
|
||||||
|
|
||||||
|
def schedule_coroutine_blocking(target, *args):
|
||||||
|
future = asyncio.run_coroutine_threadsafe(target(*args), loop)
|
||||||
|
return future.result() # This makes the call blocking
|
||||||
|
|
||||||
|
for tensor in images:
|
||||||
|
array = 255.0 * tensor.cpu().numpy()
|
||||||
|
image = Image.fromarray(np.clip(array, 0, 255).astype(np.uint8))
|
||||||
|
|
||||||
|
schedule_coroutine_blocking(send_image, [file_type, image, None, quality], client_id, output_id)
|
||||||
|
print("Image sent")
|
||||||
|
|
||||||
|
return {"ui": {}}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
NODE_CLASS_MAPPINGS = {"ComfyDeployWebscoketImageOutput": ComfyDeployWebscoketImageOutput}
|
||||||
|
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyDeployWebscoketImageOutput": "Image Websocket Output (ComfyDeploy)"}
|
||||||
+383
-137
@@ -1,38 +1,48 @@
|
|||||||
|
from io import BytesIO
|
||||||
from aiohttp import web
|
from aiohttp import web
|
||||||
import os
|
import os
|
||||||
import requests
|
import requests
|
||||||
import folder_paths
|
import folder_paths
|
||||||
import json
|
import json
|
||||||
import numpy as np
|
|
||||||
import server
|
import server
|
||||||
import re
|
|
||||||
import base64
|
|
||||||
from PIL import Image
|
from PIL import Image
|
||||||
import io
|
|
||||||
import time
|
import time
|
||||||
import execution
|
import execution
|
||||||
import random
|
import random
|
||||||
import traceback
|
import traceback
|
||||||
import uuid
|
import uuid
|
||||||
import asyncio
|
import asyncio
|
||||||
import atexit
|
|
||||||
import logging
|
import logging
|
||||||
import sys
|
|
||||||
from logging.handlers import RotatingFileHandler
|
|
||||||
from enum import Enum
|
|
||||||
from urllib.parse import quote
|
from urllib.parse import quote
|
||||||
import threading
|
import threading
|
||||||
import hashlib
|
import hashlib
|
||||||
import aiohttp
|
import aiohttp
|
||||||
import aiofiles
|
import aiofiles
|
||||||
import concurrent.futures
|
from typing import List, Union, Any, Optional
|
||||||
|
from PIL import Image
|
||||||
|
import copy
|
||||||
|
import struct
|
||||||
|
|
||||||
|
from globals import StreamingPrompt, Status, sockets, SimplePrompt, streaming_prompt_metadata, prompt_metadata
|
||||||
|
|
||||||
api = None
|
api = None
|
||||||
api_task = None
|
api_task = None
|
||||||
prompt_metadata = {}
|
|
||||||
cd_enable_log = os.environ.get('CD_ENABLE_LOG', 'false').lower() == 'true'
|
cd_enable_log = os.environ.get('CD_ENABLE_LOG', 'false').lower() == 'true'
|
||||||
cd_enable_run_log = os.environ.get('CD_ENABLE_RUN_LOG', 'false').lower() == 'true'
|
cd_enable_run_log = os.environ.get('CD_ENABLE_RUN_LOG', 'false').lower() == 'true'
|
||||||
|
|
||||||
|
def clear_current_prompt(sid):
|
||||||
|
prompt_server = server.PromptServer.instance
|
||||||
|
to_delete = list(streaming_prompt_metadata[sid].running_prompt_ids) # Convert set to list
|
||||||
|
|
||||||
|
print("clearning out prompt: ", to_delete)
|
||||||
|
for id_to_delete in to_delete:
|
||||||
|
delete_func = lambda a: a[1] == id_to_delete
|
||||||
|
prompt_server.prompt_queue.delete_queue_item(delete_func)
|
||||||
|
print("deleted prompt: ", id_to_delete, prompt_server.prompt_queue.get_tasks_remaining())
|
||||||
|
|
||||||
|
streaming_prompt_metadata[sid].running_prompt_ids.clear()
|
||||||
|
|
||||||
def post_prompt(json_data):
|
def post_prompt(json_data):
|
||||||
prompt_server = server.PromptServer.instance
|
prompt_server = server.PromptServer.instance
|
||||||
json_data = prompt_server.trigger_on_prompt(json_data)
|
json_data = prompt_server.trigger_on_prompt(json_data)
|
||||||
@@ -80,19 +90,105 @@ def randomSeed(num_digits=15):
|
|||||||
range_end = (10**num_digits) - 1
|
range_end = (10**num_digits) - 1
|
||||||
return random.randint(range_start, range_end)
|
return random.randint(range_start, range_end)
|
||||||
|
|
||||||
|
def apply_random_seed_to_workflow(workflow_api):
|
||||||
|
"""
|
||||||
|
Applies a random seed to each element in the workflow_api that has a 'seed' input.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
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();
|
||||||
|
|
||||||
|
def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
|
||||||
|
# Loop through each of the inputs and replace them
|
||||||
|
for key, value in workflow_api.items():
|
||||||
|
if 'inputs' in value:
|
||||||
|
|
||||||
|
# Support websocket
|
||||||
|
if sid is not None:
|
||||||
|
if (value["class_type"] == "ComfyDeployWebscoketImageOutput"):
|
||||||
|
value['inputs']["client_id"] = sid
|
||||||
|
if (value["class_type"] == "ComfyDeployWebscoketImageInput"):
|
||||||
|
value['inputs']["client_id"] = sid
|
||||||
|
|
||||||
|
if "input_id" in value['inputs'] and value['inputs']['input_id'] in inputs:
|
||||||
|
new_value = inputs[value['inputs']['input_id']]
|
||||||
|
|
||||||
|
# Lets skip it if its an image
|
||||||
|
if isinstance(new_value, Image.Image):
|
||||||
|
continue
|
||||||
|
|
||||||
|
# Backward compactibility
|
||||||
|
value['inputs']["input_id"] = new_value
|
||||||
|
|
||||||
|
# Fix for external text default value
|
||||||
|
if (value["class_type"] == "ComfyUIDeployExternalText"):
|
||||||
|
value['inputs']["default_value"] = new_value
|
||||||
|
|
||||||
|
if (value["class_type"] == "ComfyUIDeployExternalCheckpoint"):
|
||||||
|
value['inputs']["default_value"] = new_value
|
||||||
|
|
||||||
|
if (value["class_type"] == "ComfyUIDeployExternalImageBatch"):
|
||||||
|
value['inputs']["images"] = new_value
|
||||||
|
|
||||||
|
def send_prompt(sid: str, inputs: StreamingPrompt):
|
||||||
|
# workflow_api = inputs.workflow_api
|
||||||
|
workflow_api = copy.deepcopy(inputs.workflow_api)
|
||||||
|
|
||||||
|
# Random seed
|
||||||
|
apply_random_seed_to_workflow(workflow_api)
|
||||||
|
|
||||||
|
print("getting inputs" , inputs.inputs)
|
||||||
|
|
||||||
|
apply_inputs_to_workflow(workflow_api, inputs.inputs, sid=sid)
|
||||||
|
|
||||||
|
print(workflow_api)
|
||||||
|
|
||||||
|
prompt_id = str(uuid.uuid4())
|
||||||
|
|
||||||
|
prompt = {
|
||||||
|
"prompt": workflow_api,
|
||||||
|
"client_id": sid, #"comfy_deploy_instance", #api.client_id
|
||||||
|
"prompt_id": prompt_id
|
||||||
|
}
|
||||||
|
|
||||||
|
try:
|
||||||
|
res = post_prompt(prompt)
|
||||||
|
inputs.running_prompt_ids.add(prompt_id)
|
||||||
|
prompt_metadata[prompt_id] = SimplePrompt(
|
||||||
|
status_endpoint=inputs.status_endpoint,
|
||||||
|
file_upload_endpoint=inputs.file_upload_endpoint,
|
||||||
|
workflow_api=workflow_api,
|
||||||
|
is_realtime=True
|
||||||
|
)
|
||||||
|
except Exception as e:
|
||||||
|
error_type = type(e).__name__
|
||||||
|
stack_trace_short = traceback.format_exc().strip().split('\n')[-2]
|
||||||
|
stack_trace = traceback.format_exc().strip()
|
||||||
|
print(f"error: {error_type}, {e}")
|
||||||
|
print(f"stack trace: {stack_trace_short}")
|
||||||
|
|
||||||
@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):
|
||||||
prompt_server = server.PromptServer.instance
|
prompt_server = server.PromptServer.instance
|
||||||
data = await request.json()
|
data = await request.json()
|
||||||
|
|
||||||
workflow_api = data.get("workflow_api")
|
# 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
|
# 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")
|
||||||
|
|
||||||
for key in workflow_api:
|
# Now it handles directly in here
|
||||||
if 'inputs' in workflow_api[key] and 'seed' in workflow_api[key]['inputs']:
|
apply_random_seed_to_workflow(workflow_api)
|
||||||
workflow_api[key]['inputs']['seed'] = randomSeed()
|
apply_inputs_to_workflow(workflow_api, inputs)
|
||||||
|
|
||||||
prompt = {
|
prompt = {
|
||||||
"prompt": workflow_api,
|
"prompt": workflow_api,
|
||||||
@@ -100,11 +196,11 @@ async def comfy_deploy_run(request):
|
|||||||
"prompt_id": prompt_id
|
"prompt_id": prompt_id
|
||||||
}
|
}
|
||||||
|
|
||||||
prompt_metadata[prompt_id] = {
|
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
|
||||||
}
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
res = post_prompt(prompt)
|
res = post_prompt(prompt)
|
||||||
@@ -148,7 +244,6 @@ async def comfy_deploy_run(request):
|
|||||||
|
|
||||||
return web.json_response(res, status=status)
|
return web.json_response(res, status=status)
|
||||||
|
|
||||||
sockets = dict()
|
|
||||||
|
|
||||||
def get_comfyui_path_from_file_path(file_path):
|
def get_comfyui_path_from_file_path(file_path):
|
||||||
file_path_parts = file_path.split("\\")
|
file_path_parts = file_path.split("\\")
|
||||||
@@ -179,63 +274,22 @@ async def compute_sha256_checksum(filepath):
|
|||||||
sha256.update(chunk)
|
sha256.update(chunk)
|
||||||
return sha256.hexdigest()
|
return sha256.hexdigest()
|
||||||
|
|
||||||
# def hash_chunk(start_end, filepath):
|
@server.PromptServer.instance.routes.get('/comfyui-deploy/models')
|
||||||
# """Hash a specific chunk of the file."""
|
async def get_installed_models(request):
|
||||||
# start, end = start_end
|
# Directly return the list of paths as JSON
|
||||||
# sha256 = hashlib.sha256()
|
new_dict = {}
|
||||||
# with open(filepath, 'rb') as f:
|
for key, value in folder_paths.folder_names_and_paths.items():
|
||||||
# f.seek(start)
|
# Convert set to list for JSON compatibility
|
||||||
# chunk = f.read(end - start)
|
# for path in value[0]:
|
||||||
# sha256.update(chunk)
|
file_list = folder_paths.get_filename_list(key)
|
||||||
# return sha256.digest() # Return the digest of the chunk
|
value_json_compatible = (value[0], list(value[1]), file_list)
|
||||||
|
new_dict[key] = value_json_compatible
|
||||||
# async def compute_sha256_checksum(filepath):
|
# print(new_dict)
|
||||||
# file_size = os.path.getsize(filepath)
|
return web.json_response(new_dict)
|
||||||
# parts = 1 # Or any other division based on file size or desired concurrency
|
|
||||||
# part_size = file_size // parts
|
|
||||||
# start_end_ranges = [(i * part_size, min((i + 1) * part_size, file_size)) for i in range(parts)]
|
|
||||||
|
|
||||||
# print(start_end_ranges, file_size)
|
|
||||||
|
|
||||||
# loop = asyncio.get_running_loop()
|
|
||||||
|
|
||||||
# # Use ThreadPoolExecutor to process chunks in parallel
|
|
||||||
# with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
|
|
||||||
# futures = [loop.run_in_executor(executor, hash_chunk, start_end, filepath) for start_end in start_end_ranges]
|
|
||||||
# chunk_hashes = await asyncio.gather(*futures)
|
|
||||||
|
|
||||||
# # Combine the hashes sequentially
|
|
||||||
# final_sha256 = hashlib.sha256()
|
|
||||||
# for chunk_hash in chunk_hashes:
|
|
||||||
# final_sha256.update(chunk_hash)
|
|
||||||
|
|
||||||
# return final_sha256.hexdigest()
|
|
||||||
|
|
||||||
# def hash_chunk(filepath):
|
|
||||||
# chunk_size = 1024 * 256 # 256KB per chunk
|
|
||||||
# sha256 = hashlib.sha256()
|
|
||||||
# with open(filepath, 'rb') as f:
|
|
||||||
# while True:
|
|
||||||
# chunk = f.read(chunk_size)
|
|
||||||
# if not chunk:
|
|
||||||
# break # End of file
|
|
||||||
# sha256.update(chunk)
|
|
||||||
# return sha256.hexdigest()
|
|
||||||
|
|
||||||
# async def compute_sha256_checksum(filepath):
|
|
||||||
# print("computing sha256 checksum")
|
|
||||||
# filepath = get_comfyui_path_from_file_path(filepath)
|
|
||||||
|
|
||||||
# loop = asyncio.get_running_loop()
|
|
||||||
|
|
||||||
# with concurrent.futures.ThreadPoolExecutor(max_workers=1) as executor:
|
|
||||||
# task = loop.run_in_executor(executor, hash_chunk, filepath)
|
|
||||||
|
|
||||||
# return await task
|
|
||||||
|
|
||||||
# This is start uploading the files to Comfy Deploy
|
# This is start uploading the files to Comfy Deploy
|
||||||
@server.PromptServer.instance.routes.post('/comfyui-deploy/upload-file')
|
@server.PromptServer.instance.routes.post('/comfyui-deploy/upload-file')
|
||||||
async def upload_file(request):
|
async def upload_file_endpoint(request):
|
||||||
data = await request.json()
|
data = await request.json()
|
||||||
|
|
||||||
file_path = data.get("file_path")
|
file_path = data.get("file_path")
|
||||||
@@ -317,26 +371,58 @@ async def upload_file(request):
|
|||||||
}, status=500)
|
}, status=500)
|
||||||
|
|
||||||
|
|
||||||
|
script_dir = os.path.dirname(os.path.abspath(__file__))
|
||||||
|
# Assuming the cache file is stored in the same directory as this script
|
||||||
|
CACHE_FILE_PATH = script_dir + '/file-hash-cache.json'
|
||||||
|
|
||||||
|
# Global in-memory cache
|
||||||
|
file_hash_cache = {}
|
||||||
|
|
||||||
|
# Load cache from disk at startup
|
||||||
|
def load_cache():
|
||||||
|
global file_hash_cache
|
||||||
|
try:
|
||||||
|
with open(CACHE_FILE_PATH, 'r') as cache_file:
|
||||||
|
file_hash_cache = json.load(cache_file)
|
||||||
|
except (FileNotFoundError, json.JSONDecodeError):
|
||||||
|
file_hash_cache = {}
|
||||||
|
|
||||||
|
# Save cache to disk
|
||||||
|
def save_cache():
|
||||||
|
with open(CACHE_FILE_PATH, 'w') as cache_file:
|
||||||
|
json.dump(file_hash_cache, cache_file)
|
||||||
|
|
||||||
|
# Initialize cache on application start
|
||||||
|
load_cache()
|
||||||
|
|
||||||
@server.PromptServer.instance.routes.get('/comfyui-deploy/get-file-hash')
|
@server.PromptServer.instance.routes.get('/comfyui-deploy/get-file-hash')
|
||||||
async def get_file_hash(request):
|
async def get_file_hash(request):
|
||||||
file_path = request.rel_url.query.get('file_path', '')
|
file_path = request.rel_url.query.get('file_path', '')
|
||||||
|
|
||||||
if file_path is None:
|
if not file_path:
|
||||||
return web.json_response({
|
return web.json_response({
|
||||||
"error": "file_path is required"
|
"error": "file_path is required"
|
||||||
}, status=400)
|
}, status=400)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
base = folder_paths.base_path
|
base = folder_paths.base_path
|
||||||
file_path = os.path.join(base, file_path)
|
full_file_path = os.path.join(base, file_path)
|
||||||
# print("file_path", file_path)
|
|
||||||
start_time = time.time() # Capture the start time
|
# Check if the file hash is in the cache
|
||||||
file_hash = await compute_sha256_checksum(
|
if full_file_path in file_hash_cache:
|
||||||
file_path
|
file_hash = file_hash_cache[full_file_path]
|
||||||
)
|
else:
|
||||||
end_time = time.time() # Capture the end time after the code execution
|
start_time = time.time()
|
||||||
elapsed_time = end_time - start_time # Calculate the elapsed time
|
file_hash = await compute_sha256_checksum(full_file_path)
|
||||||
print(f"Execution time: {elapsed_time} seconds")
|
end_time = time.time()
|
||||||
|
elapsed_time = end_time - start_time
|
||||||
|
print(f"Cache miss -> Execution time: {elapsed_time} seconds")
|
||||||
|
|
||||||
|
# Update the in-memory cache
|
||||||
|
file_hash_cache[full_file_path] = file_hash
|
||||||
|
|
||||||
|
save_cache()
|
||||||
|
|
||||||
return web.json_response({
|
return web.json_response({
|
||||||
"file_hash": file_hash
|
"file_hash": file_hash
|
||||||
})
|
})
|
||||||
@@ -345,6 +431,16 @@ async def get_file_hash(request):
|
|||||||
"error": str(e)
|
"error": str(e)
|
||||||
}, status=500)
|
}, status=500)
|
||||||
|
|
||||||
|
async def update_realtime_run_status(realtime_id: str, status_endpoint: str, status: Status):
|
||||||
|
body = {
|
||||||
|
"run_id": realtime_id,
|
||||||
|
"status": status.value,
|
||||||
|
}
|
||||||
|
# requests.post(status_endpoint, json=body)
|
||||||
|
async with aiohttp.ClientSession() as session:
|
||||||
|
async with session.post(status_endpoint, json=body) as response:
|
||||||
|
pass
|
||||||
|
|
||||||
@server.PromptServer.instance.routes.get('/comfyui-deploy/ws')
|
@server.PromptServer.instance.routes.get('/comfyui-deploy/ws')
|
||||||
async def websocket_handler(request):
|
async def websocket_handler(request):
|
||||||
ws = web.WebSocketResponse()
|
ws = web.WebSocketResponse()
|
||||||
@@ -358,35 +454,134 @@ async def websocket_handler(request):
|
|||||||
|
|
||||||
sockets[sid] = ws
|
sockets[sid] = ws
|
||||||
|
|
||||||
|
auth_token = request.rel_url.query.get('token', None)
|
||||||
|
get_workflow_endpoint_url = request.rel_url.query.get('workflow_endpoint', None)
|
||||||
|
realtime_id = request.rel_url.query.get('realtime_id', None)
|
||||||
|
status_endpoint = request.rel_url.query.get('status_endpoint', None)
|
||||||
|
|
||||||
|
if auth_token is not None and get_workflow_endpoint_url is not None:
|
||||||
|
async with aiohttp.ClientSession() as session:
|
||||||
|
headers = {'Authorization': f'Bearer {auth_token}'}
|
||||||
|
async with session.get(get_workflow_endpoint_url, headers=headers) as response:
|
||||||
|
if response.status == 200:
|
||||||
|
workflow = await response.json()
|
||||||
|
|
||||||
|
print("Loaded workflow version ",workflow["version"])
|
||||||
|
|
||||||
|
streaming_prompt_metadata[sid] = StreamingPrompt(
|
||||||
|
workflow_api=workflow["workflow_api"],
|
||||||
|
auth_token=auth_token,
|
||||||
|
inputs={},
|
||||||
|
status_endpoint=status_endpoint,
|
||||||
|
file_upload_endpoint=request.rel_url.query.get('file_upload_endpoint', None),
|
||||||
|
)
|
||||||
|
|
||||||
|
await update_realtime_run_status(realtime_id, status_endpoint, Status.RUNNING)
|
||||||
|
# await send("workflow_api", workflow_api, sid)
|
||||||
|
else:
|
||||||
|
error_message = await response.text()
|
||||||
|
print(f"Failed to fetch workflow endpoint. Status: {response.status}, Error: {error_message}")
|
||||||
|
# await send("error", {"message": error_message}, sid)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Send initial state to the new client
|
# Send initial state to the new client
|
||||||
await send("status", { 'sid': sid }, sid)
|
await send("status", { 'sid': sid }, sid)
|
||||||
|
|
||||||
if cd_enable_log:
|
# Make sure when its connected via client, the full log is not being sent
|
||||||
|
if cd_enable_log and get_workflow_endpoint_url is None:
|
||||||
await send_first_time_log(sid)
|
await send_first_time_log(sid)
|
||||||
|
|
||||||
async for msg in ws:
|
async for msg in ws:
|
||||||
|
if msg.type == aiohttp.WSMsgType.TEXT:
|
||||||
|
try:
|
||||||
|
data = json.loads(msg.data)
|
||||||
|
print(data)
|
||||||
|
event_type = data.get('event')
|
||||||
|
if event_type == 'input':
|
||||||
|
print("Got input: ", data.get("inputs"))
|
||||||
|
input = data.get('inputs')
|
||||||
|
streaming_prompt_metadata[sid].inputs.update(input)
|
||||||
|
elif event_type == 'queue_prompt':
|
||||||
|
clear_current_prompt(sid)
|
||||||
|
send_prompt(sid, streaming_prompt_metadata[sid])
|
||||||
|
else:
|
||||||
|
# Handle other event types
|
||||||
|
pass
|
||||||
|
except json.JSONDecodeError:
|
||||||
|
print('Failed to decode JSON from message')
|
||||||
|
|
||||||
|
if msg.type == aiohttp.WSMsgType.BINARY:
|
||||||
|
data = msg.data
|
||||||
|
event_type, = struct.unpack("<I", data[:4])
|
||||||
|
if event_type == 0: # Image input
|
||||||
|
image_type_code, = struct.unpack("<I", data[4:8])
|
||||||
|
input_id_bytes = data[8:32] # Extract the next 24 bytes for the input ID
|
||||||
|
input_id = input_id_bytes.decode('ascii').strip() # Decode the input ID from ASCII
|
||||||
|
print(event_type)
|
||||||
|
print(image_type_code)
|
||||||
|
print(input_id)
|
||||||
|
image_data = data[32:] # The rest is the image data
|
||||||
|
if image_type_code == 1:
|
||||||
|
image_type = "JPEG"
|
||||||
|
elif image_type_code == 2:
|
||||||
|
image_type = "PNG"
|
||||||
|
elif image_type_code == 3:
|
||||||
|
image_type = "WEBP"
|
||||||
|
else:
|
||||||
|
print("Unknown image type code:", image_type_code)
|
||||||
|
return
|
||||||
|
image = Image.open(BytesIO(image_data))
|
||||||
|
# Check if the input ID already exists and replace the input with the new one
|
||||||
|
if input_id in streaming_prompt_metadata[sid].inputs:
|
||||||
|
# If the input exists, we assume it's an image and attempt to close it to free resources
|
||||||
|
try:
|
||||||
|
existing_image = streaming_prompt_metadata[sid].inputs[input_id]
|
||||||
|
if hasattr(existing_image, 'close'):
|
||||||
|
existing_image.close()
|
||||||
|
except Exception as e:
|
||||||
|
print(f"Error closing previous image for input ID {input_id}: {e}")
|
||||||
|
streaming_prompt_metadata[sid].inputs[input_id] = image
|
||||||
|
# clear_current_prompt(sid)
|
||||||
|
# send_prompt(sid, streaming_prompt_metadata[sid])
|
||||||
|
print(f"Received {image_type} image of size {image.size} with input ID {input_id}")
|
||||||
|
|
||||||
if msg.type == aiohttp.WSMsgType.ERROR:
|
if msg.type == aiohttp.WSMsgType.ERROR:
|
||||||
print('ws connection closed with exception %s' % ws.exception())
|
print('ws connection closed with exception %s' % ws.exception())
|
||||||
finally:
|
finally:
|
||||||
sockets.pop(sid, None)
|
sockets.pop(sid, None)
|
||||||
|
|
||||||
|
if realtime_id is not None:
|
||||||
|
await update_realtime_run_status(realtime_id, status_endpoint, Status.SUCCESS)
|
||||||
return ws
|
return ws
|
||||||
|
|
||||||
@server.PromptServer.instance.routes.get('/comfyui-deploy/check-status')
|
@server.PromptServer.instance.routes.get('/comfyui-deploy/check-status')
|
||||||
async def comfy_deploy_check_status(request):
|
async def comfy_deploy_check_status(request):
|
||||||
prompt_server = server.PromptServer.instance
|
|
||||||
prompt_id = request.rel_url.query.get('prompt_id', None)
|
prompt_id = request.rel_url.query.get('prompt_id', None)
|
||||||
if prompt_id in prompt_metadata and 'status' in prompt_metadata[prompt_id]:
|
if prompt_id in prompt_metadata:
|
||||||
return web.json_response({
|
return web.json_response({
|
||||||
"status": prompt_metadata[prompt_id]['status'].value
|
"status": prompt_metadata[prompt_id].status.value
|
||||||
})
|
})
|
||||||
else:
|
else:
|
||||||
return web.json_response({
|
return web.json_response({
|
||||||
"message": "prompt_id not found"
|
"message": "prompt_id not found"
|
||||||
})
|
})
|
||||||
|
|
||||||
|
@server.PromptServer.instance.routes.get('/comfyui-deploy/check-ws-status')
|
||||||
|
async def comfy_deploy_check_ws_status(request):
|
||||||
|
client_id = request.rel_url.query.get('client_id', None)
|
||||||
|
if client_id in streaming_prompt_metadata:
|
||||||
|
remaining_queue = 0 # Initialize remaining queue count
|
||||||
|
for prompt_id in streaming_prompt_metadata[client_id].running_prompt_ids:
|
||||||
|
prompt_status = prompt_metadata[prompt_id].status
|
||||||
|
if prompt_status not in [Status.FAILED, Status.SUCCESS]:
|
||||||
|
remaining_queue += 1 # Increment for each prompt still running
|
||||||
|
return web.json_response({"remaining_queue": remaining_queue})
|
||||||
|
else:
|
||||||
|
return web.json_response({"message": "client_id not found"}, status=404)
|
||||||
|
|
||||||
async def send(event, data, sid=None):
|
async def send(event, data, sid=None):
|
||||||
try:
|
try:
|
||||||
|
# message = {"event": event, "data": data}
|
||||||
if sid:
|
if sid:
|
||||||
ws = sockets.get(sid)
|
ws = sockets.get(sid)
|
||||||
if ws != None and not ws.closed: # Check if the WebSocket connection is open and not closing
|
if ws != None and not ws.closed: # Check if the WebSocket connection is open and not closing
|
||||||
@@ -402,53 +597,74 @@ async def send(event, data, sid=None):
|
|||||||
logging.basicConfig(level=logging.INFO)
|
logging.basicConfig(level=logging.INFO)
|
||||||
|
|
||||||
prompt_server = server.PromptServer.instance
|
prompt_server = server.PromptServer.instance
|
||||||
|
|
||||||
send_json = prompt_server.send_json
|
send_json = prompt_server.send_json
|
||||||
|
|
||||||
async def send_json_override(self, event, data, sid=None):
|
async def send_json_override(self, event, data, sid=None):
|
||||||
# print("INTERNAL:", event, data, sid)
|
# print("INTERNAL:", event, data, sid)
|
||||||
prompt_id = data.get('prompt_id')
|
prompt_id = data.get('prompt_id')
|
||||||
|
|
||||||
|
target_sid = sid
|
||||||
|
if target_sid == "comfy_deploy_instance":
|
||||||
|
target_sid = None
|
||||||
|
|
||||||
# now we send everything
|
# now we send everything
|
||||||
await asyncio.wait([
|
await asyncio.wait([
|
||||||
asyncio.create_task(send(event, data)),
|
asyncio.create_task(send(event, data, sid=target_sid)),
|
||||||
asyncio.create_task(self.send_json_original(event, data, sid))
|
asyncio.create_task(self.send_json_original(event, data, sid))
|
||||||
])
|
])
|
||||||
|
|
||||||
if event == 'execution_start':
|
if event == 'execution_start':
|
||||||
update_run(prompt_id, Status.RUNNING)
|
update_run(prompt_id, Status.RUNNING)
|
||||||
|
|
||||||
|
if prompt_id in prompt_metadata:
|
||||||
|
prompt_metadata[prompt_id].start_time = time.perf_counter()
|
||||||
|
|
||||||
# the last executing event is none, then the workflow is finished
|
# the last executing event is none, then the workflow is finished
|
||||||
if event == 'executing' and data.get('node') is None:
|
if event == 'executing' and data.get('node') is None:
|
||||||
mark_prompt_done(prompt_id=prompt_id)
|
mark_prompt_done(prompt_id=prompt_id)
|
||||||
if not have_pending_upload(prompt_id):
|
if not have_pending_upload(prompt_id):
|
||||||
update_run(prompt_id, Status.SUCCESS)
|
update_run(prompt_id, Status.SUCCESS)
|
||||||
|
if prompt_id in prompt_metadata:
|
||||||
|
current_time = time.perf_counter()
|
||||||
|
if prompt_metadata[prompt_id].start_time is not None:
|
||||||
|
elapsed_time = current_time - prompt_metadata[prompt_id].start_time
|
||||||
|
print(f"Elapsed time: {elapsed_time} seconds")
|
||||||
|
await send("elapsed_time", {
|
||||||
|
"prompt_id": prompt_id,
|
||||||
|
"elapsed_time": elapsed_time
|
||||||
|
}, sid=sid)
|
||||||
|
|
||||||
if event == 'executing' and data.get('node') is not None:
|
if event == 'executing' and data.get('node') is not None:
|
||||||
node = data.get('node')
|
node = data.get('node')
|
||||||
|
|
||||||
if prompt_id in prompt_metadata:
|
if prompt_id in prompt_metadata:
|
||||||
if 'progress' not in prompt_metadata[prompt_id]:
|
# if 'progress' not in prompt_metadata[prompt_id]:
|
||||||
prompt_metadata[prompt_id]["progress"] = set()
|
# prompt_metadata[prompt_id]["progress"] = set()
|
||||||
|
|
||||||
prompt_metadata[prompt_id]["progress"].add(node)
|
prompt_metadata[prompt_id].progress.add(node)
|
||||||
calculated_progress = len(prompt_metadata[prompt_id]["progress"]) / len(prompt_metadata[prompt_id]['workflow_api'])
|
calculated_progress = len(prompt_metadata[prompt_id].progress) / len(prompt_metadata[prompt_id].workflow_api)
|
||||||
# print("calculated_progress", calculated_progress)
|
# print("calculated_progress", calculated_progress)
|
||||||
|
|
||||||
if 'last_updated_node' in prompt_metadata[prompt_id] and prompt_metadata[prompt_id]['last_updated_node'] == node:
|
if prompt_metadata[prompt_id].last_updated_node is not None and prompt_metadata[prompt_id].last_updated_node == node:
|
||||||
return
|
return
|
||||||
prompt_metadata[prompt_id]['last_updated_node'] = node
|
prompt_metadata[prompt_id].last_updated_node = node
|
||||||
class_type = prompt_metadata[prompt_id]['workflow_api'][node]['class_type']
|
class_type = prompt_metadata[prompt_id].workflow_api[node]['class_type']
|
||||||
print("updating run live status", class_type)
|
print("updating run live status", class_type)
|
||||||
|
await send("live_status", {
|
||||||
|
"prompt_id": prompt_id,
|
||||||
|
"current_node": class_type,
|
||||||
|
"progress": calculated_progress,
|
||||||
|
}, sid=sid)
|
||||||
await update_run_live_status(prompt_id, "Executing " + class_type, calculated_progress)
|
await update_run_live_status(prompt_id, "Executing " + class_type, calculated_progress)
|
||||||
|
|
||||||
if event == 'execution_cached' and data.get('nodes') is not None:
|
if event == 'execution_cached' and data.get('nodes') is not None:
|
||||||
if prompt_id in prompt_metadata:
|
if prompt_id in prompt_metadata:
|
||||||
if 'progress' not in prompt_metadata[prompt_id]:
|
# if 'progress' not in prompt_metadata[prompt_id]:
|
||||||
prompt_metadata[prompt_id]["progress"] = set()
|
# prompt_metadata[prompt_id].progress = set()
|
||||||
|
|
||||||
if 'nodes' in data:
|
if 'nodes' in data:
|
||||||
for node in data.get('nodes', []):
|
for node in data.get('nodes', []):
|
||||||
prompt_metadata[prompt_id]["progress"].add(node)
|
prompt_metadata[prompt_id].progress.add(node)
|
||||||
# prompt_metadata[prompt_id]["progress"].update(data.get('nodes'))
|
# prompt_metadata[prompt_id]["progress"].update(data.get('nodes'))
|
||||||
|
|
||||||
if event == 'execution_error':
|
if event == 'execution_error':
|
||||||
@@ -458,18 +674,19 @@ async def send_json_override(self, event, data, sid=None):
|
|||||||
# await update_run_with_output(prompt_id, data)
|
# await update_run_with_output(prompt_id, data)
|
||||||
|
|
||||||
if event == 'executed' and 'node' in data and 'output' in data:
|
if event == 'executed' and 'node' in data and 'output' in data:
|
||||||
|
print("executed", data)
|
||||||
|
if prompt_id in prompt_metadata:
|
||||||
|
node = data.get('node')
|
||||||
|
class_type = prompt_metadata[prompt_id].workflow_api[node]['class_type']
|
||||||
|
print("executed", class_type)
|
||||||
|
if class_type == "PreviewImage":
|
||||||
|
print("skipping preview image")
|
||||||
|
return
|
||||||
|
|
||||||
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'))
|
||||||
# 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'))
|
# update_run_with_output(prompt_id, data.get('output'))
|
||||||
|
|
||||||
|
|
||||||
class Status(Enum):
|
|
||||||
NOT_STARTED = "not-started"
|
|
||||||
RUNNING = "running"
|
|
||||||
SUCCESS = "success"
|
|
||||||
FAILED = "failed"
|
|
||||||
UPLOADING = "uploading"
|
|
||||||
|
|
||||||
# Global variable to keep track of the last read line number
|
# Global variable to keep track of the last read line number
|
||||||
last_read_line_number = 0
|
last_read_line_number = 0
|
||||||
|
|
||||||
@@ -477,9 +694,12 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
|
|||||||
if prompt_id not in prompt_metadata:
|
if prompt_id not in prompt_metadata:
|
||||||
return
|
return
|
||||||
|
|
||||||
|
if prompt_metadata[prompt_id].is_realtime is True:
|
||||||
|
return
|
||||||
|
|
||||||
print("progress", calculated_progress)
|
print("progress", calculated_progress)
|
||||||
|
|
||||||
status_endpoint = prompt_metadata[prompt_id]['status_endpoint']
|
status_endpoint = prompt_metadata[prompt_id].status_endpoint
|
||||||
body = {
|
body = {
|
||||||
"run_id": prompt_id,
|
"run_id": prompt_id,
|
||||||
"live_status": live_status,
|
"live_status": live_status,
|
||||||
@@ -491,19 +711,25 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
def update_run(prompt_id, status: Status):
|
def update_run(prompt_id: str, status: Status):
|
||||||
global last_read_line_number
|
global last_read_line_number
|
||||||
|
|
||||||
if prompt_id not in prompt_metadata:
|
if prompt_id not in prompt_metadata:
|
||||||
return
|
return
|
||||||
|
|
||||||
if ('status' not in prompt_metadata[prompt_id] or prompt_metadata[prompt_id]['status'] != status):
|
# if prompt_metadata[prompt_id].start_time is None and status == Status.RUNNING:
|
||||||
|
# if its realtime prompt we need to skip that.
|
||||||
# when the status is already failed, we don't want to update it to success
|
if prompt_metadata[prompt_id].is_realtime is True:
|
||||||
if ('status' in prompt_metadata[prompt_id] and prompt_metadata[prompt_id]['status'] == Status.FAILED):
|
prompt_metadata[prompt_id].status = status
|
||||||
return
|
return
|
||||||
|
|
||||||
status_endpoint = prompt_metadata[prompt_id]['status_endpoint']
|
if (prompt_metadata[prompt_id].status != status):
|
||||||
|
|
||||||
|
# when the status is already failed, we don't want to update it to success
|
||||||
|
if (prompt_metadata[prompt_id].status is Status.FAILED):
|
||||||
|
return
|
||||||
|
|
||||||
|
status_endpoint = prompt_metadata[prompt_id].status_endpoint
|
||||||
body = {
|
body = {
|
||||||
"run_id": prompt_id,
|
"run_id": prompt_id,
|
||||||
"status": status.value,
|
"status": status.value,
|
||||||
@@ -544,13 +770,12 @@ def update_run(prompt_id, status: Status):
|
|||||||
except Exception as log_error:
|
except Exception as log_error:
|
||||||
print(f"Error reading log file: {log_error}")
|
print(f"Error reading log file: {log_error}")
|
||||||
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_type = type(e).__name__
|
error_type = type(e).__name__
|
||||||
stack_trace = traceback.format_exc().strip()
|
stack_trace = traceback.format_exc().strip()
|
||||||
print(f"Error occurred while updating run: {e} {stack_trace}")
|
print(f"Error occurred while updating run: {e} {stack_trace}")
|
||||||
finally:
|
finally:
|
||||||
prompt_metadata[prompt_id]['status'] = status
|
prompt_metadata[prompt_id].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"):
|
||||||
@@ -582,7 +807,7 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
|
|||||||
|
|
||||||
print("uploading file", file)
|
print("uploading file", file)
|
||||||
|
|
||||||
file_upload_endpoint = prompt_metadata[prompt_id]['file_upload_endpoint']
|
file_upload_endpoint = prompt_metadata[prompt_id].file_upload_endpoint
|
||||||
|
|
||||||
filename = quote(filename)
|
filename = quote(filename)
|
||||||
prompt_id = quote(prompt_id)
|
prompt_id = quote(prompt_id)
|
||||||
@@ -606,21 +831,35 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
|
|||||||
print("upload file response", response.status)
|
print("upload file response", response.status)
|
||||||
|
|
||||||
def have_pending_upload(prompt_id):
|
def have_pending_upload(prompt_id):
|
||||||
if 'prompt_id' in prompt_metadata and 'uploading_nodes' in prompt_metadata[prompt_id] and len(prompt_metadata[prompt_id]['uploading_nodes']) > 0:
|
if prompt_id in prompt_metadata and len(prompt_metadata[prompt_id].uploading_nodes) > 0:
|
||||||
print("have pending upload ", len(prompt_metadata[prompt_id]['uploading_nodes']))
|
print("have pending upload ", len(prompt_metadata[prompt_id].uploading_nodes))
|
||||||
return True
|
return True
|
||||||
|
|
||||||
print("no pending upload")
|
print("no pending upload")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
def mark_prompt_done(prompt_id):
|
def mark_prompt_done(prompt_id):
|
||||||
|
"""
|
||||||
|
Mark the prompt as done in the prompt metadata.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
prompt_id (str): The ID of the prompt to mark as done.
|
||||||
|
"""
|
||||||
if prompt_id in prompt_metadata:
|
if prompt_id in prompt_metadata:
|
||||||
prompt_metadata[prompt_id]["done"] = True
|
prompt_metadata[prompt_id].done = True
|
||||||
print("Prompt done")
|
print("Prompt done")
|
||||||
|
|
||||||
def is_prompt_done(prompt_id):
|
def is_prompt_done(prompt_id: str):
|
||||||
if prompt_id in prompt_metadata and "done" in prompt_metadata[prompt_id]:
|
"""
|
||||||
if prompt_metadata[prompt_id]["done"] == True:
|
Check if the prompt with the given ID is marked as done.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
prompt_id (str): The ID of the prompt to check.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
bool: True if the prompt is marked as done, False otherwise.
|
||||||
|
"""
|
||||||
|
if prompt_id in prompt_metadata and prompt_metadata[prompt_id].done is True:
|
||||||
return True
|
return True
|
||||||
|
|
||||||
return False
|
return False
|
||||||
@@ -644,17 +883,17 @@ async def handle_error(prompt_id, data, e: Exception):
|
|||||||
print(f"Error occurred while uploading file: {e}")
|
print(f"Error occurred while uploading file: {e}")
|
||||||
|
|
||||||
# Mark the current prompt requires upload, and block it from being marked as success
|
# Mark the current prompt requires upload, and block it from being marked as success
|
||||||
async def update_file_status(prompt_id, data, uploading, have_error=False, node_id=None):
|
async def update_file_status(prompt_id: str, data, uploading, have_error=False, node_id=None):
|
||||||
if 'uploading_nodes' not in prompt_metadata[prompt_id]:
|
# if 'uploading_nodes' not in prompt_metadata[prompt_id]:
|
||||||
prompt_metadata[prompt_id]['uploading_nodes'] = set()
|
# prompt_metadata[prompt_id]['uploading_nodes'] = set()
|
||||||
|
|
||||||
if node_id is not None:
|
if node_id is not None:
|
||||||
if uploading:
|
if uploading:
|
||||||
prompt_metadata[prompt_id]['uploading_nodes'].add(node_id)
|
prompt_metadata[prompt_id].uploading_nodes.add(node_id)
|
||||||
else:
|
else:
|
||||||
prompt_metadata[prompt_id]['uploading_nodes'].discard(node_id)
|
prompt_metadata[prompt_id].uploading_nodes.discard(node_id)
|
||||||
|
|
||||||
print(prompt_metadata[prompt_id]['uploading_nodes'])
|
print(prompt_metadata[prompt_id].uploading_nodes)
|
||||||
# Update the remote status
|
# Update the remote status
|
||||||
|
|
||||||
if have_error:
|
if have_error:
|
||||||
@@ -666,7 +905,7 @@ async def update_file_status(prompt_id, data, uploading, have_error=False, node_
|
|||||||
|
|
||||||
# if there are still nodes that are uploading, then we set the status to uploading
|
# if there are still nodes that are uploading, then we set the status to uploading
|
||||||
if uploading:
|
if uploading:
|
||||||
if prompt_metadata[prompt_id]['status'] != Status.UPLOADING:
|
if prompt_metadata[prompt_id].status != Status.UPLOADING:
|
||||||
update_run(prompt_id, Status.UPLOADING)
|
update_run(prompt_id, Status.UPLOADING)
|
||||||
await send("uploading", {
|
await send("uploading", {
|
||||||
"prompt_id": prompt_id,
|
"prompt_id": prompt_id,
|
||||||
@@ -680,9 +919,12 @@ async def update_file_status(prompt_id, data, uploading, have_error=False, node_
|
|||||||
"prompt_id": prompt_id,
|
"prompt_id": prompt_id,
|
||||||
})
|
})
|
||||||
|
|
||||||
async def handle_upload(prompt_id, data, key, content_type_key, default_content_type):
|
async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, default_content_type: str):
|
||||||
items = data.get(key, [])
|
items = data.get(key, [])
|
||||||
for item in items:
|
for item in items:
|
||||||
|
# # Skipping temp files
|
||||||
|
if item.get("type") == "temp":
|
||||||
|
continue
|
||||||
await upload_file(
|
await upload_file(
|
||||||
prompt_id,
|
prompt_id,
|
||||||
item.get("filename"),
|
item.get("filename"),
|
||||||
@@ -691,9 +933,8 @@ async def handle_upload(prompt_id, data, key, content_type_key, default_content_
|
|||||||
content_type=item.get(content_type_key, default_content_type)
|
content_type=item.get(content_type_key, default_content_type)
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
# Upload files in the background
|
# Upload files in the background
|
||||||
async def upload_in_background(prompt_id, data, node_id=None, have_upload=True):
|
async def upload_in_background(prompt_id: str, data, node_id=None, have_upload=True):
|
||||||
try:
|
try:
|
||||||
await handle_upload(prompt_id, data, 'images', "content_type", "image/png")
|
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, 'files', "content_type", "image/png")
|
||||||
@@ -706,8 +947,13 @@ async def upload_in_background(prompt_id, data, node_id=None, have_upload=True):
|
|||||||
await handle_error(prompt_id, data, 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):
|
||||||
if prompt_id in prompt_metadata:
|
if prompt_id not in prompt_metadata:
|
||||||
status_endpoint = prompt_metadata[prompt_id]['status_endpoint']
|
return
|
||||||
|
|
||||||
|
if prompt_metadata[prompt_id].is_realtime is True:
|
||||||
|
return
|
||||||
|
|
||||||
|
status_endpoint = prompt_metadata[prompt_id].status_endpoint
|
||||||
|
|
||||||
body = {
|
body = {
|
||||||
"run_id": prompt_id,
|
"run_id": prompt_id,
|
||||||
|
|||||||
+115
@@ -0,0 +1,115 @@
|
|||||||
|
import struct
|
||||||
|
from enum import Enum
|
||||||
|
import aiohttp
|
||||||
|
from typing import List, Union, Any, Optional
|
||||||
|
from PIL import Image, ImageOps
|
||||||
|
from io import BytesIO
|
||||||
|
from pydantic import BaseModel as PydanticBaseModel
|
||||||
|
|
||||||
|
class BaseModel(PydanticBaseModel):
|
||||||
|
class Config:
|
||||||
|
arbitrary_types_allowed = True
|
||||||
|
|
||||||
|
class Status(Enum):
|
||||||
|
NOT_STARTED = "not-started"
|
||||||
|
RUNNING = "running"
|
||||||
|
SUCCESS = "success"
|
||||||
|
FAILED = "failed"
|
||||||
|
UPLOADING = "uploading"
|
||||||
|
|
||||||
|
class StreamingPrompt(BaseModel):
|
||||||
|
workflow_api: Any
|
||||||
|
auth_token: str
|
||||||
|
inputs: dict[str, Union[str, bytes, Image.Image]]
|
||||||
|
running_prompt_ids: set[str] = set()
|
||||||
|
status_endpoint: str
|
||||||
|
file_upload_endpoint: str
|
||||||
|
|
||||||
|
class SimplePrompt(BaseModel):
|
||||||
|
status_endpoint: str
|
||||||
|
file_upload_endpoint: str
|
||||||
|
workflow_api: dict
|
||||||
|
status: Status = Status.NOT_STARTED
|
||||||
|
progress: set = set()
|
||||||
|
last_updated_node: Optional[str] = None,
|
||||||
|
uploading_nodes: set = set()
|
||||||
|
done: bool = False
|
||||||
|
is_realtime: bool = False,
|
||||||
|
start_time: Optional[float] = None,
|
||||||
|
|
||||||
|
sockets = dict()
|
||||||
|
prompt_metadata: dict[str, SimplePrompt] = {}
|
||||||
|
streaming_prompt_metadata: dict[str, StreamingPrompt] = {}
|
||||||
|
|
||||||
|
class BinaryEventTypes:
|
||||||
|
PREVIEW_IMAGE = 1
|
||||||
|
UNENCODED_PREVIEW_IMAGE = 2
|
||||||
|
|
||||||
|
max_output_id_length = 24
|
||||||
|
|
||||||
|
async def send_image(image_data, sid=None, output_id:str = None):
|
||||||
|
max_length = max_output_id_length
|
||||||
|
output_id = output_id[:max_length]
|
||||||
|
padded_output_id = output_id.ljust(max_length, '\x00')
|
||||||
|
encoded_output_id = padded_output_id.encode('ascii', 'replace')
|
||||||
|
|
||||||
|
image_type = image_data[0]
|
||||||
|
image = image_data[1]
|
||||||
|
max_size = image_data[2]
|
||||||
|
quality = image_data[3]
|
||||||
|
if max_size is not None:
|
||||||
|
if hasattr(Image, 'Resampling'):
|
||||||
|
resampling = Image.Resampling.BILINEAR
|
||||||
|
else:
|
||||||
|
resampling = Image.ANTIALIAS
|
||||||
|
|
||||||
|
image = ImageOps.contain(image, (max_size, max_size), resampling)
|
||||||
|
type_num = 1
|
||||||
|
if image_type == "JPEG":
|
||||||
|
type_num = 1
|
||||||
|
elif image_type == "PNG":
|
||||||
|
type_num = 2
|
||||||
|
elif image_type == "WEBP":
|
||||||
|
type_num = 3
|
||||||
|
|
||||||
|
bytesIO = BytesIO()
|
||||||
|
header = struct.pack(">I", type_num)
|
||||||
|
# 4 bytes for the type
|
||||||
|
bytesIO.write(header)
|
||||||
|
# 10 bytes for the output_id
|
||||||
|
position_before = bytesIO.tell()
|
||||||
|
bytesIO.write(encoded_output_id)
|
||||||
|
position_after = bytesIO.tell()
|
||||||
|
bytes_written = position_after - position_before
|
||||||
|
print(f"Bytes written: {bytes_written}")
|
||||||
|
|
||||||
|
image.save(bytesIO, format=image_type, quality=quality, compress_level=1)
|
||||||
|
preview_bytes = bytesIO.getvalue()
|
||||||
|
await send_bytes(BinaryEventTypes.PREVIEW_IMAGE, preview_bytes, sid=sid)
|
||||||
|
|
||||||
|
async def send_socket_catch_exception(function, message):
|
||||||
|
try:
|
||||||
|
await function(message)
|
||||||
|
except (aiohttp.ClientError, aiohttp.ClientPayloadError, ConnectionResetError) as err:
|
||||||
|
print("send error:", err)
|
||||||
|
|
||||||
|
def encode_bytes(event, data):
|
||||||
|
if not isinstance(event, int):
|
||||||
|
raise RuntimeError(f"Binary event types must be integers, got {event}")
|
||||||
|
|
||||||
|
packed = struct.pack(">I", event)
|
||||||
|
message = bytearray(packed)
|
||||||
|
message.extend(data)
|
||||||
|
return message
|
||||||
|
|
||||||
|
async def send_bytes(event, data, sid=None):
|
||||||
|
message = encode_bytes(event, data)
|
||||||
|
|
||||||
|
print("sending image to ", event, sid)
|
||||||
|
|
||||||
|
if sid is None:
|
||||||
|
_sockets = list(sockets.values())
|
||||||
|
for ws in _sockets:
|
||||||
|
await send_socket_catch_exception(ws.send_bytes, message)
|
||||||
|
elif sid in sockets:
|
||||||
|
await send_socket_catch_exception(sockets[sid].send_bytes, message)
|
||||||
@@ -58,6 +58,9 @@ if cd_enable_log:
|
|||||||
print("** Comfy Deploy logging enabled")
|
print("** Comfy Deploy logging enabled")
|
||||||
setup()
|
setup()
|
||||||
|
|
||||||
|
|
||||||
|
# Store the original working directory
|
||||||
|
original_cwd = os.getcwd()
|
||||||
try:
|
try:
|
||||||
# Get the absolute path of the script's directory
|
# Get the absolute path of the script's directory
|
||||||
script_dir = os.path.dirname(os.path.abspath(__file__))
|
script_dir = os.path.dirname(os.path.abspath(__file__))
|
||||||
@@ -67,3 +70,6 @@ try:
|
|||||||
print(f"** Comfy Deploy Revision: {current_git_commit}")
|
print(f"** Comfy Deploy Revision: {current_git_commit}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"** Comfy Deploy failed to get current git commit: {str(e)}")
|
print(f"** Comfy Deploy failed to get current git commit: {str(e)}")
|
||||||
|
finally:
|
||||||
|
# Change back to the original directory
|
||||||
|
os.chdir(original_cwd)
|
||||||
@@ -1 +1,2 @@
|
|||||||
aiofiles
|
aiofiles
|
||||||
|
pydantic
|
||||||
+126
-52
@@ -1,10 +1,18 @@
|
|||||||
import { app } from "./app.js";
|
import { app } from "./app.js";
|
||||||
import { api } from "./api.js";
|
import { api } from "./api.js";
|
||||||
import { ComfyWidgets, LGraphNode } from "./widgets.js";
|
import { ComfyWidgets, LGraphNode } from "./widgets.js";
|
||||||
import { generateDependencyGraph } from "https://esm.sh/[email protected].19";
|
import { generateDependencyGraph } from "https://esm.sh/[email protected].25";
|
||||||
|
|
||||||
const loadingIcon = `<svg xmlns="http://www.w3.org/2000/svg" width="32" height="32" viewBox="0 0 24 24"><g fill="none" stroke="#888888" stroke-linecap="round" stroke-width="2"><path stroke-dasharray="60" stroke-dashoffset="60" stroke-opacity=".3" d="M12 3C16.9706 3 21 7.02944 21 12C21 16.9706 16.9706 21 12 21C7.02944 21 3 16.9706 3 12C3 7.02944 7.02944 3 12 3Z"><animate fill="freeze" attributeName="stroke-dashoffset" dur="1.3s" values="60;0"/></path><path stroke-dasharray="15" stroke-dashoffset="15" d="M12 3C16.9706 3 21 7.02944 21 12"><animate fill="freeze" attributeName="stroke-dashoffset" dur="0.3s" values="15;0"/><animateTransform attributeName="transform" dur="1.5s" repeatCount="indefinite" type="rotate" values="0 12 12;360 12 12"/></path></g></svg>`;
|
const loadingIcon = `<svg xmlns="http://www.w3.org/2000/svg" width="32" height="32" viewBox="0 0 24 24"><g fill="none" stroke="#888888" stroke-linecap="round" stroke-width="2"><path stroke-dasharray="60" stroke-dashoffset="60" stroke-opacity=".3" d="M12 3C16.9706 3 21 7.02944 21 12C21 16.9706 16.9706 21 12 21C7.02944 21 3 16.9706 3 12C3 7.02944 7.02944 3 12 3Z"><animate fill="freeze" attributeName="stroke-dashoffset" dur="1.3s" values="60;0"/></path><path stroke-dasharray="15" stroke-dashoffset="15" d="M12 3C16.9706 3 21 7.02944 21 12"><animate fill="freeze" attributeName="stroke-dashoffset" dur="0.3s" values="15;0"/><animateTransform attributeName="transform" dur="1.5s" repeatCount="indefinite" type="rotate" values="0 12 12;360 12 12"/></path></g></svg>`;
|
||||||
|
|
||||||
|
function sendEventToCD(event, data) {
|
||||||
|
const message = {
|
||||||
|
type: event,
|
||||||
|
data: data,
|
||||||
|
};
|
||||||
|
window.parent.postMessage(JSON.stringify(message), "*");
|
||||||
|
}
|
||||||
|
|
||||||
/** @typedef {import('../../../web/types/comfy.js').ComfyExtension} ComfyExtension*/
|
/** @typedef {import('../../../web/types/comfy.js').ComfyExtension} ComfyExtension*/
|
||||||
/** @type {ComfyExtension} */
|
/** @type {ComfyExtension} */
|
||||||
const ext = {
|
const ext = {
|
||||||
@@ -18,6 +26,11 @@ const ext = {
|
|||||||
const auth_token = queryParams.get("auth_token");
|
const auth_token = queryParams.get("auth_token");
|
||||||
const org_display = queryParams.get("org_display");
|
const org_display = queryParams.get("org_display");
|
||||||
const origin = queryParams.get("origin");
|
const origin = queryParams.get("origin");
|
||||||
|
const workspace_mode = queryParams.get("workspace_mode");
|
||||||
|
|
||||||
|
if (workspace_mode) {
|
||||||
|
document.querySelector(".comfy-menu").style.display = "none";
|
||||||
|
}
|
||||||
|
|
||||||
const data = getData();
|
const data = getData();
|
||||||
let endpoint = data.endpoint;
|
let endpoint = data.endpoint;
|
||||||
@@ -59,8 +72,8 @@ const ext = {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Adding a delay to wait for the intial graph to load
|
// // Adding a delay to wait for the intial graph to load
|
||||||
await new Promise((resolve) => setTimeout(resolve, 2000));
|
// await new Promise((resolve) => setTimeout(resolve, 2000));
|
||||||
|
|
||||||
workflow?.nodes.forEach((x) => {
|
workflow?.nodes.forEach((x) => {
|
||||||
if (x?.type === "ComfyDeploy") {
|
if (x?.type === "ComfyDeploy") {
|
||||||
@@ -152,9 +165,32 @@ const ext = {
|
|||||||
async setup() {
|
async setup() {
|
||||||
// const graphCanvas = document.getElementById("graph-canvas");
|
// const graphCanvas = document.getElementById("graph-canvas");
|
||||||
|
|
||||||
window.addEventListener("message", (event) => {
|
window.addEventListener("message", async (event) => {
|
||||||
if (!event.data.flow || Object.entries(event.data.flow).length <= 0)
|
try {
|
||||||
return;
|
const message = JSON.parse(event.data);
|
||||||
|
if (message.type === "graph_load") {
|
||||||
|
const comfyUIWorkflow = message.data;
|
||||||
|
console.log("recieved: ", comfyUIWorkflow);
|
||||||
|
// Assuming there's a method to load the workflow data into the ComfyUI
|
||||||
|
// 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
|
||||||
|
if (comfyUIWorkflow && app && app.loadGraphData) {
|
||||||
|
app.loadGraphData(comfyUIWorkflow);
|
||||||
|
}
|
||||||
|
} else if (message.type === "deploy") {
|
||||||
|
// deployWorkflow();
|
||||||
|
const prompt = await app.graphToPrompt();
|
||||||
|
sendEventToCD("cd_plugin_onDeployChanges", prompt);
|
||||||
|
} else if (message.type === "queue_prompt") {
|
||||||
|
const prompt = await app.graphToPrompt();
|
||||||
|
sendEventToCD("cd_plugin_onQueuePrompt", prompt);
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
// console.error("Error processing message:", error);
|
||||||
|
}
|
||||||
|
|
||||||
|
// if (!event.data.flow || Object.entries(event.data.flow).length <= 0)
|
||||||
|
// return;
|
||||||
// updateBlendshapesPrompts(event.data.flow);
|
// updateBlendshapesPrompts(event.data.flow);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -167,6 +203,17 @@ const ext = {
|
|||||||
|
|
||||||
// }
|
// }
|
||||||
});
|
});
|
||||||
|
|
||||||
|
app.graph.onAfterChange = ((originalFunction) => async function () {
|
||||||
|
const prompt = await app.graphToPrompt();
|
||||||
|
sendEventToCD("cd_plugin_onAfterChange", prompt);
|
||||||
|
|
||||||
|
if (typeof originalFunction === "function") {
|
||||||
|
originalFunction.apply(this, arguments);
|
||||||
|
}
|
||||||
|
})(app.graph.onAfterChange);
|
||||||
|
|
||||||
|
sendEventToCD("cd_plugin_setup");
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -267,14 +314,9 @@ function createDynamicUIHtml(data) {
|
|||||||
return html;
|
return html;
|
||||||
}
|
}
|
||||||
|
|
||||||
function addButton() {
|
async function deployWorkflow() {
|
||||||
const menu = document.querySelector(".comfy-menu");
|
const deploy = document.getElementById("deploy-button");
|
||||||
|
|
||||||
const deploy = document.createElement("button");
|
|
||||||
deploy.style.position = "relative";
|
|
||||||
deploy.style.display = "block";
|
|
||||||
deploy.innerHTML = "<div id='button-title'>Deploy</div>";
|
|
||||||
deploy.onclick = async () => {
|
|
||||||
/** @type {LGraph} */
|
/** @type {LGraph} */
|
||||||
const graph = app.graph;
|
const graph = app.graph;
|
||||||
|
|
||||||
@@ -285,12 +327,37 @@ function addButton() {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let deployMeta = graph.findNodesByType("ComfyDeploy");
|
||||||
|
|
||||||
|
if (deployMeta.length == 0) {
|
||||||
|
const text = await inputDialog.input(
|
||||||
|
"Create your deployment",
|
||||||
|
"Workflow name",
|
||||||
|
);
|
||||||
|
if (!text) return;
|
||||||
|
console.log(text);
|
||||||
|
app.graph.beforeChange();
|
||||||
|
var node = LiteGraph.createNode("ComfyDeploy");
|
||||||
|
node.configure({
|
||||||
|
widgets_values: [text],
|
||||||
|
});
|
||||||
|
node.pos = [0, 0];
|
||||||
|
app.graph.add(node);
|
||||||
|
app.graph.afterChange();
|
||||||
|
deployMeta = [node];
|
||||||
|
}
|
||||||
|
|
||||||
|
const deployMetaNode = deployMeta[0];
|
||||||
|
|
||||||
|
const workflow_name = deployMetaNode.widgets[0].value;
|
||||||
|
const workflow_id = deployMetaNode.widgets[1].value;
|
||||||
|
|
||||||
const ok = await confirmDialog.confirm(
|
const ok = await confirmDialog.confirm(
|
||||||
`Confirm deployment`,
|
`Confirm deployment`,
|
||||||
`
|
`
|
||||||
<div>
|
<div>
|
||||||
|
|
||||||
A new version will be deployed, do you confirm?
|
A new version of <button style="font-size: 18px;">${workflow_name}</button> will be deployed, do you confirm?
|
||||||
<br><br>
|
<br><br>
|
||||||
|
|
||||||
<button style="font-size: 18px;">${displayName}</button>
|
<button style="font-size: 18px;">${displayName}</button>
|
||||||
@@ -332,31 +399,6 @@ function addButton() {
|
|||||||
|
|
||||||
const title = deploy.querySelector("#button-title");
|
const title = deploy.querySelector("#button-title");
|
||||||
|
|
||||||
let deployMeta = graph.findNodesByType("ComfyDeploy");
|
|
||||||
|
|
||||||
if (deployMeta.length == 0) {
|
|
||||||
const text = await inputDialog.input(
|
|
||||||
"Create your deployment",
|
|
||||||
"Workflow name",
|
|
||||||
);
|
|
||||||
if (!text) return;
|
|
||||||
console.log(text);
|
|
||||||
app.graph.beforeChange();
|
|
||||||
var node = LiteGraph.createNode("ComfyDeploy");
|
|
||||||
node.configure({
|
|
||||||
widgets_values: [text],
|
|
||||||
});
|
|
||||||
node.pos = [0, 0];
|
|
||||||
app.graph.add(node);
|
|
||||||
app.graph.afterChange();
|
|
||||||
deployMeta = [node];
|
|
||||||
}
|
|
||||||
|
|
||||||
const deployMetaNode = deployMeta[0];
|
|
||||||
|
|
||||||
const workflow_name = deployMetaNode.widgets[0].value;
|
|
||||||
const workflow_id = deployMetaNode.widgets[1].value;
|
|
||||||
|
|
||||||
const prompt = await app.graphToPrompt();
|
const prompt = await app.graphToPrompt();
|
||||||
let deps = undefined;
|
let deps = undefined;
|
||||||
|
|
||||||
@@ -474,7 +516,7 @@ function addButton() {
|
|||||||
<div style="position: absolute; top: 50%; left: 50%; transform: translate(-50%, -50%);">${loadingIcon}</div>
|
<div style="position: absolute; top: 50%; left: 50%; transform: translate(-50%, -50%);">${loadingIcon}</div>
|
||||||
<iframe
|
<iframe
|
||||||
style="z-index: 10; min-width: 600px; max-width: 1024px; min-height: 600px; border: none; background-color: transparent;"
|
style="z-index: 10; min-width: 600px; max-width: 1024px; min-height: 600px; border: none; background-color: transparent;"
|
||||||
src="${endpoint}/dependency-graph?deps=${encodeURIComponent(
|
src="https://www.comfydeploy.com/dependency-graph?deps=${encodeURIComponent(
|
||||||
JSON.stringify(deps),
|
JSON.stringify(deps),
|
||||||
)}" />`,
|
)}" />`,
|
||||||
// createDynamicUIHtml(deps),
|
// createDynamicUIHtml(deps),
|
||||||
@@ -555,6 +597,18 @@ function addButton() {
|
|||||||
title.style.color = "white";
|
title.style.color = "white";
|
||||||
}, 1000);
|
}, 1000);
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
function addButton() {
|
||||||
|
const menu = document.querySelector(".comfy-menu");
|
||||||
|
|
||||||
|
const deploy = document.createElement("button");
|
||||||
|
deploy.id = "deploy-button";
|
||||||
|
deploy.style.position = "relative";
|
||||||
|
deploy.style.display = "block";
|
||||||
|
deploy.innerHTML = "<div id='button-title'>Deploy</div>";
|
||||||
|
deploy.onclick = async () => {
|
||||||
|
await deployWorkflow()
|
||||||
};
|
};
|
||||||
|
|
||||||
const config = document.createElement("img");
|
const config = document.createElement("img");
|
||||||
@@ -880,7 +934,10 @@ export class ConfigDialog extends ComfyDialog {
|
|||||||
justifyContent: "flex-end",
|
justifyContent: "flex-end",
|
||||||
width: "100%",
|
width: "100%",
|
||||||
},
|
},
|
||||||
onclick: () => this.save(),
|
onclick: () => {
|
||||||
|
this.save();
|
||||||
|
this.close();
|
||||||
|
},
|
||||||
},
|
},
|
||||||
[
|
[
|
||||||
$el("button", {
|
$el("button", {
|
||||||
@@ -891,7 +948,10 @@ export class ConfigDialog extends ComfyDialog {
|
|||||||
$el("button", {
|
$el("button", {
|
||||||
type: "button",
|
type: "button",
|
||||||
textContent: "Save",
|
textContent: "Save",
|
||||||
onclick: () => this.save(),
|
onclick: () => {
|
||||||
|
this.save();
|
||||||
|
this.close();
|
||||||
|
},
|
||||||
}),
|
}),
|
||||||
],
|
],
|
||||||
),
|
),
|
||||||
@@ -905,20 +965,26 @@ export class ConfigDialog extends ComfyDialog {
|
|||||||
}
|
}
|
||||||
|
|
||||||
save(api_key, displayName) {
|
save(api_key, displayName) {
|
||||||
if (!displayName) displayName = getData().displayName;
|
|
||||||
|
|
||||||
const deployOption = this.container.querySelector("#deployOption").value;
|
const deployOption = this.container.querySelector("#deployOption").value;
|
||||||
localStorage.setItem("comfy_deploy_env", deployOption);
|
localStorage.setItem("comfy_deploy_env", deployOption);
|
||||||
|
|
||||||
const endpoint = this.container.querySelector("#endpoint").value;
|
const endpoint = this.container.querySelector("#endpoint").value;
|
||||||
const apiKey = api_key ?? this.container.querySelector("#apiKey").value;
|
const apiKey = api_key ?? this.container.querySelector("#apiKey").value;
|
||||||
|
|
||||||
|
if (!displayName) {
|
||||||
|
if (apiKey != getData().apiKey) {
|
||||||
|
displayName = "Custom";
|
||||||
|
} else {
|
||||||
|
displayName = getData().displayName;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
saveData({
|
saveData({
|
||||||
endpoint,
|
endpoint,
|
||||||
apiKey,
|
apiKey,
|
||||||
displayName,
|
displayName,
|
||||||
environment: deployOption,
|
environment: deployOption,
|
||||||
});
|
});
|
||||||
this.close();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
show() {
|
show() {
|
||||||
@@ -941,8 +1007,10 @@ export class ConfigDialog extends ComfyDialog {
|
|||||||
data.endpoint
|
data.endpoint
|
||||||
}">
|
}">
|
||||||
</label>
|
</label>
|
||||||
<label style="color: white;">
|
<div style="color: white;">
|
||||||
API Key: ${data.displayName ?? ""}
|
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="${
|
<input id="apiKey" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;" type="password" value="${
|
||||||
data.apiKey
|
data.apiKey
|
||||||
}">
|
}">
|
||||||
@@ -951,12 +1019,15 @@ export class ConfigDialog extends ComfyDialog {
|
|||||||
data.apiKey ? "Re-login with ComfyDeploy" : "Login with ComfyDeploy"
|
data.apiKey ? "Re-login with ComfyDeploy" : "Login with ComfyDeploy"
|
||||||
}
|
}
|
||||||
</button>
|
</button>
|
||||||
</label>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
`;
|
`;
|
||||||
|
|
||||||
const button = this.container.querySelector("#loginButton");
|
const button = this.container.querySelector("#loginButton");
|
||||||
button.onclick = () => {
|
button.onclick = () => {
|
||||||
|
this.save();
|
||||||
|
const data = getData();
|
||||||
|
|
||||||
const uuid =
|
const uuid =
|
||||||
Math.random().toString(36).substring(2, 15) +
|
Math.random().toString(36).substring(2, 15) +
|
||||||
Math.random().toString(36).substring(2, 15);
|
Math.random().toString(36).substring(2, 15);
|
||||||
@@ -973,17 +1044,20 @@ export class ConfigDialog extends ComfyDialog {
|
|||||||
this.poll = setInterval(() => {
|
this.poll = setInterval(() => {
|
||||||
fetch(data.endpoint + "/api/auth-response/" + uuid)
|
fetch(data.endpoint + "/api/auth-response/" + uuid)
|
||||||
.then((response) => response.json())
|
.then((response) => response.json())
|
||||||
.then((json) => {
|
.then(async (json) => {
|
||||||
if (json.api_key) {
|
if (json.api_key) {
|
||||||
this.save(json.api_key, json.name);
|
this.save(json.api_key, json.name);
|
||||||
|
this.close();
|
||||||
this.container.querySelector("#apiKey").value = json.api_key;
|
this.container.querySelector("#apiKey").value = json.api_key;
|
||||||
infoDialog.show();
|
// infoDialog.show();
|
||||||
clearInterval(this.poll);
|
clearInterval(this.poll);
|
||||||
clearTimeout(this.timeout);
|
clearTimeout(this.timeout);
|
||||||
infoDialog.showMessage(
|
// Refresh dialog
|
||||||
|
const a = await confirmDialog.confirm(
|
||||||
"Authenticated",
|
"Authenticated",
|
||||||
"You will be able to upload workflow to " + json.name,
|
`<div>You will be able to upload workflow to <button style="font-size: 18px; width: fit;">${json.name}</button></div>`,
|
||||||
);
|
);
|
||||||
|
configDialog.show();
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
.catch((error) => {
|
.catch((error) => {
|
||||||
|
|||||||
Reference in New Issue
Block a user