FIX webhook直接返回ComfyUI结果
This commit is contained in:
3
.gitignore
vendored
3
.gitignore
vendored
@@ -1,4 +1,5 @@
|
||||
.pypirc
|
||||
.env
|
||||
.idea
|
||||
test*
|
||||
test*
|
||||
.runtime.env
|
||||
@@ -73,8 +73,7 @@ async def comfyui_v1_status(task_id: str, response: Response):
|
||||
if task_status != "success":
|
||||
reason = result["msg"]
|
||||
else:
|
||||
media = MediaSource.from_str("s3://" + "/".join(
|
||||
[config.S3_region, config.S3_bucket_name, config.comfyui_s3_output, result["file_name"]]))
|
||||
media = MediaSource.from_str(result["file_name"])
|
||||
return ComfyTaskStatusResponse(taskId=task_id, status=task_status, code=code,
|
||||
error=reason, result=media.urn if media else None)
|
||||
|
||||
@@ -122,8 +121,7 @@ async def comfyui_v2_status(task_id: str, response: Response):
|
||||
if task_status != "success":
|
||||
reason = result["msg"]
|
||||
else:
|
||||
media = MediaSource.from_str("s3://" + "/".join(
|
||||
[config.S3_region, config.S3_bucket_name, config.comfyui_s3_output, result["file_name"]]))
|
||||
media = MediaSource.from_str(result["file_name"])
|
||||
return ComfyTaskStatusResponse(taskId=task_id, status=task_status, code=code,
|
||||
error=reason,
|
||||
result=media.urn if media else "")
|
||||
@@ -55,6 +55,10 @@ class SentryUtils:
|
||||
status = TaskStatus.success
|
||||
error = None
|
||||
code = ErrorCode.SUCCESS.value
|
||||
if not result:
|
||||
status = TaskStatus.failed
|
||||
error = "Internal Error"
|
||||
code = ErrorCode.BUSINESS_ERROR.value
|
||||
except Exception as e:
|
||||
logger.exception(e)
|
||||
result = None
|
||||
@@ -72,9 +76,13 @@ class SentryUtils:
|
||||
code=code,
|
||||
results=[result] if isinstance(result, str) else result,
|
||||
).model_dump())
|
||||
logger.info(f"webhook response = {response}")
|
||||
case WebhookMethodEnum.GET:
|
||||
response = client.post(url=webhook.endpoint.__str__())
|
||||
response = client.get(url=webhook.endpoint.__str__(), params=BaseFFMPEGTaskStatusResponse(
|
||||
taskId=func_id, status=status, error=error,
|
||||
code=code,
|
||||
results=[result] if isinstance(result, str) else result,
|
||||
).model_dump())
|
||||
logger.info(f"webhook response = {response}")
|
||||
response.raise_for_status()
|
||||
return result
|
||||
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
# ComfyUI模板--Base Auth
|
||||
import backoff
|
||||
import modal
|
||||
from dotenv import dotenv_values
|
||||
|
||||
@@ -67,7 +68,6 @@ app.set_description("ComfyUI v1 Server")
|
||||
with comfyui_image.imports():
|
||||
import json
|
||||
import os
|
||||
import time
|
||||
import shutil
|
||||
import subprocess
|
||||
import uuid
|
||||
@@ -153,7 +153,6 @@ with comfyui_image.imports():
|
||||
import datetime
|
||||
import tempfile
|
||||
import traceback
|
||||
import sys
|
||||
from loguru import logger
|
||||
|
||||
# 生成当前函数调用的唯一ID
|
||||
@@ -161,7 +160,7 @@ with comfyui_image.imports():
|
||||
logger.info(f"开始处理任务,会话ID: {self.session_id}, 调用ID: {function_call_id}")
|
||||
|
||||
@SentryUtils.webhook_handler(webhook=webhook, func_id=function_call_id)
|
||||
async def infer(workflow: Dict):
|
||||
async def infer(workflow: Dict) -> str | None:
|
||||
self.poll_server_health()
|
||||
self.prompt_uuid = str(uuid.uuid4())
|
||||
logger.info(f"Workflow JSON: \n{json.dumps(workflow, indent=4, ensure_ascii=False)}")
|
||||
@@ -235,7 +234,7 @@ with comfyui_image.imports():
|
||||
config.comfyui_s3_output + "/" + f)
|
||||
except:
|
||||
raise Exception("Failed to move file to S3 manually")
|
||||
return f
|
||||
return "s3://" + "/".join([config.S3_region, config.S3_bucket_name, config.comfyui_s3_output, f])
|
||||
|
||||
# 使用临时文件
|
||||
temp_log_path = None
|
||||
@@ -526,32 +525,6 @@ with comfyui_image.imports():
|
||||
cleanup_span.set_tag("success", "false")
|
||||
cleanup_successful = False
|
||||
|
||||
# logger.info("清理资源并重启ComfyUI")
|
||||
# with sentry_sdk.start_span(op="system.restart", description="重启ComfyUI") as restart_span:
|
||||
# try:
|
||||
# # 停止ComfyUI
|
||||
# with sentry_sdk.start_span(op="system.stop", description="停止ComfyUI") as stop_span:
|
||||
# cmd = "comfy stop"
|
||||
# subprocess.run(cmd, shell=True, check=True)
|
||||
# time.sleep(1)
|
||||
# stop_span.set_tag("success", "true")
|
||||
# except Exception as stop_error:
|
||||
# logger.error(f"停止ComfyUI失败: {stop_error}")
|
||||
# restart_span.set_data("stop_error", str(stop_error))
|
||||
# restart_span.set_tag("stop_success", "false")
|
||||
#
|
||||
# try:
|
||||
# # 重启ComfyUI
|
||||
# with sentry_sdk.start_span(op="system.start", description="启动ComfyUI") as start_span:
|
||||
# cmd = "comfy launch --background"
|
||||
# subprocess.run(cmd, shell=True, check=True)
|
||||
# start_span.set_tag("success", "true")
|
||||
# except Exception as start_error:
|
||||
# logger.error(f"启动ComfyUI失败: {start_error}")
|
||||
# restart_span.set_data("start_error", str(start_error))
|
||||
# restart_span.set_tag("start_success", "false")
|
||||
# modal.experimental.stop_fetching_inputs()
|
||||
|
||||
# 根据清理和重启结果设置总体清理状态
|
||||
if cleanup_successful:
|
||||
cleanup_span.set_status("ok")
|
||||
@@ -564,7 +537,6 @@ with comfyui_image.imports():
|
||||
import urllib
|
||||
|
||||
try:
|
||||
# dummy request to check if the server is healthy
|
||||
req = urllib.request.Request("http://127.0.0.1:8188/system_stats")
|
||||
urllib.request.urlopen(req, timeout=5)
|
||||
print("ComfyUI server is healthy")
|
||||
|
||||
@@ -51,7 +51,6 @@ with comfyui_latentsync_1_5_image.imports():
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
import time
|
||||
import uuid
|
||||
from typing import Dict, Tuple, Any, Optional
|
||||
import sentry_sdk
|
||||
@@ -126,7 +125,6 @@ with comfyui_latentsync_1_5_image.imports():
|
||||
import datetime
|
||||
import tempfile
|
||||
import traceback
|
||||
import sys
|
||||
from loguru import logger
|
||||
|
||||
# 生成当前函数调用的唯一ID
|
||||
@@ -134,7 +132,7 @@ with comfyui_latentsync_1_5_image.imports():
|
||||
logger.info(f"开始处理任务,会话ID: {self.session_id}, 调用ID: {function_call_id}")
|
||||
|
||||
@SentryUtils.webhook_handler(webhook=webhook, func_id=function_call_id)
|
||||
async def infer(workflow: Dict):
|
||||
async def infer(workflow: Dict) -> str | None:
|
||||
self.poll_server_health()
|
||||
self.prompt_uuid = str(uuid.uuid4())
|
||||
logger.info(f"Workflow JSON: \n{json.dumps(workflow, indent=4, ensure_ascii=False)}")
|
||||
@@ -206,7 +204,7 @@ with comfyui_latentsync_1_5_image.imports():
|
||||
config.comfyui_s3_output + "/" + f)
|
||||
except:
|
||||
raise Exception("Failed to move file to S3 manually")
|
||||
return f
|
||||
return "s3://" + "/".join([config.S3_region, config.S3_bucket_name, config.comfyui_s3_output, f])
|
||||
|
||||
# 使用临时文件
|
||||
temp_log_path = None
|
||||
@@ -495,32 +493,6 @@ with comfyui_latentsync_1_5_image.imports():
|
||||
cleanup_span.set_tag("success", "false")
|
||||
cleanup_successful = False
|
||||
|
||||
# logger.info("清理资源并重启ComfyUI")
|
||||
# with sentry_sdk.start_span(op="system.restart", description="重启ComfyUI") as restart_span:
|
||||
# try:
|
||||
# # 停止ComfyUI
|
||||
# with sentry_sdk.start_span(op="system.stop", description="停止ComfyUI") as stop_span:
|
||||
# cmd = "comfy stop"
|
||||
# subprocess.run(cmd, shell=True, check=True)
|
||||
# time.sleep(1)
|
||||
# stop_span.set_tag("success", "true")
|
||||
# except Exception as stop_error:
|
||||
# logger.error(f"停止ComfyUI失败: {stop_error}")
|
||||
# restart_span.set_data("stop_error", str(stop_error))
|
||||
# restart_span.set_tag("stop_success", "false")
|
||||
#
|
||||
# try:
|
||||
# # 重启ComfyUI
|
||||
# with sentry_sdk.start_span(op="system.start", description="启动ComfyUI") as start_span:
|
||||
# cmd = "comfy launch --background"
|
||||
# subprocess.run(cmd, shell=True, check=True)
|
||||
# start_span.set_tag("success", "true")
|
||||
# except Exception as start_error:
|
||||
# logger.error(f"启动ComfyUI失败: {start_error}")
|
||||
# restart_span.set_data("start_error", str(start_error))
|
||||
# restart_span.set_tag("start_success", "false")
|
||||
# modal.experimental.stop_fetching_inputs()
|
||||
|
||||
# 根据清理和重启结果设置总体清理状态
|
||||
if cleanup_successful:
|
||||
cleanup_span.set_status("ok")
|
||||
|
||||
Reference in New Issue
Block a user