Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
03c4e2d85d | ||
|
|
bad8b1104f | ||
|
|
ac8779dc54 | ||
|
|
2a5223f2f0 | ||
|
|
f02aef4cb3 |
+128
-86
@@ -265,7 +265,7 @@ def clear_current_prompt(sid):
|
||||
streaming_prompt_metadata[sid].running_prompt_ids.clear()
|
||||
|
||||
|
||||
def post_prompt(json_data):
|
||||
async def post_prompt(json_data):
|
||||
prompt_server = server.PromptServer.instance
|
||||
json_data = prompt_server.trigger_on_prompt(json_data)
|
||||
|
||||
@@ -281,7 +281,14 @@ def post_prompt(json_data):
|
||||
|
||||
if "prompt" in json_data:
|
||||
prompt = json_data["prompt"]
|
||||
valid = execution.validate_prompt(prompt)
|
||||
prompt_id = json_data.get("prompt_id") or str(uuid.uuid4())
|
||||
|
||||
try:
|
||||
valid = await execution.validate_prompt(prompt_id, prompt)
|
||||
except TypeError as e:
|
||||
logger.warning(f"Trying old validate_prompt signature: {e}")
|
||||
valid = execution.validate_prompt(prompt)
|
||||
|
||||
extra_data = {}
|
||||
if "extra_data" in json_data:
|
||||
extra_data = json_data["extra_data"]
|
||||
@@ -292,8 +299,6 @@ def post_prompt(json_data):
|
||||
if "client_id" in json_data:
|
||||
extra_data["client_id"] = json_data["client_id"]
|
||||
if valid[0]:
|
||||
# if the prompt id is provided
|
||||
prompt_id = json_data.get("prompt_id") or str(uuid.uuid4())
|
||||
outputs_to_execute = valid[2]
|
||||
prompt_server.prompt_queue.put(
|
||||
(number, prompt_id, prompt, extra_data, outputs_to_execute)
|
||||
@@ -500,15 +505,15 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
|
||||
|
||||
prompt_id = str(uuid.uuid4())
|
||||
|
||||
prompt = {
|
||||
"prompt": workflow_api,
|
||||
"client_id": sid, # "comfy_deploy_instance", #api.client_id
|
||||
"prompt_id": prompt_id,
|
||||
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
|
||||
}
|
||||
# prompt = {
|
||||
# "prompt": workflow_api,
|
||||
# "client_id": sid, # "comfy_deploy_instance", #api.client_id
|
||||
# "prompt_id": prompt_id,
|
||||
# "extra_data": {"extra_pnginfo": {"workflow": workflow}},
|
||||
# }
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
# res = post_prompt(prompt)
|
||||
inputs.running_prompt_ids.add(prompt_id)
|
||||
prompt_metadata[prompt_id] = SimplePrompt(
|
||||
status_endpoint=inputs.status_endpoint,
|
||||
@@ -519,7 +524,7 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
|
||||
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()
|
||||
# stack_trace = traceback.format_exc().strip()
|
||||
logger.info(f"error: {error_type}, {e}")
|
||||
logger.info(f"stack trace: {stack_trace_short}")
|
||||
|
||||
@@ -595,7 +600,7 @@ async def comfy_deploy_run(request):
|
||||
)
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
res = await post_prompt(prompt)
|
||||
except Exception as e:
|
||||
error_type = type(e).__name__
|
||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||
@@ -664,7 +669,7 @@ async def stream_prompt(data, token):
|
||||
# log('info', "Begin prompt", prompt=prompt)
|
||||
|
||||
try:
|
||||
res = post_prompt(prompt)
|
||||
res = await post_prompt(prompt)
|
||||
except Exception as e:
|
||||
error_type = type(e).__name__
|
||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||
@@ -1266,22 +1271,11 @@ def handle_execute(class_type, last_node_id, prompt_id, server, unique_id):
|
||||
|
||||
try:
|
||||
origin_execute = execution.execute
|
||||
is_async = asyncio.iscoroutinefunction(origin_execute)
|
||||
|
||||
def swizzle_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
):
|
||||
unique_id = current_item
|
||||
class_type = dynprompt.get_node(unique_id)["class_type"]
|
||||
last_node_id = server.last_node_id
|
||||
result = origin_execute(
|
||||
if is_async:
|
||||
|
||||
async def swizzle_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
@@ -1291,12 +1285,61 @@ try:
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
)
|
||||
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||
return result
|
||||
pending_async_nodes,
|
||||
):
|
||||
unique_id = current_item
|
||||
class_type = dynprompt.get_node(unique_id)["class_type"]
|
||||
last_node_id = server.last_node_id
|
||||
|
||||
result = await origin_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
pending_async_nodes,
|
||||
)
|
||||
|
||||
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||
return result
|
||||
else:
|
||||
|
||||
def swizzle_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
):
|
||||
unique_id = current_item
|
||||
class_type = dynprompt.get_node(unique_id)["class_type"]
|
||||
last_node_id = server.last_node_id
|
||||
|
||||
result = origin_execute(
|
||||
server,
|
||||
dynprompt,
|
||||
caches,
|
||||
current_item,
|
||||
extra_data,
|
||||
executed,
|
||||
prompt_id,
|
||||
execution_list,
|
||||
pending_subgraph_results,
|
||||
)
|
||||
|
||||
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||
return result
|
||||
|
||||
execution.execute = swizzle_execute
|
||||
except Exception as e:
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@@ -1440,10 +1483,6 @@ async def send_json_override(self, event, data, sid=None):
|
||||
logger.info(format_table(headers, table_data))
|
||||
# print("========================\n")
|
||||
|
||||
timeline = format_execution_timeline(NODE_EXECUTION_TIMES)
|
||||
logger.info(f"\nNode Execution Timeline:\n{timeline}")
|
||||
# Clear the execution times for the next run
|
||||
|
||||
# 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)
|
||||
@@ -2059,6 +2098,7 @@ async def upload_in_background(
|
||||
("model_file", "format", "application/octet-stream"),
|
||||
("result", "format", "application/octet-stream"),
|
||||
("text_file", "format", "text/plain"),
|
||||
("audio", "format", "audio/mpeg"),
|
||||
]:
|
||||
items = data.get(file_type, [])
|
||||
|
||||
@@ -2146,6 +2186,7 @@ async def update_run_with_output(
|
||||
or "model_file" in data
|
||||
or "result" in data
|
||||
or "text_file" in data
|
||||
or "audio" in data
|
||||
)
|
||||
if bypass_upload and have_upload_media:
|
||||
print(
|
||||
@@ -2754,59 +2795,57 @@ class UploadQueue:
|
||||
logger.error(f"Upload failed: {str(e)}")
|
||||
logger.error(traceback.format_exc())
|
||||
finally:
|
||||
# Remove this upload from tracking
|
||||
if prompt_id in self.pending_uploads:
|
||||
self.pending_uploads[prompt_id].discard(upload_id)
|
||||
# Remove from node tracking if applicable
|
||||
if (
|
||||
node_id
|
||||
and prompt_id in self.node_uploads
|
||||
and node_id in self.node_uploads[prompt_id]
|
||||
):
|
||||
self.node_uploads[prompt_id][node_id].discard(upload_id)
|
||||
async with self.lock: # Acquire lock to protect shared dict access
|
||||
if prompt_id in self.pending_uploads:
|
||||
self.pending_uploads[prompt_id].discard(upload_id)
|
||||
|
||||
# If this was the last upload for this node, clean up node data
|
||||
if not self.node_uploads[prompt_id][node_id]:
|
||||
del self.node_uploads[prompt_id][node_id]
|
||||
if prompt_id in self.node_output_data:
|
||||
if node_id in self.node_output_data[prompt_id]:
|
||||
if self.node_output_data[prompt_id][node_id][
|
||||
"data"
|
||||
]:
|
||||
# Send final node data to API before cleanup
|
||||
if prompt_metadata[
|
||||
prompt_id
|
||||
].status_endpoint:
|
||||
body = {
|
||||
"run_id": prompt_id,
|
||||
"output_data": self.node_output_data[
|
||||
if (
|
||||
node_id
|
||||
and prompt_id in self.node_uploads
|
||||
and node_id in self.node_uploads[prompt_id]
|
||||
):
|
||||
self.node_uploads[prompt_id][node_id].discard(upload_id)
|
||||
|
||||
if not self.node_uploads[prompt_id][node_id]:
|
||||
del self.node_uploads[prompt_id][node_id]
|
||||
|
||||
if (
|
||||
prompt_id in self.node_output_data
|
||||
and node_id in self.node_output_data[prompt_id]
|
||||
):
|
||||
node_data = self.node_output_data[prompt_id][
|
||||
node_id
|
||||
]
|
||||
if node_data["data"]:
|
||||
body = {
|
||||
"run_id": prompt_id,
|
||||
"output_data": node_data["data"],
|
||||
"node_meta": {"node_id": node_id},
|
||||
}
|
||||
try:
|
||||
await async_request_with_retry(
|
||||
"POST",
|
||||
prompt_metadata[
|
||||
prompt_id
|
||||
][node_id]["data"],
|
||||
"node_meta": {"node_id": node_id},
|
||||
}
|
||||
try:
|
||||
await async_request_with_retry(
|
||||
"POST",
|
||||
prompt_metadata[
|
||||
prompt_id
|
||||
].status_endpoint,
|
||||
token=prompt_metadata[
|
||||
prompt_id
|
||||
].token,
|
||||
json=body,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"Failed to send final node data: {str(e)}"
|
||||
)
|
||||
].status_endpoint,
|
||||
token=prompt_metadata[
|
||||
prompt_id
|
||||
].token,
|
||||
json=body,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.error(
|
||||
f"Failed to send final node data: {str(e)}"
|
||||
)
|
||||
|
||||
# Safe to delete now (re-check not strictly needed with lock, but harmless)
|
||||
del self.node_output_data[prompt_id][node_id]
|
||||
|
||||
# Send status update
|
||||
await self.update_queue_status(prompt_id)
|
||||
|
||||
# If no more pending uploads for this prompt and it's done, update status
|
||||
if not self.pending_uploads[prompt_id] and is_prompt_done(
|
||||
prompt_id
|
||||
if (
|
||||
prompt_id in self.pending_uploads
|
||||
and not self.pending_uploads[prompt_id]
|
||||
and is_prompt_done(prompt_id)
|
||||
):
|
||||
# Clean up all data for this prompt
|
||||
if prompt_id in self.node_uploads:
|
||||
@@ -2819,9 +2858,12 @@ class UploadQueue:
|
||||
loop.create_task(update_run(prompt_id, Status.SUCCESS))
|
||||
loop.create_task(send("success", {"prompt_id": prompt_id}))
|
||||
|
||||
# Mark task as done
|
||||
# Mark task as done (outside lock to avoid holding it unnecessarily)
|
||||
self.queue.task_done()
|
||||
|
||||
# Send status update (also outside lock)
|
||||
await self.update_queue_status(prompt_id)
|
||||
|
||||
except Exception as e:
|
||||
logger.error(f"Error in upload worker: {str(e)}")
|
||||
logger.error(traceback.format_exc())
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
[project]
|
||||
name = "comfyui-deploy"
|
||||
description = "Open source comfyui deployment platform, a vercel for generative workflow infra."
|
||||
version = "2.1.0"
|
||||
version = "2.2.1"
|
||||
license = { file = "LICENSE" }
|
||||
dependencies = ["aiofiles", "pydantic", "opencv-python", "imageio-ffmpeg", "tabulate", "brotli"]
|
||||
|
||||
|
||||
+2
-5
@@ -2000,13 +2000,10 @@ const currentOrigin = window.location.origin;
|
||||
// serverURL: `${currentOrigin}/comfydeploy/api/`,
|
||||
// });
|
||||
|
||||
// Check if the current URL hostname starts with localhost or 127.0.0.1
|
||||
const isLocalhost =
|
||||
window.location.hostname === "localhost" ||
|
||||
window.location.hostname === "127.0.0.1";
|
||||
const isComfyDeployDashboard = currentOrigin.includes("comfydeploy.com");
|
||||
|
||||
// Only register the sidebar tab if we're on localhost
|
||||
if (isLocalhost) {
|
||||
if (!isComfyDeployDashboard) {
|
||||
app.extensionManager.registerSidebarTab({
|
||||
id: "search",
|
||||
icon: "pi pi-cloud-upload",
|
||||
|
||||
Reference in New Issue
Block a user