Compare commits

..
Author SHA1 Message Date
KarrixLee 0e1f1cd6ba tweak 2025-06-28 00:33:49 +08:00
KarrixLee 57dd0d9167 tweal 2025-06-28 00:32:34 +08:00
KarrixLee b2de923440 hi 2025-06-28 00:30:09 +08:00
KarrixLee 6d2b918ef1 tweak 2025-06-28 00:27:00 +08:00
KarrixLee d8b1e88c03 hi 2025-06-28 00:23:25 +08:00
10 changed files with 475 additions and 6164 deletions
-2
View File
@@ -2,8 +2,6 @@
Open source comfyui deployment platform, a `vercel` for generative workflow infra. (serverless hosted gpu with vertical intergation with comfyui)
Check out our latest lcoal demo -> https://github.com/comfy-deploy/comfyui-api-comfydeploy
> [!NOTE]
> Im looking for creative hacker to join ComfyDeploy's core team! DM me on [twitter](https://x.com/BennyKokMusic)
+1 -1
View File
@@ -92,7 +92,7 @@ class ComfyDeployOutputText:
}
)
return {"ui": {"text_file": results}}
return {"ui": {"text_file": results, "text": [text]}}
NODE_CLASS_MAPPINGS = {"ComfyDeployOutputText": ComfyDeployOutputText}
+86 -667
View File
@@ -265,7 +265,7 @@ def clear_current_prompt(sid):
streaming_prompt_metadata[sid].running_prompt_ids.clear()
async def post_prompt(json_data):
def post_prompt(json_data):
prompt_server = server.PromptServer.instance
json_data = prompt_server.trigger_on_prompt(json_data)
@@ -281,48 +281,7 @@ async def post_prompt(json_data):
if "prompt" in json_data:
prompt = json_data["prompt"]
prompt_id = json_data.get("prompt_id") or str(uuid.uuid4())
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
valid = execution.validate_prompt(prompt)
extra_data = {}
if "extra_data" in json_data:
extra_data = json_data["extra_data"]
@@ -333,6 +292,8 @@ async 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)
@@ -539,15 +500,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,
@@ -558,7 +519,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}")
@@ -634,7 +595,7 @@ async def comfy_deploy_run(request):
)
try:
res = await post_prompt(prompt)
res = post_prompt(prompt)
except Exception as e:
error_type = type(e).__name__
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
@@ -672,14 +633,6 @@ async def comfy_deploy_run(request):
return web.json_response(res, status=status)
@server.PromptServer.instance.routes.post("/comfyui-deploy/interrupt")
async def interrupt_prompt(request):
data = await request.json()
prompt_id = data.get("prompt_id")
await update_run(prompt_id, Status.CANCELLED)
return web.json_response({"message": "Prompt interrupted"}, status=200)
async def stream_prompt(data, token):
# In older version, we use workflow_api, but this has inputs already swapped in nextjs frontend, which is tricky
workflow_api = data.get("workflow_api_raw")
@@ -711,7 +664,7 @@ async def stream_prompt(data, token):
# log('info', "Begin prompt", prompt=prompt)
try:
res = await post_prompt(prompt)
res = post_prompt(prompt)
except Exception as e:
error_type = type(e).__name__
stack_trace_short = traceback.format_exc().strip().split("\n")[-2]
@@ -1313,11 +1266,22 @@ def handle_execute(class_type, last_node_id, prompt_id, server, unique_id):
try:
origin_execute = execution.execute
is_async = asyncio.iscoroutinefunction(origin_execute)
if is_async:
async def swizzle_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(
server,
dynprompt,
caches,
@@ -1327,61 +1291,12 @@ try:
prompt_id,
execution_list,
pending_subgraph_results,
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
)
handle_execute(class_type, last_node_id, prompt_id, server, unique_id)
return result
execution.execute = swizzle_execute
except Exception:
except Exception as e:
pass
@@ -1525,6 +1440,10 @@ 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)
@@ -2140,7 +2059,6 @@ 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, [])
@@ -2228,7 +2146,6 @@ 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(
@@ -2837,57 +2754,59 @@ class UploadQueue:
logger.error(f"Upload failed: {str(e)}")
logger.error(traceback.format_exc())
finally:
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)
# 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)
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[
# 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[
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)}"
)
# Safe to delete now (re-check not strictly needed with lock, but harmless)
][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)}"
)
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 (
prompt_id in self.pending_uploads
and not self.pending_uploads[prompt_id]
and is_prompt_done(prompt_id)
if 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:
@@ -2900,12 +2819,9 @@ class UploadQueue:
loop.create_task(update_run(prompt_id, Status.SUCCESS))
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()
# 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())
@@ -3003,7 +2919,6 @@ async def create_workflow_proxy(request):
name = data.get("name")
workflow_json = data.get("workflow_json")
workflow_api = data.get("workflow_api")
machine_id = data.get("machine_id")
api_url = data.get("api_url", "https://api.comfydeploy.com")
auth_header = request.headers.get("Authorization")
@@ -3023,7 +2938,6 @@ async def create_workflow_proxy(request):
"name": name,
"workflow_json": json.dumps(workflow_json),
"workflow_api": json.dumps(workflow_api),
"machine_id": machine_id,
}
try:
@@ -3146,498 +3060,3 @@ async def get_workflow_proxy(request):
return web.json_response(json_data, status=response.status)
except Exception as e:
return web.json_response({"error": str(e)}, status=500)
# fetch workflow versions (infinite scroll support)
@server.PromptServer.instance.routes.get("/comfyui-deploy/workflow/versions")
async def get_workflow_versions_proxy(request):
workflow_id = request.rel_url.query.get("workflow_id")
api_url = request.rel_url.query.get("api_url", "https://api.comfydeploy.com")
search = request.rel_url.query.get("search", "")
limit = request.rel_url.query.get("limit", "20")
offset = request.rel_url.query.get("offset", "0")
auth_header = request.headers.get("Authorization")
if not auth_header:
return web.json_response(
{"error": "Authorization header is required"}, status=401
)
# Build target URL with query params
params = {"limit": limit, "offset": offset}
if search:
params["search"] = search
query = urlencode(params)
target_url = f"{api_url}/api/workflow/{workflow_id}/versions"
if query:
target_url += f"?{query}"
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)
# fetch a specific workflow version json
@server.PromptServer.instance.routes.get("/comfyui-deploy/workflow/version")
async def get_workflow_version_proxy(request):
workflow_id = request.rel_url.query.get("workflow_id")
version = request.rel_url.query.get("version")
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/workflow/{workflow_id}/version/{version}"
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 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)
# Proxy: generate single-part upload URL
@server.PromptServer.instance.routes.post(
"/comfyui-deploy/volume/file/generate-upload-url"
)
async def proxy_generate_upload_url(request):
data = await request.json()
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/volume/file/generate-upload-url"
try:
await ensure_client_session()
async with client_session.post(
target_url,
json={
"filename": data.get("filename"),
"contentType": data.get("contentType"),
"size": data.get("size"),
},
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)
# Proxy: initiate multipart upload
@server.PromptServer.instance.routes.post(
"/comfyui-deploy/volume/file/initiate-multipart-upload"
)
async def proxy_initiate_multipart_upload(request):
data = await request.json()
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/volume/file/initiate-multipart-upload"
try:
await ensure_client_session()
async with client_session.post(
target_url,
json={
"filename": data.get("filename"),
"contentType": data.get("contentType"),
"size": data.get("size"),
},
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)
# Proxy: generate part upload URL
@server.PromptServer.instance.routes.post(
"/comfyui-deploy/volume/file/generate-part-upload-url"
)
async def proxy_generate_part_upload_url(request):
data = await request.json()
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/volume/file/generate-part-upload-url"
try:
await ensure_client_session()
async with client_session.post(
target_url,
json={
"uploadId": data.get("uploadId"),
"key": data.get("key"),
"partNumber": data.get("partNumber"),
},
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)
# Proxy: complete multipart upload
@server.PromptServer.instance.routes.post(
"/comfyui-deploy/volume/file/complete-multipart-upload"
)
async def proxy_complete_multipart_upload(request):
data = await request.json()
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/volume/file/complete-multipart-upload"
try:
await ensure_client_session()
async with client_session.post(
target_url,
json={
"uploadId": data.get("uploadId"),
"key": data.get("key"),
"parts": data.get("parts"),
},
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)
# Proxy: abort multipart upload
@server.PromptServer.instance.routes.post(
"/comfyui-deploy/volume/file/abort-multipart-upload"
)
async def proxy_abort_multipart_upload(request):
data = await request.json()
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/volume/file/abort-multipart-upload"
try:
await ensure_client_session()
async with client_session.post(
target_url,
json={
"uploadId": data.get("uploadId"),
"key": data.get("key"),
},
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)
# Proxy: add model (unified endpoint)
@server.PromptServer.instance.routes.post("/comfyui-deploy/volume/model")
async def proxy_add_model(request):
data = await request.json()
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/volume/model"
# pass body through but remove api_url key
forward_body = dict(data)
if "api_url" in forward_body:
forward_body.pop("api_url")
try:
await ensure_client_session()
async with client_session.post(
target_url, json=forward_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)
# FS: stat file (size)
@server.PromptServer.instance.routes.get("/comfyui-deploy/fs/stat")
async def fs_stat(request):
try:
import os
file_path = request.rel_url.query.get("path")
if not file_path:
return web.json_response({"error": "path is required"}, status=400)
# Basic safeguard: ensure it's a ComfyUI models path
if "/models/" not in file_path:
return web.json_response({"error": "invalid path"}, status=400)
st = os.stat(file_path)
return web.json_response({"size": st.st_size})
except FileNotFoundError:
return web.json_response({"error": "not found"}, status=404)
except Exception as e:
return web.json_response({"error": str(e)}, status=500)
# Upload a multipart part directly from machine filesystem to S3 upload URL
@server.PromptServer.instance.routes.post(
"/comfyui-deploy/volume/file/upload-part-from-path"
)
async def upload_part_from_path(request):
try:
import os
import io
data = await request.json()
file_path = data.get("filePath")
upload_url = data.get("uploadUrl")
start = int(data.get("start", 0))
end = int(data.get("end", 0))
if not file_path or not upload_url:
return web.json_response(
{"error": "filePath and uploadUrl are required"}, status=400
)
if "/models/" not in file_path:
return web.json_response({"error": "invalid filePath"}, status=400)
if end <= start:
return web.json_response({"error": "invalid byte range"}, status=400)
size = end - start
await ensure_client_session()
# Important: S3 pre-signed part uploads do not support chunked transfer
# Buffer the exact part into memory to provide a Content-Length header
buffer = bytearray()
chunk_size = 4 * 1024 * 1024
with open(file_path, "rb") as f:
f.seek(start)
remaining = size
while remaining > 0:
to_read = chunk_size if remaining >= chunk_size else remaining
chunk = f.read(to_read)
if not chunk:
break
buffer.extend(chunk)
remaining -= len(chunk)
if len(buffer) != size:
return web.json_response(
{
"error": "read size mismatch",
"expected": size,
"actual": len(buffer),
},
status=500,
)
headers = {
"Content-Length": str(size),
"Content-Type": "application/octet-stream",
}
async with client_session.put(
upload_url, data=bytes(buffer), headers=headers
) as resp:
text = await resp.text()
if resp.status < 200 or resp.status >= 300:
return web.json_response(
{"error": f"upload failed: {resp.status}", "body": text},
status=resp.status,
)
etag = resp.headers.get("ETag") or resp.headers.get("etag") or ""
etag = etag.replace('"', "")
return web.json_response({"eTag": etag, "bytesSent": size})
except Exception as e:
return web.json_response({"error": str(e)}, status=500)
-1
View File
@@ -18,7 +18,6 @@ class Status(Enum):
SUCCESS = "success"
FAILED = "failed"
UPLOADING = "uploading"
CANCELLED = "cancelled"
class StreamingPrompt(BaseModel):
+1 -1
View File
@@ -1,7 +1,7 @@
[project]
name = "comfyui-deploy"
description = "Open source comfyui deployment platform, a vercel for generative workflow infra."
version = "2.3.4"
version = "2.1.0"
license = { file = "LICENSE" }
dependencies = ["aiofiles", "pydantic", "opencv-python", "imageio-ffmpeg", "tabulate", "brotli"]
+384 -1065
View File
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
-82
View File
@@ -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();
}
File diff suppressed because it is too large Load Diff