diff --git a/.gitignore b/.gitignore index b361dc6..3df1f6e 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,5 @@ .pypirc .env .idea -test* \ No newline at end of file +test* +.runtime.env \ No newline at end of file diff --git a/src/BowongModalFunctions/router/comfyui.py b/src/BowongModalFunctions/router/comfyui.py index 35125a5..0c7732c 100644 --- a/src/BowongModalFunctions/router/comfyui.py +++ b/src/BowongModalFunctions/router/comfyui.py @@ -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 "") \ No newline at end of file diff --git a/src/BowongModalFunctions/utils/SentryUtils.py b/src/BowongModalFunctions/utils/SentryUtils.py index 599bc8d..f5b1817 100644 --- a/src/BowongModalFunctions/utils/SentryUtils.py +++ b/src/BowongModalFunctions/utils/SentryUtils.py @@ -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 diff --git a/src/cluster/comfyui_v1.py b/src/cluster/comfyui_v1.py index 11ffd47..27e2d50 100644 --- a/src/cluster/comfyui_v1.py +++ b/src/cluster/comfyui_v1.py @@ -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") diff --git a/src/cluster/comfyui_v2.py b/src/cluster/comfyui_v2.py index 3830626..d445b5f 100644 --- a/src/cluster/comfyui_v2.py +++ b/src/cluster/comfyui_v2.py @@ -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")