Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
94210021e3 | ||
|
|
6d2b918ef1 | ||
|
|
d8b1e88c03 |
@@ -92,7 +92,7 @@ class ComfyDeployOutputText:
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
return {"ui": {"text_file": results}}
|
return {"ui": {"text_file": results, "text": [text]}}
|
||||||
|
|
||||||
|
|
||||||
NODE_CLASS_MAPPINGS = {"ComfyDeployOutputText": ComfyDeployOutputText}
|
NODE_CLASS_MAPPINGS = {"ComfyDeployOutputText": ComfyDeployOutputText}
|
||||||
|
|||||||
+86
-308
@@ -265,7 +265,7 @@ def clear_current_prompt(sid):
|
|||||||
streaming_prompt_metadata[sid].running_prompt_ids.clear()
|
streaming_prompt_metadata[sid].running_prompt_ids.clear()
|
||||||
|
|
||||||
|
|
||||||
async def post_prompt(json_data):
|
def post_prompt(json_data):
|
||||||
prompt_server = server.PromptServer.instance
|
prompt_server = server.PromptServer.instance
|
||||||
json_data = prompt_server.trigger_on_prompt(json_data)
|
json_data = prompt_server.trigger_on_prompt(json_data)
|
||||||
|
|
||||||
@@ -281,48 +281,7 @@ async def post_prompt(json_data):
|
|||||||
|
|
||||||
if "prompt" in json_data:
|
if "prompt" in json_data:
|
||||||
prompt = json_data["prompt"]
|
prompt = json_data["prompt"]
|
||||||
prompt_id = json_data.get("prompt_id") or str(uuid.uuid4())
|
valid = execution.validate_prompt(prompt)
|
||||||
|
|
||||||
partial_execution_targets = None
|
|
||||||
if "partial_execution_targets" in json_data:
|
|
||||||
partial_execution_targets = json_data["partial_execution_targets"]
|
|
||||||
|
|
||||||
# Handle different validate_prompt signatures (newest to oldest)
|
|
||||||
valid = None
|
|
||||||
last_error = None
|
|
||||||
|
|
||||||
# v0.3.48 (3 args)
|
|
||||||
try:
|
|
||||||
valid = await execution.validate_prompt(
|
|
||||||
prompt_id, prompt, partial_execution_targets
|
|
||||||
)
|
|
||||||
except TypeError as e:
|
|
||||||
last_error = e
|
|
||||||
logger.debug(
|
|
||||||
f"validate_prompt with 3 params not supported, trying with 2. Debug: {last_error}"
|
|
||||||
)
|
|
||||||
|
|
||||||
# v0.3.45 - 0.3.47 (2 args)
|
|
||||||
if valid is None:
|
|
||||||
try:
|
|
||||||
valid = await execution.validate_prompt(prompt_id, prompt)
|
|
||||||
except TypeError as e:
|
|
||||||
last_error = e
|
|
||||||
logger.debug(
|
|
||||||
f"validate_prompt with 2 params not supported, trying legacy signature. Debug: {last_error}"
|
|
||||||
)
|
|
||||||
|
|
||||||
# v0.3.44 or older (1 arg)
|
|
||||||
if valid is None:
|
|
||||||
try:
|
|
||||||
valid = execution.validate_prompt(prompt)
|
|
||||||
except TypeError as e:
|
|
||||||
last_error = e
|
|
||||||
logger.error(
|
|
||||||
f"validate_prompt failed with all signatures. Last error: {last_error}"
|
|
||||||
)
|
|
||||||
raise
|
|
||||||
|
|
||||||
extra_data = {}
|
extra_data = {}
|
||||||
if "extra_data" in json_data:
|
if "extra_data" in json_data:
|
||||||
extra_data = json_data["extra_data"]
|
extra_data = json_data["extra_data"]
|
||||||
@@ -333,6 +292,8 @@ async def post_prompt(json_data):
|
|||||||
if "client_id" in json_data:
|
if "client_id" in json_data:
|
||||||
extra_data["client_id"] = json_data["client_id"]
|
extra_data["client_id"] = json_data["client_id"]
|
||||||
if valid[0]:
|
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]
|
outputs_to_execute = valid[2]
|
||||||
prompt_server.prompt_queue.put(
|
prompt_server.prompt_queue.put(
|
||||||
(number, prompt_id, prompt, extra_data, outputs_to_execute)
|
(number, prompt_id, prompt, extra_data, outputs_to_execute)
|
||||||
@@ -539,15 +500,15 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
|
|||||||
|
|
||||||
prompt_id = str(uuid.uuid4())
|
prompt_id = str(uuid.uuid4())
|
||||||
|
|
||||||
# prompt = {
|
prompt = {
|
||||||
# "prompt": workflow_api,
|
"prompt": workflow_api,
|
||||||
# "client_id": sid, # "comfy_deploy_instance", #api.client_id
|
"client_id": sid, # "comfy_deploy_instance", #api.client_id
|
||||||
# "prompt_id": prompt_id,
|
"prompt_id": prompt_id,
|
||||||
# "extra_data": {"extra_pnginfo": {"workflow": workflow}},
|
"extra_data": {"extra_pnginfo": {"workflow": workflow}},
|
||||||
# }
|
}
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# res = post_prompt(prompt)
|
res = post_prompt(prompt)
|
||||||
inputs.running_prompt_ids.add(prompt_id)
|
inputs.running_prompt_ids.add(prompt_id)
|
||||||
prompt_metadata[prompt_id] = SimplePrompt(
|
prompt_metadata[prompt_id] = SimplePrompt(
|
||||||
status_endpoint=inputs.status_endpoint,
|
status_endpoint=inputs.status_endpoint,
|
||||||
@@ -558,7 +519,7 @@ def send_prompt(sid: str, inputs: StreamingPrompt):
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_type = type(e).__name__
|
error_type = type(e).__name__
|
||||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
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"error: {error_type}, {e}")
|
||||||
logger.info(f"stack trace: {stack_trace_short}")
|
logger.info(f"stack trace: {stack_trace_short}")
|
||||||
|
|
||||||
@@ -634,7 +595,7 @@ async def comfy_deploy_run(request):
|
|||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
res = await post_prompt(prompt)
|
res = post_prompt(prompt)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_type = type(e).__name__
|
error_type = type(e).__name__
|
||||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||||
@@ -703,7 +664,7 @@ async def stream_prompt(data, token):
|
|||||||
# log('info', "Begin prompt", prompt=prompt)
|
# log('info', "Begin prompt", prompt=prompt)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
res = await post_prompt(prompt)
|
res = post_prompt(prompt)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_type = type(e).__name__
|
error_type = type(e).__name__
|
||||||
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
|
||||||
@@ -1305,11 +1266,22 @@ def handle_execute(class_type, last_node_id, prompt_id, server, unique_id):
|
|||||||
|
|
||||||
try:
|
try:
|
||||||
origin_execute = execution.execute
|
origin_execute = execution.execute
|
||||||
is_async = asyncio.iscoroutinefunction(origin_execute)
|
|
||||||
|
|
||||||
if is_async:
|
def swizzle_execute(
|
||||||
|
server,
|
||||||
async def swizzle_execute(
|
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,
|
server,
|
||||||
dynprompt,
|
dynprompt,
|
||||||
caches,
|
caches,
|
||||||
@@ -1319,61 +1291,12 @@ try:
|
|||||||
prompt_id,
|
prompt_id,
|
||||||
execution_list,
|
execution_list,
|
||||||
pending_subgraph_results,
|
pending_subgraph_results,
|
||||||
pending_async_nodes,
|
)
|
||||||
):
|
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
|
||||||
unique_id = current_item
|
return result
|
||||||
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
|
execution.execute = swizzle_execute
|
||||||
except Exception:
|
except Exception as e:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
|
|
||||||
@@ -1517,6 +1440,10 @@ async def send_json_override(self, event, data, sid=None):
|
|||||||
logger.info(format_table(headers, table_data))
|
logger.info(format_table(headers, table_data))
|
||||||
# print("========================\n")
|
# 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
|
# the last executing event is none, then the workflow is finished
|
||||||
if event == "executing" and data.get("node") is None:
|
if event == "executing" and data.get("node") is None:
|
||||||
mark_prompt_done(prompt_id=prompt_id)
|
mark_prompt_done(prompt_id=prompt_id)
|
||||||
@@ -2132,7 +2059,6 @@ async def upload_in_background(
|
|||||||
("model_file", "format", "application/octet-stream"),
|
("model_file", "format", "application/octet-stream"),
|
||||||
("result", "format", "application/octet-stream"),
|
("result", "format", "application/octet-stream"),
|
||||||
("text_file", "format", "text/plain"),
|
("text_file", "format", "text/plain"),
|
||||||
("audio", "format", "audio/mpeg"),
|
|
||||||
]:
|
]:
|
||||||
items = data.get(file_type, [])
|
items = data.get(file_type, [])
|
||||||
|
|
||||||
@@ -2220,7 +2146,6 @@ async def update_run_with_output(
|
|||||||
or "model_file" in data
|
or "model_file" in data
|
||||||
or "result" in data
|
or "result" in data
|
||||||
or "text_file" in data
|
or "text_file" in data
|
||||||
or "audio" in data
|
|
||||||
)
|
)
|
||||||
if bypass_upload and have_upload_media:
|
if bypass_upload and have_upload_media:
|
||||||
print(
|
print(
|
||||||
@@ -2829,57 +2754,59 @@ class UploadQueue:
|
|||||||
logger.error(f"Upload failed: {str(e)}")
|
logger.error(f"Upload failed: {str(e)}")
|
||||||
logger.error(traceback.format_exc())
|
logger.error(traceback.format_exc())
|
||||||
finally:
|
finally:
|
||||||
async with self.lock: # Acquire lock to protect shared dict access
|
# Remove this upload from tracking
|
||||||
if prompt_id in self.pending_uploads:
|
if prompt_id in self.pending_uploads:
|
||||||
self.pending_uploads[prompt_id].discard(upload_id)
|
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)
|
||||||
|
|
||||||
if (
|
# If this was the last upload for this node, clean up node data
|
||||||
node_id
|
if not self.node_uploads[prompt_id][node_id]:
|
||||||
and prompt_id in self.node_uploads
|
del self.node_uploads[prompt_id][node_id]
|
||||||
and node_id in self.node_uploads[prompt_id]
|
if prompt_id in self.node_output_data:
|
||||||
):
|
if node_id in self.node_output_data[prompt_id]:
|
||||||
self.node_uploads[prompt_id][node_id].discard(upload_id)
|
if self.node_output_data[prompt_id][node_id][
|
||||||
|
"data"
|
||||||
if not self.node_uploads[prompt_id][node_id]:
|
]:
|
||||||
del self.node_uploads[prompt_id][node_id]
|
# Send final node data to API before cleanup
|
||||||
|
if prompt_metadata[
|
||||||
if (
|
prompt_id
|
||||||
prompt_id in self.node_output_data
|
].status_endpoint:
|
||||||
and node_id in self.node_output_data[prompt_id]
|
body = {
|
||||||
):
|
"run_id": prompt_id,
|
||||||
node_data = self.node_output_data[prompt_id][
|
"output_data": self.node_output_data[
|
||||||
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
|
prompt_id
|
||||||
].status_endpoint,
|
][node_id]["data"],
|
||||||
token=prompt_metadata[
|
"node_meta": {"node_id": node_id},
|
||||||
prompt_id
|
}
|
||||||
].token,
|
try:
|
||||||
json=body,
|
await async_request_with_retry(
|
||||||
)
|
"POST",
|
||||||
except Exception as e:
|
prompt_metadata[
|
||||||
logger.error(
|
prompt_id
|
||||||
f"Failed to send final node data: {str(e)}"
|
].status_endpoint,
|
||||||
)
|
token=prompt_metadata[
|
||||||
|
prompt_id
|
||||||
# Safe to delete now (re-check not strictly needed with lock, but harmless)
|
].token,
|
||||||
|
json=body,
|
||||||
|
)
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(
|
||||||
|
f"Failed to send final node data: {str(e)}"
|
||||||
|
)
|
||||||
del self.node_output_data[prompt_id][node_id]
|
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 no more pending uploads for this prompt and it's done, update status
|
||||||
if (
|
if not self.pending_uploads[prompt_id] and is_prompt_done(
|
||||||
prompt_id in self.pending_uploads
|
prompt_id
|
||||||
and not self.pending_uploads[prompt_id]
|
|
||||||
and is_prompt_done(prompt_id)
|
|
||||||
):
|
):
|
||||||
# Clean up all data for this prompt
|
# Clean up all data for this prompt
|
||||||
if prompt_id in self.node_uploads:
|
if prompt_id in self.node_uploads:
|
||||||
@@ -2892,12 +2819,9 @@ class UploadQueue:
|
|||||||
loop.create_task(update_run(prompt_id, Status.SUCCESS))
|
loop.create_task(update_run(prompt_id, Status.SUCCESS))
|
||||||
loop.create_task(send("success", {"prompt_id": prompt_id}))
|
loop.create_task(send("success", {"prompt_id": prompt_id}))
|
||||||
|
|
||||||
# Mark task as done (outside lock to avoid holding it unnecessarily)
|
# Mark task as done
|
||||||
self.queue.task_done()
|
self.queue.task_done()
|
||||||
|
|
||||||
# Send status update (also outside lock)
|
|
||||||
await self.update_queue_status(prompt_id)
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error in upload worker: {str(e)}")
|
logger.error(f"Error in upload worker: {str(e)}")
|
||||||
logger.error(traceback.format_exc())
|
logger.error(traceback.format_exc())
|
||||||
@@ -2995,7 +2919,6 @@ async def create_workflow_proxy(request):
|
|||||||
name = data.get("name")
|
name = data.get("name")
|
||||||
workflow_json = data.get("workflow_json")
|
workflow_json = data.get("workflow_json")
|
||||||
workflow_api = data.get("workflow_api")
|
workflow_api = data.get("workflow_api")
|
||||||
machine_id = data.get("machine_id")
|
|
||||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
||||||
|
|
||||||
auth_header = request.headers.get("Authorization")
|
auth_header = request.headers.get("Authorization")
|
||||||
@@ -3015,7 +2938,6 @@ async def create_workflow_proxy(request):
|
|||||||
"name": name,
|
"name": name,
|
||||||
"workflow_json": json.dumps(workflow_json),
|
"workflow_json": json.dumps(workflow_json),
|
||||||
"workflow_api": json.dumps(workflow_api),
|
"workflow_api": json.dumps(workflow_api),
|
||||||
"machine_id": machine_id,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -3138,147 +3060,3 @@ async def get_workflow_proxy(request):
|
|||||||
return web.json_response(json_data, status=response.status)
|
return web.json_response(json_data, status=response.status)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
return web.json_response({"error": str(e)}, status=500)
|
return web.json_response({"error": str(e)}, status=500)
|
||||||
|
|
||||||
|
|
||||||
# for getting a machine by id
|
|
||||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/machine")
|
|
||||||
async def get_machine_proxy(request):
|
|
||||||
machine_id = request.rel_url.query.get("machine_id")
|
|
||||||
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
|
|
||||||
auth_header = request.headers.get("Authorization")
|
|
||||||
|
|
||||||
if not auth_header:
|
|
||||||
return web.json_response(
|
|
||||||
{"error": "Authorization header is required"}, status=401
|
|
||||||
)
|
|
||||||
|
|
||||||
target_url = f"{api_url}/api/machine/{machine_id}"
|
|
||||||
|
|
||||||
try:
|
|
||||||
await ensure_client_session()
|
|
||||||
async with client_session.get(
|
|
||||||
target_url, headers={"Authorization": auth_header}
|
|
||||||
) as response:
|
|
||||||
json_data = await response.json()
|
|
||||||
return web.json_response(json_data, status=response.status)
|
|
||||||
except Exception as e:
|
|
||||||
return web.json_response({"error": str(e)}, status=500)
|
|
||||||
|
|
||||||
|
|
||||||
# for fetching docker steps from current snapshot
|
|
||||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/snapshot-to-docker")
|
|
||||||
async def snapshot_to_docker_proxy(request):
|
|
||||||
data = await request.json()
|
|
||||||
snapshot = data.get("snapshot")
|
|
||||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
|
||||||
auth_header = request.headers.get("Authorization")
|
|
||||||
|
|
||||||
if not auth_header:
|
|
||||||
return web.json_response(
|
|
||||||
{"error": "Authorization header is required"}, status=401
|
|
||||||
)
|
|
||||||
|
|
||||||
target_url = f"{api_url}/api/snapshot-to-docker"
|
|
||||||
|
|
||||||
request_body = snapshot
|
|
||||||
|
|
||||||
try:
|
|
||||||
await ensure_client_session()
|
|
||||||
async with client_session.post(
|
|
||||||
target_url, json=request_body, headers={"Authorization": auth_header}
|
|
||||||
) as response:
|
|
||||||
json_data = await response.json()
|
|
||||||
return web.json_response(json_data, status=response.status)
|
|
||||||
except Exception as e:
|
|
||||||
return web.json_response({"error": str(e)}, status=500)
|
|
||||||
|
|
||||||
|
|
||||||
# update a serverless machine with machine id
|
|
||||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/machine/update")
|
|
||||||
async def update_machine_proxy(request):
|
|
||||||
data = await request.json()
|
|
||||||
machine_id = data.get("machine_id")
|
|
||||||
comfyui_version = data.get("comfyui_version", None)
|
|
||||||
docker_steps = data.get("docker_steps")
|
|
||||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
|
||||||
auth_header = request.headers.get("Authorization")
|
|
||||||
|
|
||||||
if not auth_header:
|
|
||||||
return web.json_response(
|
|
||||||
{"error": "Authorization header is required"}, status=401
|
|
||||||
)
|
|
||||||
|
|
||||||
target_url = f"{api_url}/api/machine/serverless/{machine_id}"
|
|
||||||
|
|
||||||
request_body = {"docker_command_steps": docker_steps}
|
|
||||||
|
|
||||||
if comfyui_version:
|
|
||||||
request_body["comfyui_version"] = comfyui_version
|
|
||||||
|
|
||||||
try:
|
|
||||||
await ensure_client_session()
|
|
||||||
async with client_session.patch(
|
|
||||||
target_url, json=request_body, headers={"Authorization": auth_header}
|
|
||||||
) as response:
|
|
||||||
json_data = await response.json()
|
|
||||||
return web.json_response(json_data, status=response.status)
|
|
||||||
except Exception as e:
|
|
||||||
return web.json_response({"error": str(e)}, status=500)
|
|
||||||
|
|
||||||
|
|
||||||
@server.PromptServer.instance.routes.post("/comfyui-deploy/machine/create")
|
|
||||||
async def create_machine_proxy(request):
|
|
||||||
data = await request.json()
|
|
||||||
name = data.get("name")
|
|
||||||
docker_command_steps = data.get("docker_command_steps")
|
|
||||||
comfyui_version = data.get("comfyui_version")
|
|
||||||
api_url = data.get("api_url", "https://api.comfydeploy.com")
|
|
||||||
auth_header = request.headers.get("Authorization")
|
|
||||||
|
|
||||||
if not auth_header:
|
|
||||||
return web.json_response(
|
|
||||||
{"error": "Authorization header is required"}, status=401
|
|
||||||
)
|
|
||||||
|
|
||||||
target_url = f"{api_url}/api/machine/serverless"
|
|
||||||
|
|
||||||
request_body = {
|
|
||||||
"name": name,
|
|
||||||
"docker_command_steps": docker_command_steps,
|
|
||||||
"comfyui_version": comfyui_version,
|
|
||||||
"gpu": "A10G",
|
|
||||||
}
|
|
||||||
|
|
||||||
try:
|
|
||||||
await ensure_client_session()
|
|
||||||
async with client_session.post(
|
|
||||||
target_url, json=request_body, headers={"Authorization": auth_header}
|
|
||||||
) as response:
|
|
||||||
json_data = await response.json()
|
|
||||||
return web.json_response(json_data, status=response.status)
|
|
||||||
except Exception as e:
|
|
||||||
return web.json_response({"error": str(e)}, status=500)
|
|
||||||
|
|
||||||
|
|
||||||
# get latest comfyui version
|
|
||||||
@server.PromptServer.instance.routes.get("/comfyui-deploy/comfyui-version")
|
|
||||||
async def get_comfyui_version_proxy(request):
|
|
||||||
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
|
|
||||||
auth_header = request.headers.get("Authorization")
|
|
||||||
|
|
||||||
if not auth_header:
|
|
||||||
return web.json_response(
|
|
||||||
{"error": "Authorization header is required"}, status=401
|
|
||||||
)
|
|
||||||
|
|
||||||
target_url = f"{api_url}/api/latest-hashes"
|
|
||||||
|
|
||||||
try:
|
|
||||||
await ensure_client_session()
|
|
||||||
async with client_session.get(
|
|
||||||
target_url, headers={"Authorization": auth_header}
|
|
||||||
) as response:
|
|
||||||
json_data = await response.json()
|
|
||||||
return web.json_response(json_data, status=response.status)
|
|
||||||
except Exception as e:
|
|
||||||
return web.json_response({"error": str(e)}, status=500)
|
|
||||||
|
|||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
[project]
|
[project]
|
||||||
name = "comfyui-deploy"
|
name = "comfyui-deploy"
|
||||||
description = "Open source comfyui deployment platform, a vercel for generative workflow infra."
|
description = "Open source comfyui deployment platform, a vercel for generative workflow infra."
|
||||||
version = "2.3.2"
|
version = "2.1.0"
|
||||||
license = { file = "LICENSE" }
|
license = { file = "LICENSE" }
|
||||||
dependencies = ["aiofiles", "pydantic", "opencv-python", "imageio-ffmpeg", "tabulate", "brotli"]
|
dependencies = ["aiofiles", "pydantic", "opencv-python", "imageio-ffmpeg", "tabulate", "brotli"]
|
||||||
|
|
||||||
|
|||||||
+354
-894
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -1,82 +0,0 @@
|
|||||||
// Snapshot Utilities
|
|
||||||
// Centralized snapshot fetching with ComfyUI version fallback
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Fetches the current snapshot with ComfyUI version fallback
|
|
||||||
* If the snapshot response has null comfyui field, it will fetch the latest ComfyUI version
|
|
||||||
* and update the snapshot with the comfyui_hash
|
|
||||||
*
|
|
||||||
* @param {Function} getDataFn - Function that returns { apiKey, apiUrl } for ComfyUI version API calls
|
|
||||||
* @returns {Promise<Object>} - The snapshot data with comfyui field populated
|
|
||||||
*/
|
|
||||||
export async function fetchSnapshot(getDataFn = null) {
|
|
||||||
try {
|
|
||||||
// Fetch the current snapshot
|
|
||||||
const response = await fetch("/snapshot/get_current");
|
|
||||||
if (!response.ok) {
|
|
||||||
throw new Error(`Snapshot fetch failed: ${response.status}`);
|
|
||||||
}
|
|
||||||
|
|
||||||
const snapshot = await response.json();
|
|
||||||
|
|
||||||
// Check if comfyui field is null and we have getDataFn for fallback
|
|
||||||
if (snapshot.comfyui === null && getDataFn) {
|
|
||||||
console.log(
|
|
||||||
"ComfyUI version is null in snapshot, fetching latest version..."
|
|
||||||
);
|
|
||||||
|
|
||||||
try {
|
|
||||||
const data = getDataFn();
|
|
||||||
if (data && data.apiKey) {
|
|
||||||
const comfyuiVersionResponse = await fetch(
|
|
||||||
`/comfyui-deploy/comfyui-version?api_url=${encodeURIComponent(
|
|
||||||
data.apiUrl || "https://api.comfydeploy.com"
|
|
||||||
)}`,
|
|
||||||
{
|
|
||||||
headers: {
|
|
||||||
Authorization: `Bearer ${data.apiKey}`,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
);
|
|
||||||
|
|
||||||
if (comfyuiVersionResponse.ok) {
|
|
||||||
const versionData = await comfyuiVersionResponse.json();
|
|
||||||
if (versionData.comfyui_hash) {
|
|
||||||
console.log(
|
|
||||||
`Using ComfyUI hash from API: ${versionData.comfyui_hash}`
|
|
||||||
);
|
|
||||||
snapshot.comfyui = versionData.comfyui_hash;
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
console.warn(
|
|
||||||
"Failed to fetch ComfyUI version from API:",
|
|
||||||
comfyuiVersionResponse.status
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
} catch (error) {
|
|
||||||
console.warn("Error fetching ComfyUI version fallback:", error);
|
|
||||||
// Continue with original snapshot even if fallback fails
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return snapshot;
|
|
||||||
} catch (error) {
|
|
||||||
console.error("Error fetching snapshot:", error);
|
|
||||||
throw error;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Simple snapshot fetch without ComfyUI version fallback
|
|
||||||
* Use this when you don't need the ComfyUI version fallback logic
|
|
||||||
*
|
|
||||||
* @returns {Promise<Object>} - The snapshot data as-is
|
|
||||||
*/
|
|
||||||
export async function fetchSnapshotSimple() {
|
|
||||||
const response = await fetch("/snapshot/get_current");
|
|
||||||
if (!response.ok) {
|
|
||||||
throw new Error(`Snapshot fetch failed: ${response.status}`);
|
|
||||||
}
|
|
||||||
return response.json();
|
|
||||||
}
|
|
||||||
@@ -49,7 +49,6 @@ async function fetchWorkflows(getData, offset = 0, limit = 20, search = "") {
|
|||||||
|
|
||||||
function createWorkflowItem(workflow, getTimeAgo, getData) {
|
function createWorkflowItem(workflow, getTimeAgo, getData) {
|
||||||
const li = document.createElement("li");
|
const li = document.createElement("li");
|
||||||
let loadingToast = null;
|
|
||||||
li.style.cssText = `
|
li.style.cssText = `
|
||||||
border-bottom: 1px solid #444;
|
border-bottom: 1px solid #444;
|
||||||
background: transparent;
|
background: transparent;
|
||||||
@@ -75,7 +74,7 @@ function createWorkflowItem(workflow, getTimeAgo, getData) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Show loading toast
|
// Show loading toast
|
||||||
loadingToast = window.app.extensionManager.toast.add({
|
const loadingToast = window.app.extensionManager.toast.add({
|
||||||
severity: "info",
|
severity: "info",
|
||||||
summary: "Loading workflow...",
|
summary: "Loading workflow...",
|
||||||
detail: `Loading "${workflow.name}"`,
|
detail: `Loading "${workflow.name}"`,
|
||||||
@@ -109,33 +108,6 @@ function createWorkflowItem(workflow, getTimeAgo, getData) {
|
|||||||
// Load the workflow
|
// Load the workflow
|
||||||
window.app.loadGraphData(latestVersion.workflow);
|
window.app.loadGraphData(latestVersion.workflow);
|
||||||
|
|
||||||
// Wait a bit for the graph to fully load before checking for ComfyDeploy node
|
|
||||||
await new Promise((resolve) => setTimeout(resolve, 100));
|
|
||||||
|
|
||||||
// Check if ComfyDeploy node exists, if not add it back
|
|
||||||
const graph = window.app.graph;
|
|
||||||
let deployMeta = graph.findNodesByType("ComfyDeploy");
|
|
||||||
|
|
||||||
if (deployMeta.length === 0) {
|
|
||||||
// Add ComfyDeploy node with workflow metadata
|
|
||||||
graph.beforeChange();
|
|
||||||
const node = LiteGraph.createNode("ComfyDeploy");
|
|
||||||
node.configure({
|
|
||||||
widgets_values: [
|
|
||||||
workflow.name, // workflow_name
|
|
||||||
workflow.id, // workflow_id
|
|
||||||
latestVersion.version, // version
|
|
||||||
],
|
|
||||||
});
|
|
||||||
node.pos = [0, 0];
|
|
||||||
graph.add(node);
|
|
||||||
graph.afterChange();
|
|
||||||
|
|
||||||
console.log(
|
|
||||||
`Added ComfyDeploy node with: name="${workflow.name}", id="${workflow.id}", version="${latestVersion.version}"`
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Show success toast
|
// Show success toast
|
||||||
window.app.extensionManager.toast.add({
|
window.app.extensionManager.toast.add({
|
||||||
severity: "success",
|
severity: "success",
|
||||||
@@ -155,9 +127,7 @@ function createWorkflowItem(workflow, getTimeAgo, getData) {
|
|||||||
life: 5000,
|
life: 5000,
|
||||||
});
|
});
|
||||||
} finally {
|
} finally {
|
||||||
if (loadingToast) {
|
loadingToast.close();
|
||||||
loadingToast.close();
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -309,7 +279,7 @@ async function initializeWorkflowsList(element, getData, getTimeAgo) {
|
|||||||
list-style-type: none;
|
list-style-type: none;
|
||||||
padding: 0;
|
padding: 0;
|
||||||
margin: 0;
|
margin: 0;
|
||||||
height: calc(100vh - 550px);
|
height: calc(100vh - 350px);
|
||||||
overflow-y: auto;
|
overflow-y: auto;
|
||||||
scrollbar-width: thin;
|
scrollbar-width: thin;
|
||||||
scrollbar-color: #666 transparent;
|
scrollbar-color: #666 transparent;
|
||||||
|
|||||||
Reference in New Issue
Block a user