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": {
|
||||
"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"
|
||||
|
||||
def run(self, input_id, default_checkpoint_name=None):
|
||||
def run(self, input_id, default_value=None):
|
||||
import requests
|
||||
import os
|
||||
import uuid
|
||||
|
||||
if input_id and input_id.startswith('http'):
|
||||
if default_value.startswith('http'):
|
||||
unique_filename = str(uuid.uuid4()) + ".safetensors"
|
||||
print(unique_filename)
|
||||
print(folder_paths.folder_names_and_paths["checkpoints"][0][0])
|
||||
@@ -59,7 +59,7 @@ class ComfyUIDeployExternalCheckpoint:
|
||||
out_file.write(chunk)
|
||||
return (unique_filename,)
|
||||
else:
|
||||
return (default_checkpoints_name,)
|
||||
return (default_value,)
|
||||
|
||||
|
||||
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 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"
|
||||
print(unique_filename)
|
||||
print(folder_paths.folder_names_and_paths["loras"][0][0])
|
||||
@@ -44,6 +49,8 @@ class ComfyUIDeployExternalLora:
|
||||
out_file.write(response.content)
|
||||
return (unique_filename,)
|
||||
else:
|
||||
return (input_id,)
|
||||
|
||||
return (default_lora_name,)
|
||||
|
||||
|
||||
|
||||
@@ -16,7 +16,7 @@ class ComfyUIDeployExternalNumber:
|
||||
"optional": {
|
||||
"default_value": (
|
||||
"FLOAT",
|
||||
{"multiline": True, "display": "number", "default": 0},
|
||||
{"multiline": True, "display": "number", "default": 0, "step": 0.01},
|
||||
),
|
||||
}
|
||||
}
|
||||
@@ -29,9 +29,12 @@ class ComfyUIDeployExternalNumber:
|
||||
CATEGORY = "number"
|
||||
|
||||
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 [int(input_id)]
|
||||
|
||||
|
||||
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
|
||||
import os
|
||||
import requests
|
||||
import folder_paths
|
||||
import json
|
||||
import numpy as np
|
||||
import server
|
||||
import re
|
||||
import base64
|
||||
from PIL import Image
|
||||
import io
|
||||
import time
|
||||
import execution
|
||||
import random
|
||||
import traceback
|
||||
import uuid
|
||||
import asyncio
|
||||
import atexit
|
||||
import logging
|
||||
import sys
|
||||
from logging.handlers import RotatingFileHandler
|
||||
from enum import Enum
|
||||
from urllib.parse import quote
|
||||
import threading
|
||||
import hashlib
|
||||
import aiohttp
|
||||
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_task = None
|
||||
prompt_metadata = {}
|
||||
|
||||
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'
|
||||
|
||||
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):
|
||||
prompt_server = server.PromptServer.instance
|
||||
json_data = prompt_server.trigger_on_prompt(json_data)
|
||||
@@ -80,19 +90,105 @@ def randomSeed(num_digits=15):
|
||||
range_end = (10**num_digits) - 1
|
||||
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")
|
||||
async def comfy_deploy_run(request):
|
||||
prompt_server = server.PromptServer.instance
|
||||
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
|
||||
prompt_id = data.get("prompt_id")
|
||||
inputs = data.get("inputs")
|
||||
|
||||
for key in workflow_api:
|
||||
if 'inputs' in workflow_api[key] and 'seed' in workflow_api[key]['inputs']:
|
||||
workflow_api[key]['inputs']['seed'] = randomSeed()
|
||||
# Now it handles directly in here
|
||||
apply_random_seed_to_workflow(workflow_api)
|
||||
apply_inputs_to_workflow(workflow_api, inputs)
|
||||
|
||||
prompt = {
|
||||
"prompt": workflow_api,
|
||||
@@ -100,11 +196,11 @@ async def comfy_deploy_run(request):
|
||||
"prompt_id": prompt_id
|
||||
}
|
||||
|
||||
prompt_metadata[prompt_id] = {
|
||||
'status_endpoint': data.get('status_endpoint'),
|
||||
'file_upload_endpoint': data.get('file_upload_endpoint'),
|
||||
'workflow_api': workflow_api
|
||||
}
|
||||
prompt_metadata[prompt_id] = SimplePrompt(
|
||||
status_endpoint=data.get('status_endpoint'),
|
||||
file_upload_endpoint=data.get('file_upload_endpoint'),
|
||||
workflow_api=workflow_api
|
||||
)
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
@@ -148,7 +244,6 @@ async def comfy_deploy_run(request):
|
||||
|
||||
return web.json_response(res, status=status)
|
||||
|
||||
sockets = dict()
|
||||
|
||||
def get_comfyui_path_from_file_path(file_path):
|
||||
file_path_parts = file_path.split("\\")
|
||||
@@ -179,63 +274,22 @@ async def compute_sha256_checksum(filepath):
|
||||
sha256.update(chunk)
|
||||
return sha256.hexdigest()
|
||||
|
||||
# def hash_chunk(start_end, filepath):
|
||||
# """Hash a specific chunk of the file."""
|
||||
# start, end = start_end
|
||||
# sha256 = hashlib.sha256()
|
||||
# with open(filepath, 'rb') as f:
|
||||
# f.seek(start)
|
||||
# chunk = f.read(end - start)
|
||||
# sha256.update(chunk)
|
||||
# return sha256.digest() # Return the digest of the chunk
|
||||
|
||||
# async def compute_sha256_checksum(filepath):
|
||||
# file_size = os.path.getsize(filepath)
|
||||
# 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
|
||||
@server.PromptServer.instance.routes.get('/comfyui-deploy/models')
|
||||
async def get_installed_models(request):
|
||||
# Directly return the list of paths as JSON
|
||||
new_dict = {}
|
||||
for key, value in folder_paths.folder_names_and_paths.items():
|
||||
# Convert set to list for JSON compatibility
|
||||
# for path in value[0]:
|
||||
file_list = folder_paths.get_filename_list(key)
|
||||
value_json_compatible = (value[0], list(value[1]), file_list)
|
||||
new_dict[key] = value_json_compatible
|
||||
# print(new_dict)
|
||||
return web.json_response(new_dict)
|
||||
|
||||
# This is start uploading the files to Comfy Deploy
|
||||
@server.PromptServer.instance.routes.post('/comfyui-deploy/upload-file')
|
||||
async def upload_file(request):
|
||||
async def upload_file_endpoint(request):
|
||||
data = await request.json()
|
||||
|
||||
file_path = data.get("file_path")
|
||||
@@ -317,26 +371,58 @@ async def upload_file(request):
|
||||
}, 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')
|
||||
async def get_file_hash(request):
|
||||
file_path = request.rel_url.query.get('file_path', '')
|
||||
|
||||
if file_path is None:
|
||||
if not file_path:
|
||||
return web.json_response({
|
||||
"error": "file_path is required"
|
||||
}, status=400)
|
||||
|
||||
try:
|
||||
base = folder_paths.base_path
|
||||
file_path = os.path.join(base, file_path)
|
||||
# print("file_path", file_path)
|
||||
start_time = time.time() # Capture the start time
|
||||
file_hash = await compute_sha256_checksum(
|
||||
file_path
|
||||
)
|
||||
end_time = time.time() # Capture the end time after the code execution
|
||||
elapsed_time = end_time - start_time # Calculate the elapsed time
|
||||
print(f"Execution time: {elapsed_time} seconds")
|
||||
full_file_path = os.path.join(base, file_path)
|
||||
|
||||
# Check if the file hash is in the cache
|
||||
if full_file_path in file_hash_cache:
|
||||
file_hash = file_hash_cache[full_file_path]
|
||||
else:
|
||||
start_time = time.time()
|
||||
file_hash = await compute_sha256_checksum(full_file_path)
|
||||
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({
|
||||
"file_hash": file_hash
|
||||
})
|
||||
@@ -345,6 +431,16 @@ async def get_file_hash(request):
|
||||
"error": str(e)
|
||||
}, 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')
|
||||
async def websocket_handler(request):
|
||||
ws = web.WebSocketResponse()
|
||||
@@ -358,35 +454,134 @@ async def websocket_handler(request):
|
||||
|
||||
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:
|
||||
# Send initial state to the new client
|
||||
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)
|
||||
|
||||
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:
|
||||
print('ws connection closed with exception %s' % ws.exception())
|
||||
finally:
|
||||
sockets.pop(sid, None)
|
||||
|
||||
if realtime_id is not None:
|
||||
await update_realtime_run_status(realtime_id, status_endpoint, Status.SUCCESS)
|
||||
return ws
|
||||
|
||||
@server.PromptServer.instance.routes.get('/comfyui-deploy/check-status')
|
||||
async def comfy_deploy_check_status(request):
|
||||
prompt_server = server.PromptServer.instance
|
||||
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({
|
||||
"status": prompt_metadata[prompt_id]['status'].value
|
||||
"status": prompt_metadata[prompt_id].status.value
|
||||
})
|
||||
else:
|
||||
return web.json_response({
|
||||
"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):
|
||||
try:
|
||||
# message = {"event": event, "data": data}
|
||||
if sid:
|
||||
ws = sockets.get(sid)
|
||||
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)
|
||||
|
||||
prompt_server = server.PromptServer.instance
|
||||
|
||||
send_json = prompt_server.send_json
|
||||
|
||||
async def send_json_override(self, event, data, sid=None):
|
||||
# print("INTERNAL:", event, data, sid)
|
||||
prompt_id = data.get('prompt_id')
|
||||
|
||||
target_sid = sid
|
||||
if target_sid == "comfy_deploy_instance":
|
||||
target_sid = None
|
||||
|
||||
# now we send everything
|
||||
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))
|
||||
])
|
||||
|
||||
if event == 'execution_start':
|
||||
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
|
||||
if event == 'executing' and data.get('node') is None:
|
||||
mark_prompt_done(prompt_id=prompt_id)
|
||||
if not have_pending_upload(prompt_id):
|
||||
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:
|
||||
node = data.get('node')
|
||||
|
||||
if prompt_id in prompt_metadata:
|
||||
if 'progress' not in prompt_metadata[prompt_id]:
|
||||
prompt_metadata[prompt_id]["progress"] = set()
|
||||
# if 'progress' not in prompt_metadata[prompt_id]:
|
||||
# prompt_metadata[prompt_id]["progress"] = set()
|
||||
|
||||
prompt_metadata[prompt_id]["progress"].add(node)
|
||||
calculated_progress = len(prompt_metadata[prompt_id]["progress"]) / len(prompt_metadata[prompt_id]['workflow_api'])
|
||||
prompt_metadata[prompt_id].progress.add(node)
|
||||
calculated_progress = len(prompt_metadata[prompt_id].progress) / len(prompt_metadata[prompt_id].workflow_api)
|
||||
# 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
|
||||
prompt_metadata[prompt_id]['last_updated_node'] = node
|
||||
class_type = prompt_metadata[prompt_id]['workflow_api'][node]['class_type']
|
||||
prompt_metadata[prompt_id].last_updated_node = node
|
||||
class_type = prompt_metadata[prompt_id].workflow_api[node]['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)
|
||||
|
||||
if event == 'execution_cached' and data.get('nodes') is not None:
|
||||
if prompt_id in prompt_metadata:
|
||||
if 'progress' not in prompt_metadata[prompt_id]:
|
||||
prompt_metadata[prompt_id]["progress"] = set()
|
||||
# if 'progress' not in prompt_metadata[prompt_id]:
|
||||
# prompt_metadata[prompt_id].progress = set()
|
||||
|
||||
if 'nodes' in data:
|
||||
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'))
|
||||
|
||||
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)
|
||||
|
||||
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'))
|
||||
# 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
|
||||
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:
|
||||
return
|
||||
|
||||
if prompt_metadata[prompt_id].is_realtime is True:
|
||||
return
|
||||
|
||||
print("progress", calculated_progress)
|
||||
|
||||
status_endpoint = prompt_metadata[prompt_id]['status_endpoint']
|
||||
status_endpoint = prompt_metadata[prompt_id].status_endpoint
|
||||
body = {
|
||||
"run_id": prompt_id,
|
||||
"live_status": live_status,
|
||||
@@ -491,19 +711,25 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
|
||||
pass
|
||||
|
||||
|
||||
def update_run(prompt_id, status: Status):
|
||||
def update_run(prompt_id: str, status: Status):
|
||||
global last_read_line_number
|
||||
|
||||
if prompt_id not in prompt_metadata:
|
||||
return
|
||||
|
||||
if ('status' not in prompt_metadata[prompt_id] or prompt_metadata[prompt_id]['status'] != status):
|
||||
|
||||
# when the status is already failed, we don't want to update it to success
|
||||
if ('status' in prompt_metadata[prompt_id] and prompt_metadata[prompt_id]['status'] == Status.FAILED):
|
||||
# if prompt_metadata[prompt_id].start_time is None and status == Status.RUNNING:
|
||||
# if its realtime prompt we need to skip that.
|
||||
if prompt_metadata[prompt_id].is_realtime is True:
|
||||
prompt_metadata[prompt_id].status = status
|
||||
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 = {
|
||||
"run_id": prompt_id,
|
||||
"status": status.value,
|
||||
@@ -544,13 +770,12 @@ def update_run(prompt_id, status: Status):
|
||||
except Exception as log_error:
|
||||
print(f"Error reading log file: {log_error}")
|
||||
|
||||
|
||||
except Exception as e:
|
||||
error_type = type(e).__name__
|
||||
stack_trace = traceback.format_exc().strip()
|
||||
print(f"Error occurred while updating run: {e} {stack_trace}")
|
||||
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"):
|
||||
@@ -582,7 +807,7 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
|
||||
|
||||
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)
|
||||
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)
|
||||
|
||||
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:
|
||||
print("have pending upload ", len(prompt_metadata[prompt_id]['uploading_nodes']))
|
||||
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))
|
||||
return True
|
||||
|
||||
print("no pending upload")
|
||||
return False
|
||||
|
||||
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:
|
||||
prompt_metadata[prompt_id]["done"] = True
|
||||
prompt_metadata[prompt_id].done = True
|
||||
print("Prompt done")
|
||||
|
||||
def is_prompt_done(prompt_id):
|
||||
if prompt_id in prompt_metadata and "done" in prompt_metadata[prompt_id]:
|
||||
if prompt_metadata[prompt_id]["done"] == True:
|
||||
def is_prompt_done(prompt_id: str):
|
||||
"""
|
||||
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 False
|
||||
@@ -644,17 +883,17 @@ async def handle_error(prompt_id, data, e: Exception):
|
||||
print(f"Error occurred while uploading file: {e}")
|
||||
|
||||
# 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):
|
||||
if 'uploading_nodes' not in prompt_metadata[prompt_id]:
|
||||
prompt_metadata[prompt_id]['uploading_nodes'] = set()
|
||||
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]:
|
||||
# prompt_metadata[prompt_id]['uploading_nodes'] = set()
|
||||
|
||||
if node_id is not None:
|
||||
if uploading:
|
||||
prompt_metadata[prompt_id]['uploading_nodes'].add(node_id)
|
||||
prompt_metadata[prompt_id].uploading_nodes.add(node_id)
|
||||
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
|
||||
|
||||
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 uploading:
|
||||
if prompt_metadata[prompt_id]['status'] != Status.UPLOADING:
|
||||
if prompt_metadata[prompt_id].status != Status.UPLOADING:
|
||||
update_run(prompt_id, Status.UPLOADING)
|
||||
await send("uploading", {
|
||||
"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,
|
||||
})
|
||||
|
||||
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, [])
|
||||
for item in items:
|
||||
# # Skipping temp files
|
||||
if item.get("type") == "temp":
|
||||
continue
|
||||
await upload_file(
|
||||
prompt_id,
|
||||
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)
|
||||
)
|
||||
|
||||
|
||||
# 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:
|
||||
await handle_upload(prompt_id, data, 'images', "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)
|
||||
|
||||
async def update_run_with_output(prompt_id, data, node_id=None):
|
||||
if prompt_id in prompt_metadata:
|
||||
status_endpoint = prompt_metadata[prompt_id]['status_endpoint']
|
||||
if prompt_id not in prompt_metadata:
|
||||
return
|
||||
|
||||
if prompt_metadata[prompt_id].is_realtime is True:
|
||||
return
|
||||
|
||||
status_endpoint = prompt_metadata[prompt_id].status_endpoint
|
||||
|
||||
body = {
|
||||
"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")
|
||||
setup()
|
||||
|
||||
|
||||
# Store the original working directory
|
||||
original_cwd = os.getcwd()
|
||||
try:
|
||||
# Get the absolute path of the script's directory
|
||||
script_dir = os.path.dirname(os.path.abspath(__file__))
|
||||
@@ -67,3 +70,6 @@ try:
|
||||
print(f"** Comfy Deploy Revision: {current_git_commit}")
|
||||
except Exception as 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
|
||||
pydantic
|
||||
+126
-52
@@ -1,10 +1,18 @@
|
||||
import { app } from "./app.js";
|
||||
import { api } from "./api.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>`;
|
||||
|
||||
function sendEventToCD(event, data) {
|
||||
const message = {
|
||||
type: event,
|
||||
data: data,
|
||||
};
|
||||
window.parent.postMessage(JSON.stringify(message), "*");
|
||||
}
|
||||
|
||||
/** @typedef {import('../../../web/types/comfy.js').ComfyExtension} ComfyExtension*/
|
||||
/** @type {ComfyExtension} */
|
||||
const ext = {
|
||||
@@ -18,6 +26,11 @@ const ext = {
|
||||
const auth_token = queryParams.get("auth_token");
|
||||
const org_display = queryParams.get("org_display");
|
||||
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();
|
||||
let endpoint = data.endpoint;
|
||||
@@ -59,8 +72,8 @@ const ext = {
|
||||
return;
|
||||
}
|
||||
|
||||
// Adding a delay to wait for the intial graph to load
|
||||
await new Promise((resolve) => setTimeout(resolve, 2000));
|
||||
// // Adding a delay to wait for the intial graph to load
|
||||
// await new Promise((resolve) => setTimeout(resolve, 2000));
|
||||
|
||||
workflow?.nodes.forEach((x) => {
|
||||
if (x?.type === "ComfyDeploy") {
|
||||
@@ -152,9 +165,32 @@ const ext = {
|
||||
async setup() {
|
||||
// const graphCanvas = document.getElementById("graph-canvas");
|
||||
|
||||
window.addEventListener("message", (event) => {
|
||||
if (!event.data.flow || Object.entries(event.data.flow).length <= 0)
|
||||
return;
|
||||
window.addEventListener("message", async (event) => {
|
||||
try {
|
||||
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);
|
||||
});
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
function addButton() {
|
||||
const menu = document.querySelector(".comfy-menu");
|
||||
async function deployWorkflow() {
|
||||
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} */
|
||||
const graph = app.graph;
|
||||
|
||||
@@ -285,12 +327,37 @@ function addButton() {
|
||||
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(
|
||||
`Confirm deployment`,
|
||||
`
|
||||
<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>
|
||||
|
||||
<button style="font-size: 18px;">${displayName}</button>
|
||||
@@ -332,31 +399,6 @@ function addButton() {
|
||||
|
||||
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();
|
||||
let deps = undefined;
|
||||
|
||||
@@ -474,7 +516,7 @@ function addButton() {
|
||||
<div style="position: absolute; top: 50%; left: 50%; transform: translate(-50%, -50%);">${loadingIcon}</div>
|
||||
<iframe
|
||||
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),
|
||||
)}" />`,
|
||||
// createDynamicUIHtml(deps),
|
||||
@@ -555,6 +597,18 @@ function addButton() {
|
||||
title.style.color = "white";
|
||||
}, 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");
|
||||
@@ -880,7 +934,10 @@ export class ConfigDialog extends ComfyDialog {
|
||||
justifyContent: "flex-end",
|
||||
width: "100%",
|
||||
},
|
||||
onclick: () => this.save(),
|
||||
onclick: () => {
|
||||
this.save();
|
||||
this.close();
|
||||
},
|
||||
},
|
||||
[
|
||||
$el("button", {
|
||||
@@ -891,7 +948,10 @@ export class ConfigDialog extends ComfyDialog {
|
||||
$el("button", {
|
||||
type: "button",
|
||||
textContent: "Save",
|
||||
onclick: () => this.save(),
|
||||
onclick: () => {
|
||||
this.save();
|
||||
this.close();
|
||||
},
|
||||
}),
|
||||
],
|
||||
),
|
||||
@@ -905,20 +965,26 @@ export class ConfigDialog extends ComfyDialog {
|
||||
}
|
||||
|
||||
save(api_key, displayName) {
|
||||
if (!displayName) displayName = getData().displayName;
|
||||
|
||||
const deployOption = this.container.querySelector("#deployOption").value;
|
||||
localStorage.setItem("comfy_deploy_env", deployOption);
|
||||
|
||||
const endpoint = this.container.querySelector("#endpoint").value;
|
||||
const apiKey = api_key ?? this.container.querySelector("#apiKey").value;
|
||||
|
||||
if (!displayName) {
|
||||
if (apiKey != getData().apiKey) {
|
||||
displayName = "Custom";
|
||||
} else {
|
||||
displayName = getData().displayName;
|
||||
}
|
||||
}
|
||||
|
||||
saveData({
|
||||
endpoint,
|
||||
apiKey,
|
||||
displayName,
|
||||
environment: deployOption,
|
||||
});
|
||||
this.close();
|
||||
}
|
||||
|
||||
show() {
|
||||
@@ -941,8 +1007,10 @@ export class ConfigDialog extends ComfyDialog {
|
||||
data.endpoint
|
||||
}">
|
||||
</label>
|
||||
<label style="color: white;">
|
||||
API Key: ${data.displayName ?? ""}
|
||||
<div style="color: white;">
|
||||
API Key: User / Org <button style="font-size: 18px;">${
|
||||
data.displayName ?? ""
|
||||
}</button>
|
||||
<input id="apiKey" style="margin-top: 8px; width: 100%; height:40px; box-sizing: border-box; padding: 0px 6px;" type="password" value="${
|
||||
data.apiKey
|
||||
}">
|
||||
@@ -951,12 +1019,15 @@ export class ConfigDialog extends ComfyDialog {
|
||||
data.apiKey ? "Re-login with ComfyDeploy" : "Login with ComfyDeploy"
|
||||
}
|
||||
</button>
|
||||
</label>
|
||||
</div>
|
||||
</div>
|
||||
`;
|
||||
|
||||
const button = this.container.querySelector("#loginButton");
|
||||
button.onclick = () => {
|
||||
this.save();
|
||||
const data = getData();
|
||||
|
||||
const uuid =
|
||||
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(() => {
|
||||
fetch(data.endpoint + "/api/auth-response/" + uuid)
|
||||
.then((response) => response.json())
|
||||
.then((json) => {
|
||||
.then(async (json) => {
|
||||
if (json.api_key) {
|
||||
this.save(json.api_key, json.name);
|
||||
this.close();
|
||||
this.container.querySelector("#apiKey").value = json.api_key;
|
||||
infoDialog.show();
|
||||
// infoDialog.show();
|
||||
clearInterval(this.poll);
|
||||
clearTimeout(this.timeout);
|
||||
infoDialog.showMessage(
|
||||
// Refresh dialog
|
||||
const a = await confirmDialog.confirm(
|
||||
"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) => {
|
||||
|
||||
Reference in New Issue
Block a user