Compare commits

..
Author SHA1 Message Date
BennyKok 8c8f2abc16 Merge branch 'main' into dev 2024-07-07 22:06:54 -07:00
nick c4d1b09a24 custom route 2024-06-19 16:52:17 -07:00
bennykok c70e08a706 chore(plugin): add log 2024-06-11 17:43:03 -07:00
bennykok daf1669e70 fix: node_error proxy 2024-06-11 17:43:02 -07:00
bennykok 62df715655 fix: prompt error 2024-06-11 17:43:02 -07:00
bennykok 04fd08d5ba fix: streaming event format 2024-06-11 17:43:02 -07:00
bennykok 4a8ef7c77c fix(plugin): event 2024-06-11 17:43:02 -07:00
bennykok 5b8dac37fb feat(plugin): add dispatchAPIEventData 2024-06-11 17:43:02 -07:00
bennykok 875f7f24d1 fix: run issues 2024-06-11 17:43:02 -07:00
bennykok af0fac7afc feat: add streaming endpoint 2024-06-11 17:43:02 -07:00
18 changed files with 264 additions and 1157 deletions
+1 -11
View File
@@ -8,16 +8,6 @@ class ComfyUIDeployExternalBoolean:
{"multiline": False, "default": "input_bool"},
),
"default_value": ("BOOLEAN", {"default": False})
},
"optional": {
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -26,7 +16,7 @@ class ComfyUIDeployExternalBoolean:
FUNCTION = "run"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
print(f"Node '{input_id}' processing with switch set to {default_value}")
return [default_value]
+2 -16
View File
@@ -5,12 +5,6 @@ import torch
import folder_paths
from tqdm import tqdm
class AnyType(str):
def __ne__(self, __value: object) -> bool:
return False
WILDCARD = AnyType("*")
class ComfyUIDeployExternalCheckpoint:
@classmethod
def INPUT_TYPES(s):
@@ -23,25 +17,17 @@ class ComfyUIDeployExternalCheckpoint:
},
"optional": {
"default_value": (folder_paths.get_filename_list("checkpoints"), ),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
RETURN_TYPES = (WILDCARD,)
RETURN_TYPES = (folder_paths.get_filename_list("checkpoints"),)
RETURN_NAMES = ("path",)
FUNCTION = "run"
CATEGORY = "deploy"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
import requests
import os
import uuid
+1 -9
View File
@@ -15,14 +15,6 @@ class ComfyUIDeployExternalImage:
},
"optional": {
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -33,7 +25,7 @@ class ComfyUIDeployExternalImage:
CATEGORY = "image"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
image = default_value
try:
if input_id.startswith('http'):
+1 -9
View File
@@ -15,14 +15,6 @@ class ComfyUIDeployExternalImageAlpha:
},
"optional": {
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -33,7 +25,7 @@ class ComfyUIDeployExternalImageAlpha:
CATEGORY = "image"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
image = default_value
try:
if input_id.startswith('http'):
+1 -9
View File
@@ -21,14 +21,6 @@ class ComfyUIDeployExternalImageBatch:
},
"optional": {
"default_value": ("IMAGE",),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -39,7 +31,7 @@ class ComfyUIDeployExternalImageBatch:
CATEGORY = "image"
def run(self, input_id, images=None, default_value=None, display_name=None, description=None):
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
+7 -38
View File
@@ -5,14 +5,6 @@ import torch
import folder_paths
class AnyType(str):
def __ne__(self, __value: object) -> bool:
return False
WILDCARD = AnyType("*")
class ComfyUIDeployExternalLora:
@classmethod
def INPUT_TYPES(s):
@@ -25,50 +17,27 @@ class ComfyUIDeployExternalLora:
},
"optional": {
"default_lora_name": (folder_paths.get_filename_list("loras"),),
"lora_save_name": ( # if `default_lora_name` is a link to download a file, we will attempt to save it with this name
"STRING",
{"multiline": False, "default": ""},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
"lora_url": (
"STRING",
{"multiline": False, "default": ""},
),
},
}
RETURN_TYPES = (WILDCARD,)
RETURN_TYPES = (folder_paths.get_filename_list("loras"),)
RETURN_NAMES = ("path",)
FUNCTION = "run"
CATEGORY = "deploy"
def run(self, input_id, default_lora_name=None, lora_save_name=None, display_name=None, description=None, lora_url=None):
def run(self, input_id, default_lora_name=None):
import requests
import os
import uuid
if lora_url and lora_url.startswith("http"):
if lora_save_name:
existing_loras = folder_paths.get_filename_list("loras")
# Check if lora_save_name exists in the list
if lora_save_name in existing_loras:
print(f"using lora: {lora_save_name}")
return (lora_save_name,)
else:
lora_save_name = str(uuid.uuid4()) + ".safetensors"
print(lora_save_name)
if default_lora_name.startswith("http"):
unique_filename = str(uuid.uuid4()) + ".safetensors"
print(unique_filename)
print(folder_paths.folder_names_and_paths["loras"][0][0])
destination_path = os.path.join(
folder_paths.folder_names_and_paths["loras"][0][0], lora_save_name
folder_paths.folder_names_and_paths["loras"][0][0], unique_filename
)
print(destination_path)
print("Downloading external lora - " + input_id + " to " + destination_path)
@@ -79,7 +48,7 @@ class ComfyUIDeployExternalLora:
)
with open(destination_path, "wb") as out_file:
out_file.write(response.content)
return (lora_save_name,)
return (unique_filename,)
else:
print(f"using lora: {default_lora_name}")
return (default_lora_name,)
+2 -10
View File
@@ -16,15 +16,7 @@ class ComfyUIDeployExternalNumber:
"optional": {
"default_value": (
"FLOAT",
{"multiline": True, "display": "number", "default": 0, "min": -2147483647, "max": 2147483647, "step": 0.01},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "display": "number", "default": 0, "step": 0.01},
),
}
}
@@ -36,7 +28,7 @@ class ComfyUIDeployExternalNumber:
CATEGORY = "number"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
try:
float_value = float(input_id)
print("my number", float_value)
+2 -10
View File
@@ -16,15 +16,7 @@ class ComfyUIDeployExternalNumberInt:
"optional": {
"default_value": (
"INT",
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 0},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "display": "number", "default": 0},
),
}
}
@@ -36,7 +28,7 @@ class ComfyUIDeployExternalNumberInt:
CATEGORY = "number"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
if not input_id or (isinstance(input_id, str) and not input_id.strip().isdigit()):
return [default_value]
return [int(input_id)]
+4 -12
View File
@@ -11,23 +11,15 @@ class ComfyUIDeployExternalNumberSlider:
"optional": {
"default_value": (
"FLOAT",
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 0.5, "step": 0.01},
{"multiline": True, "display": "number", "default": 0.5, "step": 0.01},
),
"min_value": (
"FLOAT",
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 0, "step": 0.01},
{"multiline": True, "display": "number", "default": 0, "step": 0.01},
),
"max_value": (
"FLOAT",
{"multiline": True, "display": "number", "min": -2147483647, "max": 2147483647, "default": 1, "step": 0.01},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
{"multiline": True, "display": "number", "default": 1, "step": 0.01},
),
}
}
@@ -39,7 +31,7 @@ class ComfyUIDeployExternalNumberSlider:
CATEGORY = "number"
def run(self, input_id, default_value=None, min_value=0, max_value=1, display_name=None, description=None):
def run(self, input_id, default_value=None, min_value=0, max_value=1):
try:
float_value = float(input_id)
if min_value <= float_value <= max_value:
+1 -9
View File
@@ -18,14 +18,6 @@ class ComfyUIDeployExternalText:
"STRING",
{"multiline": True, "default": ""},
),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
@@ -36,7 +28,7 @@ class ComfyUIDeployExternalText:
CATEGORY = "text"
def run(self, input_id, default_value=None, display_name=None, description=None):
def run(self, input_id, default_value=None):
return [default_value]
-52
View File
@@ -1,52 +0,0 @@
import folder_paths
from PIL import Image, ImageOps
import numpy as np
import torch
import json
class ComfyUIDeployExternalTextList:
@classmethod
def INPUT_TYPES(s):
return {
"required": {
"input_id": (
"STRING",
{"multiline": False, "default": 'input_text_list'},
),
"text": (
"STRING",
{"multiline": True, "default": "[]"},
),
},
"optional": {
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
}
}
RETURN_TYPES = ("STRING",)
RETURN_NAMES = ("text",)
OUTPUT_IS_LIST = (True,)
FUNCTION = "run"
CATEGORY = "text"
def run(self, input_id, text=None, display_name=None, description=None):
text_list = []
try:
text_list = json.loads(text) # Assuming text is a JSON array string
except Exception as e:
print(f"Error processing images: {e}")
pass
return ([text_list],)
NODE_CLASS_MAPPINGS = {"ComfyUIDeployExternalTextList": ComfyUIDeployExternalTextList}
NODE_DISPLAY_NAME_MAPPINGS = {"ComfyUIDeployExternalTextList": "External Text List (ComfyUI Deploy)"}
+69 -339
View File
@@ -1,15 +1,10 @@
# credit goes to https://github.com/Kosinkadink/ComfyUI-VideoHelperSuite
# Intended to work with https://github.com/NicholasKao1029/ComfyUI-VideoHelperSuite/tree/main
# credit goes to https://github.com/Kosinkadink/ComfyUI-VideoHelperSuite and is meant to work with
import os
import itertools
import numpy as np
import torch
from typing import Union
from torch import Tensor
import cv2
import psutil
from collections.abc import Mapping
import folder_paths
from comfy.utils import common_upscale
@@ -95,25 +90,13 @@ if gifski_path is None:
gifski_path = shutil.which("gifski")
def is_safe_path(path):
if "VHS_STRICT_PATHS" not in os.environ:
return True
basedir = os.path.abspath(".")
try:
common_path = os.path.commonpath([basedir, path])
except:
# Different drive on windows
return False
return common_path == basedir
def get_sorted_dir_files_from_directory(
directory: str,
skip_first_images: int = 0,
select_every_nth: int = 1,
extensions: Iterable = None,
):
directory = strip_path(directory)
directory = directory.strip()
dir_files = os.listdir(directory)
dir_files = sorted(dir_files)
dir_files = [os.path.join(directory, x) for x in dir_files]
@@ -194,59 +177,18 @@ def requeue_workflow(requeue_required=(-1, True)):
def get_audio(file, start_time=0, duration=0):
args = [ffmpeg_path, "-i", file]
args = [ffmpeg_path, "-v", "error", "-i", file]
if start_time > 0:
args += ["-ss", str(start_time)]
if duration > 0:
args += ["-t", str(duration)]
try:
# TODO: scan for sample rate and maintain
res = subprocess.run(
args + ["-f", "f32le", "-"], capture_output=True, check=True
)
audio = torch.frombuffer(bytearray(res.stdout), dtype=torch.float32)
match = re.search(", (\\d+) Hz, (\\w+), ", res.stderr.decode("utf-8"))
args + ["-f", "wav", "-"], stdout=subprocess.PIPE, check=True
).stdout
except subprocess.CalledProcessError as e:
raise Exception(
f"VHS failed to extract audio from {file}:\n" + e.stderr.decode("utf-8")
)
if match:
ar = int(match.group(1))
# NOTE: Just throwing an error for other channel types right now
# Will deal with issues if they come
ac = {"mono": 1, "stereo": 2}[match.group(2)]
else:
ar = 44100
ac = 2
audio = audio.reshape((-1, ac)).transpose(0, 1).unsqueeze(0)
return {"waveform": audio, "sample_rate": ar}
class LazyAudioMap(Mapping):
def __init__(self, file, start_time, duration):
self.file = file
self.start_time = start_time
self.duration = duration
self._dict = None
def __getitem__(self, key):
if self._dict is None:
self._dict = get_audio(self.file, self.start_time, self.duration)
return self._dict[key]
def __iter__(self):
if self._dict is None:
self._dict = get_audio(self.file, self.start_time, self.duration)
return iter(self._dict)
def __len__(self):
if self._dict is None:
self._dict = get_audio(self.file, self.start_time, self.duration)
return len(self._dict)
def lazy_get_audio(file, start_time=0, duration=0):
return LazyAudioMap(file, start_time, duration)
return False
return res
def lazy_eval(func):
@@ -288,19 +230,6 @@ def validate_sequence(path):
return False
def strip_path(path):
# This leaves whitespace inside quotes and only a single "
# thus ' ""test"' -> '"test'
# consider path.strip(string.whitespace+"\"")
# or weightier re.fullmatch("[\\s\"]*(.+?)[\\s\"]*", path).group(1)
path = path.strip()
if path.startswith('"'):
path = path[1:]
if path.endswith('"'):
path = path[:-1]
return path
def hash_path(path):
if path is None:
return "input"
@@ -357,145 +286,6 @@ def target_size(
return (width, height)
def validate_index(
index: int,
length: int = 0,
is_range: bool = False,
allow_negative=False,
allow_missing=False,
) -> int:
# if part of range, do nothing
if is_range:
return index
# otherwise, validate index
# validate not out of range - only when latent_count is passed in
if length > 0 and index > length - 1 and not allow_missing:
raise IndexError(f"Index '{index}' out of range for {length} item(s).")
# if negative, validate not out of range
if index < 0:
if not allow_negative:
raise IndexError(f"Negative indeces not allowed, but was '{index}'.")
conv_index = length + index
if conv_index < 0 and not allow_missing:
raise IndexError(
f"Index '{index}', converted to '{conv_index}' out of range for {length} item(s)."
)
index = conv_index
return index
def convert_to_index_int(
raw_index: str,
length: int = 0,
is_range: bool = False,
allow_negative=False,
allow_missing=False,
) -> int:
try:
return validate_index(
int(raw_index),
length=length,
is_range=is_range,
allow_negative=allow_negative,
allow_missing=allow_missing,
)
except ValueError as e:
raise ValueError(f"Index '{raw_index}' must be an integer.", e)
def convert_str_to_indexes(
indexes_str: str, length: int = 0, allow_missing=False
) -> list[int]:
if not indexes_str:
return []
int_indexes = list(range(0, length))
allow_negative = length > 0
chosen_indexes = []
# parse string - allow positive ints, negative ints, and ranges separated by ':'
groups = indexes_str.split(",")
groups = [g.strip() for g in groups]
for g in groups:
# parse range of indeces (e.g. 2:16)
if ":" in g:
index_range = g.split(":", 2)
index_range = [r.strip() for r in index_range]
start_index = index_range[0]
if len(start_index) > 0:
start_index = convert_to_index_int(
start_index,
length=length,
is_range=True,
allow_negative=allow_negative,
allow_missing=allow_missing,
)
else:
start_index = 0
end_index = index_range[1]
if len(end_index) > 0:
end_index = convert_to_index_int(
end_index,
length=length,
is_range=True,
allow_negative=allow_negative,
allow_missing=allow_missing,
)
else:
end_index = length
# support step as well, to allow things like reversing, every-other, etc.
step = 1
if len(index_range) > 2:
step = index_range[2]
if len(step) > 0:
step = convert_to_index_int(
step,
length=length,
is_range=True,
allow_negative=True,
allow_missing=True,
)
else:
step = 1
# if latents were passed in, base indeces on known latent count
if len(int_indexes) > 0:
chosen_indexes.extend(int_indexes[start_index:end_index][::step])
# otherwise, assume indeces are valid
else:
chosen_indexes.extend(list(range(start_index, end_index, step)))
# parse individual indeces
else:
chosen_indexes.append(
convert_to_index_int(
g,
length=length,
allow_negative=allow_negative,
allow_missing=allow_missing,
)
)
return chosen_indexes
def select_indexes(input_obj: Union[Tensor, list], idxs: list):
if type(input_obj) == Tensor:
return input_obj[idxs]
else:
return [input_obj[i] for i in idxs]
def select_indexes_from_str(
input_obj: Union[Tensor, list], indexes: str, err_if_missing=True, err_if_empty=True
):
real_idxs = convert_str_to_indexes(
indexes, len(input_obj), allow_missing=not err_if_missing
)
if err_if_empty and len(real_idxs) == 0:
raise Exception(f"Nothing was selected based on indexes found in '{indexes}'.")
return select_indexes(input_obj, real_idxs)
###
def cv_frame_generator(
video,
force_rate,
@@ -505,10 +295,9 @@ def cv_frame_generator(
meta_batch=None,
unique_id=None,
):
video_cap = cv2.VideoCapture(strip_path(video))
video_cap = cv2.VideoCapture(video)
if not video_cap.isOpened():
raise ValueError(f"{video} could not be loaded with cv.")
pbar = None
# extract video metadata
fps = video_cap.get(cv2.CAP_PROP_FPS)
@@ -530,8 +319,6 @@ def cv_frame_generator(
target_frame_time = 1 / force_rate
yield (width, height, fps, duration, total_frames, target_frame_time)
if meta_batch is not None:
yield min(frame_load_cap, total_frames)
time_offset = target_frame_time - base_frame_time
while video_cap.isOpened():
@@ -562,8 +349,7 @@ def cv_frame_generator(
frame = cv2.cvtColor(frame, cv2.COLOR_BGR2RGB)
# convert frame to comfyui's expected format
# TODO: frame contains no exif information. Check if opencv2 has already applied
frame = np.array(frame, dtype=np.float32)
torch.from_numpy(frame).div_(255)
frame = np.array(frame, dtype=np.float32) / 255.0
if prev_frame is not None:
inp = yield prev_frame
if inp is not None:
@@ -571,8 +357,6 @@ def cv_frame_generator(
return
prev_frame = frame
frames_added += 1
if pbar is not None:
pbar.update_absolute(frames_added, frame_load_cap)
# if cap exists and we've reached it, stop processing frames
if frame_load_cap > 0 and frames_added >= frame_load_cap:
break
@@ -583,17 +367,6 @@ def cv_frame_generator(
yield prev_frame
def batched(it, n):
while batch := tuple(itertools.islice(it, n)):
yield batch
def batched_vae_encode(images, vae, frames_per_batch):
for batch in batched(images, frames_per_batch):
image_batch = torch.from_numpy(np.array(batch))
yield from vae.encode(image_batch).numpy()
def load_video_cv(
video: str,
force_rate: int,
@@ -605,8 +378,6 @@ def load_video_cv(
select_every_nth: int,
meta_batch=None,
unique_id=None,
memory_limit_mb=None,
vae=None,
):
if meta_batch is None or unique_id not in meta_batch.inputs:
gen = cv_frame_generator(
@@ -630,89 +401,30 @@ def load_video_cv(
total_frames,
target_frame_time,
)
meta_batch.total_frames = min(meta_batch.total_frames, next(gen))
else:
(gen, width, height, fps, duration, total_frames, target_frame_time) = (
meta_batch.inputs[unique_id]
)
memory_limit = None
if memory_limit_mb is not None:
memory_limit *= 2**20
else:
# TODO: verify if garbage collection should be performed here.
# leaves ~128 MB unreserved for safety
try:
memory_limit = (
psutil.virtual_memory().available + psutil.swap_memory().free
) - 2**27
except:
print(
"Failed to calculate available memory. Memory load limit has been disabled"
)
if memory_limit is not None:
if vae is not None:
# space required to load as f32, exist as latent with wiggle room, decode to f32
max_loadable_frames = int(
memory_limit // (width * height * 3 * (4 + 4 + 1 / 10))
)
else:
# TODO: use better estimate for when vae is not None
# Consider completely ignoring for load_latent case?
max_loadable_frames = int(memory_limit // (width * height * 3 * (0.1)))
if meta_batch is not None:
if meta_batch.frames_per_batch > max_loadable_frames:
raise RuntimeError(
f"Meta Batch set to {meta_batch.frames_per_batch} frames but only {max_loadable_frames} can fit in memory"
)
gen = itertools.islice(gen, meta_batch.frames_per_batch)
else:
original_gen = gen
gen = itertools.islice(gen, max_loadable_frames)
downscale_ratio = getattr(vae, "downscale_ratio", 8)
frames_per_batch = (1920 * 1080 * 16) // (width * height) or 1
if force_size != "Disabled" or vae is not None:
new_size = target_size(
width, height, force_size, custom_width, custom_height, downscale_ratio
)
if new_size[0] != width or new_size[1] != height:
def rescale(frame):
s = torch.from_numpy(
np.fromiter(frame, np.dtype((np.float32, (height, width, 3))))
)
s = s.movedim(-1, 1)
s = common_upscale(s, new_size[0], new_size[1], "lanczos", "center")
return s.movedim(1, -1).numpy()
gen = itertools.chain.from_iterable(
map(rescale, batched(gen, frames_per_batch))
)
else:
new_size = width, height
if vae is not None:
gen = batched_vae_encode(gen, vae, frames_per_batch)
vw, vh = new_size[0] // downscale_ratio, new_size[1] // downscale_ratio
images = torch.from_numpy(np.fromiter(gen, np.dtype((np.float32, (4, vh, vw)))))
else:
# Some minor wizardry to eliminate a copy and reduce max memory by a factor of ~2
images = torch.from_numpy(
np.fromiter(gen, np.dtype((np.float32, (new_size[1], new_size[0], 3))))
np.fromiter(gen, np.dtype((np.float32, (height, width, 3))))
)
if meta_batch is None and memory_limit is not None:
try:
next(original_gen)
raise RuntimeError(
f"Memory limit hit after loading {len(images)} frames. Stopping execution."
)
except StopIteration:
pass
if len(images) == 0:
raise RuntimeError("No frames generated")
if force_size != "Disabled":
new_size = target_size(width, height, force_size, custom_width, custom_height)
if new_size[0] != width or new_size[1] != height:
s = images.movedim(-1, 1)
s = common_upscale(s, new_size[0], new_size[1], "lanczos", "center")
images = s.movedim(1, -1)
# Setup lambda for lazy audio capture
audio = lazy_get_audio(
audio = lambda: get_audio(
video,
skip_first_frames * target_frame_time,
frame_load_cap * target_frame_time * select_every_nth,
@@ -728,16 +440,13 @@ def load_video_cv(
"loaded_fps": 1 / target_frame_time,
"loaded_frame_count": len(images),
"loaded_duration": len(images) * target_frame_time,
"loaded_width": new_size[0],
"loaded_height": new_size[1],
"loaded_width": images.shape[2],
"loaded_height": images.shape[1],
}
if vae is None:
return (images, len(images), audio, video_info, None)
else:
return (None, len(images), audio, video_info, {"samples": images})
return (images, len(images), lazy_eval(audio), video_info)
# modeled after Video upload node
class ComfyUIDeployExternalVideo:
@classmethod
def INPUT_TYPES(s):
@@ -748,46 +457,68 @@ class ComfyUIDeployExternalVideo:
file_parts = f.split(".")
if len(file_parts) > 1 and (file_parts[-1] in video_extensions):
files.append(f)
return {"required": {
return {
"required": {
"input_id": (
"STRING",
{"multiline": False, "default": "input_video"},
),
"force_rate": ("INT", {"default": 0, "min": 0, "max": 60, "step": 1}),
"force_size": (["Disabled", "Custom Height", "Custom Width", "Custom", "256x?", "?x256", "256x256", "512x?", "?x512", "512x512"],),
"custom_width": ("INT", {"default": 512, "min": 0, "max": DIMMAX, "step": 8}),
"custom_height": ("INT", {"default": 512, "min": 0, "max": DIMMAX, "step": 8}),
"frame_load_cap": ("INT", {"default": 0, "min": 0, "max": BIGMAX, "step": 1}),
"skip_first_frames": ("INT", {"default": 0, "min": 0, "max": BIGMAX, "step": 1}),
"select_every_nth": ("INT", {"default": 1, "min": 1, "max": BIGMAX, "step": 1}),
"force_size": (
[
"Disabled",
"Custom Height",
"Custom Width",
"Custom",
"256x?",
"?x256",
"256x256",
"512x?",
"?x512",
"512x512",
],
),
"custom_width": (
"INT",
{"default": 512, "min": 0, "max": DIMMAX, "step": 8},
),
"custom_height": (
"INT",
{"default": 512, "min": 0, "max": DIMMAX, "step": 8},
),
"frame_load_cap": (
"INT",
{"default": 0, "min": 0, "max": BIGMAX, "step": 1},
),
"skip_first_frames": (
"INT",
{"default": 0, "min": 0, "max": BIGMAX, "step": 1},
),
"select_every_nth": (
"INT",
{"default": 1, "min": 1, "max": BIGMAX, "step": 1},
),
},
"optional": {
"meta_batch": ("VHS_BatchManager",),
"vae": ("VAE",),
"default_video": (sorted(files),),
"display_name": (
"STRING",
{"multiline": False, "default": ""},
),
"description": (
"STRING",
{"multiline": True, "default": ""},
),
},
"hidden": {
"unique_id": "UNIQUE_ID"
"default_value": (sorted(files),),
},
"hidden": {"unique_id": "UNIQUE_ID"},
}
CATEGORY = "Video Helper Suite 🎥🅥🅗🅢"
RETURN_TYPES = ("IMAGE", "INT", "AUDIO", "VHS_VIDEOINFO", "LATENT")
RETURN_TYPES = (
"IMAGE",
"INT",
"VHS_AUDIO",
"VHS_VIDEOINFO",
)
RETURN_NAMES = (
"IMAGE",
"frame_count",
"audio",
"video_info",
"LATENT",
)
FUNCTION = "load_video"
@@ -804,6 +535,8 @@ class ComfyUIDeployExternalVideo:
meta_batch = kwargs.get("meta_batch")
unique_id = kwargs.get("unique_id")
video = kwargs.get("default_value")
video_path = folder_paths.get_annotated_filepath(video.strip('"'))
input_dir = folder_paths.get_input_directory()
if input_id.startswith("http"):
@@ -833,11 +566,8 @@ class ComfyUIDeployExternalVideo:
leave=True,
):
out_file.write(chunk)
else:
video = kwargs.get("default_video", "")
if video is None:
raise "No default video given and no external video provided"
video_path = folder_paths.get_annotated_filepath(video.strip('"'))
print("video path: ", video_path)
return load_video_cv(
video=video_path,
+70 -232
View File
@@ -17,120 +17,22 @@ from urllib.parse import quote
import threading
import hashlib
import aiohttp
from aiohttp import ClientSession, web
import aiofiles
from typing import Dict, List, Union, Any, Optional
from PIL import Image
import copy
import struct
from aiohttp import web, ClientSession, ClientError, ClientTimeout
import atexit
# Global session
client_session = None
# def create_client_session():
# global client_session
# if client_session is None:
# client_session = aiohttp.ClientSession()
async def ensure_client_session():
global client_session
if client_session is None:
client_session = aiohttp.ClientSession()
async def cleanup():
global client_session
if client_session:
await client_session.close()
def exit_handler():
print("Exiting the application. Initiating cleanup...")
loop = asyncio.get_event_loop()
loop.run_until_complete(cleanup())
atexit.register(exit_handler)
max_retries = int(os.environ.get('MAX_RETRIES', '5'))
retry_delay_multiplier = float(os.environ.get('RETRY_DELAY_MULTIPLIER', '2'))
print(f"max_retries: {max_retries}, retry_delay_multiplier: {retry_delay_multiplier}")
async def async_request_with_retry(method, url, disable_timeout=False, **kwargs):
global client_session
await ensure_client_session()
# async with aiohttp.ClientSession() as client_session:
retry_delay = 1 # Start with 1 second delay
initial_timeout = 5 # 5 seconds timeout for the initial connection
for attempt in range(max_retries):
try:
# Set a timeout for the initial connection
if not disable_timeout:
timeout = ClientTimeout(total=None, connect=initial_timeout)
kwargs['timeout'] = timeout
async with client_session.request(method, url, **kwargs) as response:
response.raise_for_status()
if method.upper() == 'GET':
await response.read()
return response
except asyncio.TimeoutError:
logger.warning(f"Request timed out after {initial_timeout} seconds (attempt {attempt + 1}/{max_retries})")
except ClientError as e:
if attempt == max_retries - 1:
logger.error(f"Request failed after {max_retries} attempts: {e}")
# raise
logger.warning(f"Request failed (attempt {attempt + 1}/{max_retries}): {e}")
# Wait before retrying
await asyncio.sleep(retry_delay)
retry_delay *= retry_delay_multiplier # Exponential backoff
# If all retries fail, raise an exception
raise Exception(f"Request failed after {max_retries} attempts")
from logging import basicConfig, getLogger
# Check for an environment variable to enable/disable Logfire
use_logfire = os.environ.get('USE_LOGFIRE', 'false').lower() == 'true'
if use_logfire:
try:
import logfire
logfire.configure(
import logfire
# if os.environ.get('LOGFIRE_TOKEN', None) is not None:
logfire.configure(
send_to_logfire="if-token-present"
)
logger = logfire
except ImportError:
print("Logfire not installed or disabled. Using standard Python logger.")
use_logfire = False
if not use_logfire:
# Use a standard Python logger when Logfire is disabled or not available
logger = getLogger("comfy-deploy")
basicConfig(level="INFO") # You can adjust the logging level as needed
def log(level, message, **kwargs):
if use_logfire:
getattr(logger, level)(message, **kwargs)
else:
getattr(logger, level)(f"{message} {kwargs}")
# For a span, you might need to create a context manager
from contextlib import contextmanager
@contextmanager
def log_span(name):
if use_logfire:
with logger.span(name):
yield
else:
yield
# logger.info(f"Start: {name}")
# yield
# logger.info(f"End: {name}")
)
# basicConfig(handlers=[logfire.LogfireLoggingHandler()])
logfire_handler = logfire.LogfireLoggingHandler()
logger = getLogger("comfy-deploy")
logger.addHandler(logfire_handler)
from globals import StreamingPrompt, Status, sockets, SimplePrompt, streaming_prompt_metadata, prompt_metadata
@@ -171,11 +73,11 @@ def clear_current_prompt(sid):
prompt_server = server.PromptServer.instance
to_delete = list(streaming_prompt_metadata[sid].running_prompt_ids) # Convert set to list
logger.info(f"clearing out prompt: {to_delete}")
logger.info("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)
logger.info(f"deleted prompt: {id_to_delete}, remaining tasks: {prompt_server.prompt_queue.get_tasks_remaining()}")
logger.info("deleted prompt: ", id_to_delete, prompt_server.prompt_queue.get_tasks_remaining())
streaming_prompt_metadata[sid].running_prompt_ids.clear()
@@ -234,30 +136,13 @@ def apply_random_seed_to_workflow(workflow_api):
workflow_api (dict): The workflow API dictionary to modify.
"""
for key in workflow_api:
if 'inputs' in workflow_api[key]:
if 'seed' in workflow_api[key]['inputs']:
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)
logger.info(f"Applied random seed {workflow_api[key]['inputs']['seed']} to PromptExpansion")
continue
workflow_api[key]['inputs']['seed'] = randomSeed()
logger.info(f"Applied random seed {workflow_api[key]['inputs']['seed']} to {workflow_api[key]['class_type']}")
if 'noise_seed' in workflow_api[key]['inputs']:
if workflow_api[key]['class_type'] == "RandomNoise":
workflow_api[key]['inputs']['noise_seed'] = randomSeed()
logger.info(f"Applied random noise_seed {workflow_api[key]['inputs']['noise_seed']} to RandomNoise")
continue
if workflow_api[key]['class_type'] == "KSamplerAdvanced":
workflow_api[key]['inputs']['noise_seed'] = randomSeed()
logger.info(f"Applied random noise_seed {workflow_api[key]['inputs']['noise_seed']} to KSamplerAdvanced")
continue
if workflow_api[key]['class_type'] == "SamplerCustom":
workflow_api[key]['inputs']['noise_seed'] = randomSeed()
logger.info(f"Applied random noise_seed {workflow_api[key]['inputs']['noise_seed']} to SamplerCustom")
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
@@ -292,7 +177,7 @@ def apply_inputs_to_workflow(workflow_api: Any, inputs: Any, sid: str = None):
value['inputs']["images"] = new_value
if value["class_type"] == "ComfyUIDeployExternalLora":
value["inputs"]["lora_url"] = new_value
value["inputs"]["default_lora_name"] = new_value
if value["class_type"] == "ComfyUIDeployExternalSlider":
value["inputs"]["default_value"] = new_value
@@ -383,7 +268,7 @@ async def comfy_deploy_run(request):
status = 200
if "node_errors" in res and res["node_errors"] is not None and len(res["node_errors"]) > 0:
if "node_errors" in res and res["node_errors"]:
# Even tho there are node_errors it can still be run
status = 400
await update_run_with_output(prompt_id, {
@@ -421,7 +306,7 @@ async def stream_prompt(data):
workflow_api=workflow_api
)
# log('info', "Begin prompt", prompt=prompt)
logfire.info("Begin prompt", prompt=prompt)
try:
res = post_prompt(prompt)
@@ -444,7 +329,7 @@ async def stream_prompt(data):
status = 200
if "node_errors" in res and res["node_errors"] is not None and len(res["node_errors"]) > 0:
if "node_errors" in res and res["node_errors"]:
# Even tho there are node_errors it can still be run
status = 400
await update_run_with_output(prompt_id, {
@@ -474,8 +359,8 @@ async def stream_response(request):
prompt_id = data.get("prompt_id")
comfy_message_queues[prompt_id] = asyncio.Queue()
with log_span('Streaming Run'):
log('info', 'Streaming prompt')
with logfire.span('Streaming Run'):
logfire.info('Streaming prompt')
try:
result = await stream_prompt(data=data)
@@ -488,7 +373,7 @@ async def stream_response(request):
if not comfy_message_queues[prompt_id].empty():
data = await comfy_message_queues[prompt_id].get()
# log('info', data["event"], data=json.dumps(data))
logfire.info(data["event"], data=json.dumps(data))
# logger.info("listener", data)
await response.write(f"event: event_update\ndata: {json.dumps(data)}\n\n".encode('utf-8'))
await response.drain() # Ensure the buffer is flushed
@@ -499,10 +384,10 @@ async def stream_response(request):
await asyncio.sleep(0.1) # Adjust the sleep duration as needed
except asyncio.CancelledError:
log('info', "Streaming was cancelled")
logfire.info("Streaming was cancelled")
raise
except Exception as e:
log('error', "Streaming error", error=e)
logfire.error("Streaming error", error=e)
finally:
# event_emitter.off("send_json", task)
await response.write_eof()
@@ -597,9 +482,10 @@ async def upload_file_endpoint(request):
if get_url:
try:
async with aiohttp.ClientSession() as session:
headers = {'Authorization': f'Bearer {token}'}
params = {'file_size': file_size, 'type': file_type}
response = await async_request_with_retry('GET', get_url, params=params, headers=headers)
async with session.get(get_url, params=params, headers=headers) as response:
if response.status == 200:
content = await response.json()
upload_url = content["upload_url"]
@@ -610,7 +496,7 @@ async def upload_file_endpoint(request):
# "x-amz-acl": "public-read",
"Content-Length": str(file_size)
}
upload_response = await async_request_with_retry('PUT', upload_url, data=f, headers=headers)
async with session.put(upload_url, data=f, headers=headers) as upload_response:
if upload_response.status == 200:
return web.json_response({
"message": "File uploaded successfully",
@@ -702,7 +588,9 @@ async def update_realtime_run_status(realtime_id: str, status_endpoint: str, sta
if (status_endpoint is None):
return
# requests.post(status_endpoint, json=body)
await async_request_with_retry('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):
@@ -723,8 +611,9 @@ async def websocket_handler(request):
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}'}
response = await async_request_with_retry('GET', get_workflow_endpoint_url, headers=headers)
async with session.get(get_workflow_endpoint_url, headers=headers) as response:
if response.status == 200:
workflow = await response.json()
@@ -856,50 +745,6 @@ async def send(event, data, sid=None):
logger.info(f"Exception: {e}")
traceback.print_exc()
@server.PromptServer.instance.routes.get('/comfydeploy/{tail:.*}')
@server.PromptServer.instance.routes.post('/comfydeploy/{tail:.*}')
async def proxy_to_comfydeploy(request):
# Get the base URL
base_url = f'https://www.comfydeploy.com/{request.match_info["tail"]}'
# Get all query parameters
query_params = request.query_string
# Construct the full target URL with query parameters
target_url = f"{base_url}?{query_params}" if query_params else base_url
# print(f"Proxying request to: {target_url}")
try:
# Create a new ClientSession for each request
async with ClientSession() as client_session:
# Forward the request
client_req = await client_session.request(
method=request.method,
url=target_url,
headers={k: v for k, v in request.headers.items() if k.lower() not in ('host', 'content-length')},
data=await request.read(),
allow_redirects=False,
)
# Read the entire response content
content = await client_req.read()
# Try to decode the content as JSON
try:
json_data = json.loads(content)
# If successful, return a JSON response
return web.json_response(json_data, status=client_req.status)
except json.JSONDecodeError:
# If it's not valid JSON, return the content as-is
return web.Response(body=content, status=client_req.status, headers=client_req.headers)
except ClientError as e:
print(f"Client error occurred while proxying request: {str(e)}")
return web.Response(status=502, text=f"Bad Gateway: {str(e)}")
except Exception as e:
print(f"Error occurred while proxying request: {str(e)}")
return web.Response(status=500, text=f"Internal Server Error: {str(e)}")
prompt_server = server.PromptServer.instance
@@ -960,14 +805,13 @@ async def send_json_override(self, event, data, sid=None):
prompt_metadata[prompt_id].progress.add(node)
calculated_progress = len(prompt_metadata[prompt_id].progress) / len(prompt_metadata[prompt_id].workflow_api)
calculated_progress = round(calculated_progress, 2)
# logger.info("calculated_progress", calculated_progress)
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']
logger.info(f"At: {calculated_progress * 100}% - {class_type}")
logger.info(f"updating run live status {class_type}")
await send("live_status", {
"prompt_id": prompt_id,
"current_node": class_type,
@@ -992,22 +836,16 @@ async def send_json_override(self, event, data, sid=None):
# await update_run_with_output(prompt_id, data)
if event == 'executed' and 'node' in data and 'output' in data:
node_meta = None
logger.info(f"executed {data}")
if prompt_id in prompt_metadata:
node = data.get('node')
class_type = prompt_metadata[prompt_id].workflow_api[node]['class_type']
logger.info(f"Executed {class_type} {data}")
node_meta = {
"node_id": node,
"node_class": class_type,
}
logger.info(f"executed {class_type}")
if class_type == "PreviewImage":
logger.info("Skipping preview image")
logger.info("skipping preview image")
return
else:
logger.info(f"Executed {data}")
await update_run_with_output(prompt_id, data.get('output'), node_id=data.get('node'), node_meta=node_meta)
await update_run_with_output(prompt_id, data.get('output'), node_id=data.get('node'))
# await update_run_with_output(prompt_id, data.get('output'), node_id=data.get('node'))
# update_run_with_output(prompt_id, data.get('output'))
@@ -1026,7 +864,7 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
if (status_endpoint is None):
return
# logger.info(f"progress {calculated_progress}")
logger.info(f"progress {calculated_progress}")
body = {
"run_id": prompt_id,
@@ -1045,7 +883,9 @@ async def update_run_live_status(prompt_id, live_status, calculated_progress: fl
})
# requests.post(status_endpoint, json=body)
await async_request_with_retry('POST', status_endpoint, json=body)
async with aiohttp.ClientSession() as session:
async with session.post(status_endpoint, json=body) as response:
pass
async def update_run(prompt_id: str, status: Status):
@@ -1076,7 +916,9 @@ async def update_run(prompt_id: str, status: Status):
try:
# requests.post(status_endpoint, json=body)
if (status_endpoint is not None):
await async_request_with_retry('POST', status_endpoint, json=body)
async with aiohttp.ClientSession() as session:
async with session.post(status_endpoint, json=body) as response:
pass
if (status_endpoint is not None) and cd_enable_run_log and (status == Status.SUCCESS or status == Status.FAILED):
try:
@@ -1106,7 +948,9 @@ async def update_run(prompt_id: str, status: Status):
]
}
await async_request_with_retry('POST', status_endpoint, json=body)
async with aiohttp.ClientSession() as session:
async with session.post(status_endpoint, json=body) as response:
pass
# requests.post(status_endpoint, json=body)
except Exception as log_error:
logger.info(f"Error reading log file: {log_error}")
@@ -1154,7 +998,7 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
filename = os.path.basename(filename)
file = os.path.join(output_dir, filename)
logger.info(f"Uploading file {file}")
logger.info(f"uploading file {file}")
file_upload_endpoint = prompt_metadata[prompt_id].file_upload_endpoint
@@ -1162,35 +1006,36 @@ async def upload_file(prompt_id, filename, subfolder=None, content_type="image/p
prompt_id = quote(prompt_id)
content_type = quote(content_type)
target_url = f"{file_upload_endpoint}?file_name={filename}&run_id={prompt_id}&type={content_type}&version=v2"
target_url = f"{file_upload_endpoint}?file_name={filename}&run_id={prompt_id}&type={content_type}"
start_time = time.time() # Start timing here
result = await async_request_with_retry("GET", target_url, disable_timeout=True)
result = requests.get(target_url)
end_time = time.time() # End timing after the request is complete
logger.info("Time taken for getting file upload endpoint: {:.2f} seconds".format(end_time - start_time))
ok = await result.json()
ok = result.json()
start_time = time.time() # Start timing here
async with aiofiles.open(file, 'rb') as f:
data = await f.read()
with open(file, 'rb') as f:
data = f.read()
headers = {
# "x-amz-acl": "public-read",
"Content-Type": content_type,
"Content-Length": str(len(data)),
}
# response = requests.put(ok.get("url"), headers=headers, data=data)
response = await async_request_with_retry('PUT', ok.get("url"), headers=headers, data=data)
async with aiohttp.ClientSession() as session:
async with session.put(ok.get("url"), headers=headers, data=data) as response:
logger.info(f"Upload file response status: {response.status}, status text: {response.reason}")
end_time = time.time() # End timing after the request is complete
logger.info("Upload time: {:.2f} seconds".format(end_time - start_time))
def have_pending_upload(prompt_id):
if prompt_id in prompt_metadata and len(prompt_metadata[prompt_id].uploading_nodes) > 0:
logger.info(f"Have pending upload {len(prompt_metadata[prompt_id].uploading_nodes)}")
logger.info(f"have pending upload {len(prompt_metadata[prompt_id].uploading_nodes)}")
return True
logger.info("No pending upload")
logger.info("no pending upload")
return False
def mark_prompt_done(prompt_id):
@@ -1248,7 +1093,7 @@ async def update_file_status(prompt_id: str, data, uploading, have_error=False,
else:
prompt_metadata[prompt_id].uploading_nodes.discard(node_id)
logger.info(f"Remaining uploads: {prompt_metadata[prompt_id].uploading_nodes}")
logger.info(prompt_metadata[prompt_id].uploading_nodes)
# Update the remote status
if have_error:
@@ -1276,10 +1121,8 @@ async def update_file_status(prompt_id: str, data, uploading, have_error=False,
async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, default_content_type: str):
items = data.get(key, [])
upload_tasks = []
for item in items:
# Skipping temp files
# # Skipping temp files
if item.get("type") == "temp":
continue
@@ -1292,35 +1135,29 @@ async def handle_upload(prompt_id: str, data, key: str, content_type_key: str, d
elif file_extension == '.webp':
file_type = 'image/webp'
upload_tasks.append(upload_file(
await upload_file(
prompt_id,
item.get("filename"),
subfolder=item.get("subfolder"),
type=item.get("type"),
content_type=file_type
))
# Execute all upload tasks concurrently
await asyncio.gather(*upload_tasks)
)
# Upload files in the background
async def upload_in_background(prompt_id: str, data, node_id=None, have_upload=True):
try:
upload_tasks = [
handle_upload(prompt_id, data, 'images', "content_type", "image/png"),
handle_upload(prompt_id, data, 'files', "content_type", "image/png"),
handle_upload(prompt_id, data, 'gifs', "format", "image/gif"),
handle_upload(prompt_id, data, 'mesh', "format", "application/octet-stream")
]
await asyncio.gather(*upload_tasks)
await handle_upload(prompt_id, data, 'images', "content_type", "image/png")
await handle_upload(prompt_id, data, 'files', "content_type", "image/png")
# This will also be mp4
await handle_upload(prompt_id, data, 'gifs', "format", "image/gif")
await handle_upload(prompt_id, data, 'mesh', "format", "application/octet-stream")
if have_upload:
await update_file_status(prompt_id, data, False, node_id=node_id)
except Exception as e:
await handle_error(prompt_id, data, e)
async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
async def update_run_with_output(prompt_id, data, node_id=None):
if prompt_id not in prompt_metadata:
return
@@ -1331,8 +1168,7 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
body = {
"run_id": prompt_id,
"output_data": data,
"node_meta": node_meta,
"output_data": data
}
have_upload_media = 'images' in data or 'files' in data or 'gifs' in data or 'mesh' in data
if bypass_upload and have_upload_media:
@@ -1341,7 +1177,7 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
if have_upload_media:
try:
logger.info(f"\nHave_upload {have_upload_media} Node Id: {node_id}")
logger.info(f"\nhave_upload {have_upload} {node_id}")
if have_upload_media:
await update_file_status(prompt_id, data, True, node_id=node_id)
@@ -1354,7 +1190,9 @@ async def update_run_with_output(prompt_id, data, node_id=None, node_meta=None):
# requests.post(status_endpoint, json=body)
if status_endpoint is not None:
await async_request_with_retry('POST', status_endpoint, json=body)
async with aiohttp.ClientSession() as session:
async with session.post(status_endpoint, json=body) as response:
pass
await send('outputs_uploaded', {
"prompt_id": prompt_id
+1 -2
View File
@@ -2,5 +2,4 @@ aiofiles
pydantic
opencv-python
imageio-ffmpeg
brotli
# logfire
logfire
+36 -331
View File
@@ -2,7 +2,6 @@ import { app } from "./app.js";
import { api } from "./api.js";
import { ComfyWidgets, LGraphNode } from "./widgets.js";
import { generateDependencyGraph } from "https://esm.sh/[email protected]";
import { ComfyDeploy } from "https://esm.sh/[email protected]";
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>`;
@@ -20,8 +19,11 @@ function dispatchAPIEventData(data) {
// Custom parse error
if (msg.error) {
let message = msg.error.message;
if (msg.error.details) message += ": " + msg.error.details;
for (const [nodeID, nodeError] of Object.entries(msg.node_errors)) {
if (msg.error.details)
message += ": " + msg.error.details;
for (const [nodeID, nodeError] of Object.entries(
msg.node_errors,
)) {
message += "\n" + nodeError.class_type + ":";
for (const errorReason of nodeError.errors) {
message +=
@@ -207,26 +209,14 @@ const ext = {
ComfyWidgets.STRING(
this,
"workflow_name",
[
"",
{
default: this.properties.workflow_name,
multiline: false,
},
],
["", { default: this.properties.workflow_name, multiline: false }],
app,
);
ComfyWidgets.STRING(
this,
"workflow_id",
[
"",
{
default: this.properties.workflow_id,
multiline: false,
},
],
["", { default: this.properties.workflow_id, multiline: false }],
app,
);
@@ -271,103 +261,26 @@ const ext = {
// const graphCanvas = document.getElementById("graph-canvas");
window.addEventListener("message", async (event) => {
// console.log("message", event);
try {
const message = JSON.parse(event.data);
if (message.type === "graph_load") {
const comfyUIWorkflow = message.data;
// console.log("recieved: ", comfyUIWorkflow);
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) {
console.log("loadGraphData");
app.loadGraphData(comfyUIWorkflow);
}
} else if (message.type === "deploy") {
// deployWorkflow();
const prompt = await app.graphToPrompt();
// api.handlePromptGenerated(prompt);
sendEventToCD("cd_plugin_onDeployChanges", prompt);
} else if (message.type === "queue_prompt") {
const prompt = await app.graphToPrompt();
if (typeof api.handlePromptGenerated === "function") {
api.handlePromptGenerated(prompt);
} else {
console.warn("api.handlePromptGenerated is not a function");
}
sendEventToCD("cd_plugin_onQueuePrompt", prompt);
} else if (message.type === "get_prompt") {
const prompt = await app.graphToPrompt();
sendEventToCD("cd_plugin_onGetPrompt", prompt);
} else if (message.type === "event") {
dispatchAPIEventData(message.data);
} else if (message.type === "add_node") {
console.log("add node", message.data);
app.graph.beforeChange();
var node = LiteGraph.createNode(message.data.type);
node.configure({
widgets_values: message.data.widgets_values,
});
console.log("node", node);
const graphMouse = app.canvas.graph_mouse;
node.pos = [graphMouse[0], graphMouse[1]];
app.graph.add(node);
app.graph.afterChange();
} else if (message.type === "zoom_to_node") {
const nodeId = message.data.nodeId;
const position = message.data.position;
const node = app.graph.getNodeById(nodeId);
if (!node) return;
const canvas = app.canvas;
const targetScale = 1;
const targetOffsetX =
canvas.canvas.width / 4 - position[0] - node.size[0] / 2;
const targetOffsetY =
canvas.canvas.height / 4 - position[1] - node.size[1] / 2;
const startScale = canvas.ds.scale;
const startOffsetX = canvas.ds.offset[0];
const startOffsetY = canvas.ds.offset[1];
const duration = 400; // Animation duration in milliseconds
const startTime = Date.now();
function easeOutCubic(t) {
return 1 - Math.pow(1 - t, 3);
}
function lerp(start, end, t) {
return start * (1 - t) + end * t;
}
function animate() {
const currentTime = Date.now();
const elapsedTime = currentTime - startTime;
const t = Math.min(elapsedTime / duration, 1);
const easedT = easeOutCubic(t);
const currentScale = lerp(startScale, targetScale, easedT);
const currentOffsetX = lerp(startOffsetX, targetOffsetX, easedT);
const currentOffsetY = lerp(startOffsetY, targetOffsetY, easedT);
canvas.setZoom(currentScale);
canvas.ds.offset = [currentOffsetX, currentOffsetY];
canvas.draw(true, true);
if (t < 1) {
requestAnimationFrame(animate);
}
}
animate();
}
// else if (message.type === "refresh") {
// sendEventToCD("cd_plugin_onRefresh");
@@ -375,6 +288,10 @@ const ext = {
} catch (error) {
// console.error("Error processing message:", error);
}
// if (!event.data.flow || Object.entries(event.data.flow).length <= 0)
// return;
// updateBlendshapesPrompts(event.data.flow);
});
api.addEventListener("executed", (evt) => {
@@ -498,7 +415,6 @@ function createDynamicUIHtml(data) {
return html;
}
// Modify the existing deployWorkflow function
async function deployWorkflow() {
const deploy = document.getElementById("deploy-button");
@@ -645,30 +561,30 @@ async function deployWorkflow() {
console.log(hash);
return hash.file_hash;
},
// handleFileUpload: async (file, hash, prevhash) => {
// console.log("Uploading ", file);
// loadingDialog.showLoading("Uploading file", file);
// try {
// const { download_url } = await fetch(`/comfyui-deploy/upload-file`, {
// method: "POST",
// body: JSON.stringify({
// file_path: file,
// token: apiKey,
// url: endpoint + "/api/upload-url",
// }),
// })
// .then((x) => x.json())
// .catch(() => {
// loadingDialog.close();
// confirmDialog.confirm("Error", "Unable to upload file " + file);
// });
// loadingDialog.showLoading("Uploaded file", file);
// console.log(download_url);
// return download_url;
// } catch (error) {
// return undefined;
// }
// },
handleFileUpload: async (file, hash, prevhash) => {
console.log("Uploading ", file);
loadingDialog.showLoading("Uploading file", file);
try {
const { download_url } = await fetch(`/comfyui-deploy/upload-file`, {
method: "POST",
body: JSON.stringify({
file_path: file,
token: apiKey,
url: endpoint + "/api/upload-url",
}),
})
.then((x) => x.json())
.catch(() => {
loadingDialog.close();
confirmDialog.confirm("Error", "Unable to upload file " + file);
});
loadingDialog.showLoading("Uploaded file", file);
console.log(download_url);
return download_url;
} catch (error) {
return undefined;
}
},
existingDependencies: existing_workflow.dependencies,
});
@@ -693,15 +609,6 @@ async function deployWorkflow() {
"Check dependencies",
// JSON.stringify(deps, null, 2),
`
<div>
You will need to create a cloud machine with the following configuration on ComfyDeploy
<ol style="text-align: left; margin-top: 10px;">
<li>Review the dependencies listed in the graph below</li>
<li>Create a new cloud machine with the required configuration</li>
<li>Install missing models and check missing files</li>
<li>Deploy your workflow to the newly created machine</li>
</ol>
</div>
<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;"
@@ -771,14 +678,6 @@ async function deployWorkflow() {
`<span style="color:green;">Deployed successfully!</span> <a style="color:white;" target="_blank" href=${endpoint}/workflows/${data.workflow_id}>-> View here</a> <br/> <br/> Workflow ID: ${data.workflow_id} <br/> Workflow Name: ${workflow_name} <br/> Workflow Version: ${data.version} <br/>`,
);
// // Refresh the workflows list in the sidebar
// const sidebarEl = document.querySelector(
// '.comfy-sidebar-tab[data-id="search"]',
// );
// if (sidebarEl) {
// refreshWorkflowsList(sidebarEl);
// }
setTimeout(() => {
title.textContent = "Deploy";
title.style.color = "white";
@@ -796,85 +695,6 @@ async function deployWorkflow() {
}
}
// Add this function to refresh the workflows list
function refreshWorkflowsList(el) {
const workflowsList = el.querySelector("#workflows-list");
const workflowsLoading = el.querySelector("#workflows-loading");
workflowsLoading.style.display = "flex";
workflowsList.style.display = "none";
workflowsList.innerHTML = "";
client.workflows
.getAll({
page: "1",
pageSize: "10",
})
.then((result) => {
workflowsLoading.style.display = "none";
workflowsList.style.display = "block";
if (result.length === 0) {
workflowsList.innerHTML =
"<li style='color: #bdbdbd;'>No workflows found</li>";
return;
}
result.forEach((workflow) => {
const li = document.createElement("li");
li.style.marginBottom = "15px";
li.style.padding = "15px";
li.style.backgroundColor = "#2a2a2a";
li.style.borderRadius = "8px";
li.style.boxShadow = "0 2px 4px rgba(0,0,0,0.1)";
const lastRun = workflow.runs[0];
const lastRunStatus = lastRun ? lastRun.status : "No runs";
const statusColor =
lastRunStatus === "success"
? "#4CAF50"
: lastRunStatus === "error"
? "#F44336"
: "#FFC107";
const timeAgo = getTimeAgo(new Date(workflow.updatedAt));
li.innerHTML = `
<div style="display: flex; justify-content: space-between; align-items: center; margin-bottom: 10px;">
<div style="flex: 1; overflow: hidden; text-overflow: ellipsis; white-space: nowrap;">
<strong style="font-size: 18px; color: #e0e0e0;">${workflow.name}</strong>
</div>
<span style="font-size: 12px; color: ${statusColor}; margin-left: 10px;">Last run: ${lastRunStatus}</span>
</div>
<div style="font-size: 14px; color: #bdbdbd; margin-bottom: 10px;">Last updated ${timeAgo}</div>
<div style="display: flex; gap: 10px;">
<button class="open-cloud-btn" style="padding: 5px 10px; background-color: #4CAF50; color: white; border: none; border-radius: 4px; cursor: pointer;">Open in Cloud</button>
<button class="load-api-btn" style="padding: 5px 10px; background-color: #2196F3; color: white; border: none; border-radius: 4px; cursor: pointer;">Load Workflow</button>
</div>
`;
const openCloudBtn = li.querySelector(".open-cloud-btn");
openCloudBtn.onclick = () =>
window.open(
`${getData().endpoint}/workflows/${workflow.id}?workspace=true`,
"_blank",
);
const loadApiBtn = li.querySelector(".load-api-btn");
loadApiBtn.onclick = () => loadWorkflowApi(workflow.versions[0].id);
workflowsList.appendChild(li);
});
})
.catch((error) => {
console.error("Error fetching workflows:", error);
workflowsLoading.style.display = "none";
workflowsList.style.display = "block";
workflowsList.innerHTML =
"<li style='color: #F44336;'>Error fetching workflows</li>";
});
}
function addButton() {
const menu = document.querySelector(".comfy-menu");
@@ -1367,118 +1187,3 @@ export class ConfigDialog extends ComfyDialog {
}
export const configDialog = new ConfigDialog();
const currentOrigin = window.location.origin;
const client = new ComfyDeploy({
bearerAuth: getData().apiKey,
serverURL: `${currentOrigin}/comfydeploy/api/`,
});
app.extensionManager.registerSidebarTab({
id: "search",
icon: "pi pi-cloud-upload",
title: "Deploy",
tooltip: "Deploy and Configure",
type: "custom",
render: (el) => {
el.innerHTML = `
<div style="padding: 20px;">
<h3>Comfy Deploy</h3>
<div id="deploy-container" style="margin-bottom: 20px;"></div>
<div id="workflows-container">
<h4>Your Workflows</h4>
<div id="workflows-loading" style="display: flex; justify-content: center; align-items: center; height: 100px;">
${loadingIcon}
</div>
<ul id="workflows-list" style="list-style-type: none; padding: 0; display: none;"></ul>
</div>
<div id="config-container"></div>
</div>
`;
// Add deploy button
const deployContainer = el.querySelector("#deploy-container");
const deployButton = document.createElement("button");
deployButton.id = "sidebar-deploy-button";
deployButton.style.display = "flex";
deployButton.style.alignItems = "center";
deployButton.style.justifyContent = "center";
deployButton.style.width = "100%";
deployButton.style.marginBottom = "10px";
deployButton.style.padding = "10px";
deployButton.style.fontSize = "16px";
deployButton.style.fontWeight = "bold";
deployButton.style.backgroundColor = "#4CAF50";
deployButton.style.color = "white";
deployButton.style.border = "none";
deployButton.style.borderRadius = "5px";
deployButton.style.cursor = "pointer";
deployButton.innerHTML = `<i class="pi pi-cloud-upload" style="margin-right: 8px;"></i><div id='sidebar-button-title'>Deploy</div>`;
deployButton.onclick = async () => {
await deployWorkflow();
// Refresh the workflows list after deployment
refreshWorkflowsList(el);
};
deployContainer.appendChild(deployButton);
// Add config button
const configContainer = el.querySelector("#config-container");
const configButton = document.createElement("button");
configButton.style.display = "flex";
configButton.style.alignItems = "center";
configButton.style.justifyContent = "center";
configButton.style.width = "100%";
configButton.style.padding = "8px";
configButton.style.fontSize = "14px";
configButton.style.backgroundColor = "#f0f0f0";
configButton.style.color = "#333";
configButton.style.border = "1px solid #ccc";
configButton.style.borderRadius = "5px";
configButton.style.cursor = "pointer";
configButton.innerHTML = `<i class="pi pi-cog" style="margin-right: 8px;"></i>Configure`;
configButton.onclick = () => {
configDialog.show();
};
deployContainer.appendChild(configButton);
// Fetch and display workflows
const workflowsList = el.querySelector("#workflows-list");
const workflowsLoading = el.querySelector("#workflows-loading");
refreshWorkflowsList(el);
},
});
function getTimeAgo(date) {
const seconds = Math.floor((new Date() - date) / 1000);
let interval = seconds / 31536000;
if (interval > 1) return Math.floor(interval) + " years ago";
interval = seconds / 2592000;
if (interval > 1) return Math.floor(interval) + " months ago";
interval = seconds / 86400;
if (interval > 1) return Math.floor(interval) + " days ago";
interval = seconds / 3600;
if (interval > 1) return Math.floor(interval) + " hours ago";
interval = seconds / 60;
if (interval > 1) return Math.floor(interval) + " minutes ago";
return Math.floor(seconds) + " seconds ago";
}
async function loadWorkflowApi(versionId) {
try {
const response = await client.comfyui.getWorkflowVersionVersionId({
versionId: versionId,
});
// Implement the logic to load the workflow API into the ComfyUI interface
console.log("Workflow API loaded:", response);
await window["app"].ui.settings.setSettingValueAsync(
"Comfy.Validation.Workflows",
false,
);
app.loadGraphData(response.workflow);
// You might want to update the UI or trigger some action in ComfyUI here
} catch (error) {
console.error("Error loading workflow API:", error);
// Show an error message to the user
}
}
+1 -1
View File
@@ -74,7 +74,7 @@
"mitata": "^0.1.6",
"ms": "^2.1.3",
"nanoid": "^5.0.4",
"next": "14.2",
"next": "14.1",
"next-plausible": "^3.12.0",
"next-themes": "^0.2.1",
"next-usequerystate": "^1.13.2",
+1 -3
View File
@@ -51,9 +51,7 @@ const createRunRoute = createRoute({
export const registerCreateRunRoute = (app: App) => {
app.openapi(createRunRoute, async (c) => {
const data = c.req.valid("json");
const proto = c.req.headers.get('x-forwarded-proto') || "http";
const host = c.req.headers.get('x-forwarded-host') || c.req.headers.get('host');
const origin = `${proto}://${host}` || new URL(c.req.url).origin;
const origin = new URL(c.req.url).origin;
const apiKeyTokenData = c.get("apiKeyTokenData")!;
const { deployment_id, inputs } = data;
+1 -1
View File
@@ -102,7 +102,7 @@ export const createRun = withServerPromise(
let prompt_id: string | undefined = undefined;
const shareData = {
workflow_api_raw: workflow_api,
workflow_api: workflow_api,
status_endpoint: `${origin}/api/update-run`,
file_upload_endpoint: `${origin}/api/file-upload`,
};