diff --git a/backend/app/core/ap_scheduler.py b/backend/app/core/ap_scheduler.py index aacfa79b..1b16191b 100644 --- a/backend/app/core/ap_scheduler.py +++ b/backend/app/core/ap_scheduler.py @@ -143,7 +143,8 @@ class SchedulerUtil: status="running", ) else: - log.warning(f"任务 {job_id} 提交执行,但未找到任务信息") + # 任务可能已经被移除(一次性任务),尝试从事件中获取信息 + log.debug(f"任务 {job_id} 提交执行,但未找到任务信息(可能已被移除)") @classmethod def _handle_job_executed(cls, event: JobExecutionEvent) -> None: @@ -643,7 +644,7 @@ class SchedulerUtil: job = cls.get_job(job_id=job_id) next_run_time = str(job.next_run_time) if job and job.next_run_time else None - job_state = cls._get_job_state(job) + job_state = cls._get_job_state(job) if job else None with Session(engine) as session: job_log = JobModel( @@ -669,7 +670,7 @@ class SchedulerUtil: job = cls.get_job(job_id=job_id) next_run_time = str(job.next_run_time) if job and job.next_run_time else None - job_state = cls._get_job_state(job) + job_state = cls._get_job_state(job) if job else None with Session(engine) as session: job_log = ( @@ -689,6 +690,8 @@ class SchedulerUtil: if error: job_log.error = error session.commit() + else: + log.warning(f"未找到任务 {job_id} 的待执行或运行中日志记录") @classmethod def _update_latest_job_log(cls, job_id: str, status: str, result: str | None = None, error: str | None = None) -> None: @@ -702,7 +705,7 @@ class SchedulerUtil: job = cls.get_job(job_id=job_id) next_run_time = str(job.next_run_time) if job and job.next_run_time else None - job_state = cls._get_job_state(job) + job_state = cls._get_job_state(job) if job else None with Session(engine) as session: job_log = ( @@ -722,6 +725,8 @@ class SchedulerUtil: if error: job_log.error = error session.commit() + else: + log.warning(f"未找到任务 {job_id} 的运行中日志记录") @classmethod def _update_job_log_on_removed(cls, job_id: str) -> None: @@ -771,8 +776,18 @@ class SchedulerUtil: """ 立即执行任务(添加到调度器并立即运行) """ - trigger = DateTrigger(run_date=datetime.now()) - return cls._add_job_with_trigger(job_info, trigger) + # 使用稍微延迟的时间,确保事件监听器能够捕获事件 + from datetime import timedelta + trigger = DateTrigger(run_date=datetime.now() + timedelta(seconds=0.1)) + job = cls._add_job_with_trigger(job_info, trigger) + # 手动创建执行日志,确保调试时也能生成记录 + cls._create_job_log( + job_id=str(job_info.id), + job_name=job_info.name, + trigger_type="manual", + status="running", + ) + return job @classmethod def add_cron_job( diff --git a/backend/app/plugin/module_ai/chat/ws.py b/backend/app/plugin/module_ai/chat/ws.py index 238cfa22..0ce16acf 100644 --- a/backend/app/plugin/module_ai/chat/ws.py +++ b/backend/app/plugin/module_ai/chat/ws.py @@ -83,8 +83,12 @@ async def websocket_chat_controller( chat_result = ChatService.chat_query(query=query) async for chunk in chat_result: if chunk: - await websocket.send_text(chunk) - full_response += chunk + try: + await websocket.send_text(chunk) + full_response += chunk + except RuntimeError: + log.warning("WebSocket连接已关闭,停止发送消息") + break # 保存AI回复到数据库(使用独立的事务) if query.session_id and full_response: @@ -106,17 +110,39 @@ async def websocket_chat_controller( log.warning(f"未提供会话ID或AI回复为空,跳过保存AI回复: session_id={query.session_id}, full_response_length={len(full_response)}") except json.JSONDecodeError: log.warning(f"收到非JSON消息: {data}") - await websocket.send_text("消息格式错误,请发送JSON格式的消息") + try: + await websocket.send_text("消息格式错误,请发送JSON格式的消息") + except RuntimeError: + log.warning("WebSocket连接已关闭,无法发送错误消息") + break except Exception as e: log.error(f"处理消息时出错: {e}") - await websocket.send_text(f"处理消息时出错: {str(e)}") + try: + await websocket.send_text(f"处理消息时出错: {str(e)}") + except RuntimeError: + log.warning("WebSocket连接已关闭,无法发送错误消息") + break except Exception as e: log.warning(f"WebSocket认证失败或聊天出错: {e}") - await websocket.send_text(f"错误: {str(e)}") - await websocket.close() + try: + await websocket.send_text(f"错误: {str(e)}") + except RuntimeError: + log.warning("WebSocket连接已关闭,无法发送错误消息") + finally: + try: + await websocket.close() + except RuntimeError: + pass return else: log.warning(f"WebSocket连接未提供token: {websocket.client}") - await websocket.send_text("未提供认证token,请重新登录") - await websocket.close() + try: + await websocket.send_text("未提供认证token,请重新登录") + except RuntimeError: + log.warning("WebSocket连接已关闭,无法发送错误消息") + finally: + try: + await websocket.close() + except RuntimeError: + pass return diff --git a/backend/app/plugin/module_task/node/model.py b/backend/app/plugin/module_task/node/model.py index 33895864..c4d0b458 100644 --- a/backend/app/plugin/module_task/node/model.py +++ b/backend/app/plugin/module_task/node/model.py @@ -1,20 +1,9 @@ -import enum - from sqlalchemy import Boolean, Integer, String, Text from sqlalchemy.orm import Mapped, mapped_column from app.core.base_model import ModelMixin, UserMixin -class NodeCategoryEnum(enum.Enum): - """节点分类枚举""" - - TRIGGER = "trigger" - ACTION = "action" - CONDITION = "condition" - CONTROL = "control" - - class NodeModel(ModelMixin, UserMixin): """ 节点类型模型 - 动态定义节点类型 @@ -26,7 +15,6 @@ class NodeModel(ModelMixin, UserMixin): name: Mapped[str] = mapped_column(String(64), nullable=False, comment="节点名称") code: Mapped[str] = mapped_column(String(32), nullable=False, unique=True, comment="节点编码") - category: Mapped[str] = mapped_column(String(32), default=NodeCategoryEnum.ACTION.value, comment="节点分类") jobstore: Mapped[str | None] = mapped_column(String(64), nullable=True, default="default", comment="存储器") executor: Mapped[str | None] = mapped_column(String(64), nullable=True, default="default", comment="执行器") trigger: Mapped[str | None] = mapped_column(String(64), nullable=True, comment="触发器") diff --git a/backend/app/plugin/module_task/node/schema.py b/backend/app/plugin/module_task/node/schema.py index da6c3b64..f88a6a68 100644 --- a/backend/app/plugin/module_task/node/schema.py +++ b/backend/app/plugin/module_task/node/schema.py @@ -28,7 +28,6 @@ class NodeCreateSchema(BaseModel): start_date: str | None = Field(default=None, description="开始时间") end_date: str | None = Field(default=None, description="结束时间") code: str | None = Field(default=None, description="节点编码") - category: str | None = Field(default=None, description="节点分类") @model_validator(mode="after") def _validate_func(self): diff --git a/backend/app/plugin/module_task/node/service.py b/backend/app/plugin/module_task/node/service.py index dae62224..c4feddc4 100644 --- a/backend/app/plugin/module_task/node/service.py +++ b/backend/app/plugin/module_task/node/service.py @@ -2,6 +2,7 @@ from app.api.v1.module_system.auth.schema import AuthSchema from app.core.ap_scheduler import SchedulerUtil from app.core.exceptions import CustomException from app.utils.cron_util import CronUtil +from apscheduler.jobstores.base import JobLookupError from .crud import NodeCRUD from .schema import ( @@ -40,7 +41,6 @@ class NodeService: "id": obj.id, "name": obj.name, "code": obj.code, - "category": obj.category, "func": obj.func, "args": obj.args, "kwargs": obj.kwargs, @@ -142,7 +142,11 @@ class NodeService: exist_obj = await NodeCRUD(auth).get_obj_by_id_crud(id=id) if not exist_obj: raise CustomException(msg="删除失败,该节点不存在") - SchedulerUtil.remove_job(job_id=id) + try: + SchedulerUtil.remove_job(job_id=id) + except JobLookupError: + # 作业不存在,忽略异常,继续删除数据库记录 + pass await NodeCRUD(auth).delete_obj_crud(ids=ids) @classmethod diff --git a/frontend/src/api/module_task/node.ts b/frontend/src/api/module_task/node.ts index 3b5651ff..8eeebe1b 100644 --- a/frontend/src/api/module_task/node.ts +++ b/frontend/src/api/module_task/node.ts @@ -70,7 +70,6 @@ export default NodeAPI; export interface NodePageQuery extends PageQuery { name?: string; code?: string; - category?: string; created_id?: number; updated_id?: number; created_time?: string[]; @@ -95,7 +94,6 @@ export interface ExecuteNodeResult { export interface NodeTable extends BaseType { name: string; code: string; - category?: string; jobstore?: string; executor?: string; trigger?: TriggerType; @@ -115,7 +113,6 @@ export interface NodeForm { id?: number; name: string; code?: string; - category?: string; jobstore?: string; executor?: string; func?: string; @@ -131,7 +128,6 @@ export interface NodeType { id: number; name: string; code: string; - category?: string; func?: string; args?: string; kwargs?: string; diff --git a/frontend/src/views/module_task/node/index.vue b/frontend/src/views/module_task/node/index.vue index 4405df7f..7f0ab440 100644 --- a/frontend/src/views/module_task/node/index.vue +++ b/frontend/src/views/module_task/node/index.vue @@ -14,19 +14,6 @@ - - - - - - - - - - - @@ -201,14 +181,6 @@ - - - - - - - - ({ page_size: 10, name: undefined, code: undefined, - category: undefined, }); const defaultCodeBlock = `def handler(*args, **kwargs) -> None: @@ -507,7 +480,6 @@ const formData = reactive({ id: undefined, name: "", code: undefined, - category: undefined, jobstore: "default", executor: "default", func: defaultCodeBlock, @@ -552,36 +524,6 @@ const executeRules = reactive({ trigger_args: [{ required: true, message: "请设置执行参数", trigger: "blur" }], }); -function getCategoryType(category: string | undefined) { - switch (category) { - case "trigger": - return "primary"; - case "action": - return "success"; - case "condition": - return "warning"; - case "control": - return "danger"; - default: - return "info"; - } -} - -function getCategoryLabel(category: string | undefined) { - switch (category) { - case "trigger": - return "触发器节点"; - case "action": - return "动作节点"; - case "condition": - return "条件节点"; - case "control": - return "控制节点"; - default: - return "未分类"; - } -} - async function handleRefresh() { await loadingData(); } @@ -614,14 +556,13 @@ const initialFormData: Partial = { id: undefined, name: "", code: undefined, - category: undefined, jobstore: "sqlalchemy", executor: "default", func: defaultCodeBlock, args: undefined, kwargs: undefined, coalesce: false, - max_instances: 1, + max_instances: 5, start_date: undefined, end_date: undefined, }; @@ -780,9 +721,40 @@ async function handleExecuteNode() { } await NodeAPI.executeNode(currentExecuteNode.value?.id as number, params); + ElMessage.success({ + message: `节点调试${executeFormData.trigger === "now" ? "已启动" : "已创建"}`, + type: "success", + duration: 2000, + }); + handleCloseExecuteDialog(); - loadingData(); + + // 重新加载数据 + await loadingData(); + + // 如果是立即执行,提示用户查看执行记录 + if (executeFormData.trigger === "now") { + ElMessageBox.confirm("调试任务已启动,是否跳转到执行记录页面查看执行结果?", "提示", { + confirmButtonText: "查看记录", + cancelButtonText: "稍后查看", + type: "info", + }) + .then(() => { + // 跳转到执行记录页面 + router.push({ + path: "/task/job", + }); + }) + .catch(() => { + // 取消操作 + }); + } } catch (error: any) { + ElMessage.error({ + message: error.response?.data?.msg || "调试失败", + type: "error", + duration: 3000, + }); console.error(error); } finally { loading.value = false;