# ============================================================================= # 企微IT智能服务台 — 阶段5 自动化 动作执行引擎 # ============================================================================= # 说明:按场景动作计划顺序执行,落实「分级执行」策略: # - read / low 风险 + real_exec 模式 → 自动执行 # - high 风险 / plan_only 模式 → 挂起并生成审批单(坐席审批 或 员工 H5 确认) # 执行中任一动作失败 → 触发回滚补偿 + 转人工接管。 # 全部成功 → 处置成功(resolved),调度静默关单。 # # 设计要点: # - 主循环 run() 在遇到首个「需审批」动作时挂起并返回,等待 resume() 续跑 # - resume() 更新审批单后再次调用 run() 继续后续动作(支持多步审批) # ============================================================================= from __future__ import annotations import logging from datetime import datetime, timezone from typing import Any, Dict, List, Optional from sqlalchemy import select from app.constants import ( AUTOMATION_AUTO_EXECUTABLE_RISKS, AUTOMATION_SESSION_TERMINAL_STATES, AutomationErrorCode, ) from app.database import _get_session_factory from app.integrations.base import BaseClientError from app.integrations.factory import ( build_dify_client, build_ehr_client, build_huorong_client, build_lianruan_client, ) from app.models.automation import AutoAction, AutoSession from app.services.automation.action_registry import ActionContext, get_handler from app.services.automation.approval import ApprovalService from app.services.automation.exception_handler import AutomationException from app.services.automation.progress_publisher import ( cancel_silent_close, publish_progress, publish_resolved, publish_takeover, schedule_silent_close, ) from app.services.automation.rollback import RollbackService logger = logging.getLogger(__name__) # 静默关单默认 TTL(秒):处置成功后 10 分钟内员工无异议自动关单 SILENT_CLOSE_TTL = 600 class ActionExecutor: """动作执行引擎。""" def __init__(self, db: Any, redis: Any = None, audit: Any = None): self.db = db self.redis = redis self.audit = audit self.approval = ApprovalService(db) self.rollback = RollbackService(db) # -------------------------------------------------------------------------- # 数据加载 # -------------------------------------------------------------------------- async def _load(self, session_id: str): """加载会话及其动作(按动作顺序排序)。""" stmt = select(AutoSession).where(AutoSession.id == session_id) session = (await self.db.execute(stmt)).scalar_one_or_none() if session is None: return None, [] act_stmt = select(AutoAction).where(AutoAction.session_id == session_id) actions = list((await self.db.execute(act_stmt)).scalars().all()) actions.sort(key=lambda a: a.action_index) return session, actions async def _build_clients(self) -> Dict[str, Any]: """构建外部客户端(缺失时置 None,不阻断主流程)。""" clients: Dict[str, Any] = {} for name, builder in ( ("huorong", build_huorong_client), ("lianruan", build_lianruan_client), ("ehr", build_ehr_client), ("dify", build_dify_client), ): try: clients[name] = await builder(self.db, audit=self.audit) if name != "dify" else await builder(audit=self.audit) except Exception as e: # noqa: BLE001 logger.warning(f"构建外部客户端失败 {name}: {e}") clients[name] = None return clients # -------------------------------------------------------------------------- # 主循环 # -------------------------------------------------------------------------- async def run(self, session_id: str) -> None: """执行会话动作计划(遇到审批闸门挂起返回,等待 resume)。""" session, actions = await self._load(session_id) if session is None: logger.warning(f"executor.run 会话不存在: {session_id}") return if session.status in AUTOMATION_SESSION_TERMINAL_STATES: logger.info(f"会话已终态,跳过执行: {session_id} status={session.status}") return mapping = (session.meta or {}).get("mapping") or {} clients = await self._build_clients() for action in actions: if action.status in ("success", "failed", "rejected", "skipped"): continue # 已批准 → 直接执行 if action.status == "approved": await self._do_execute(session, action, mapping, clients) continue # 待决(pending / await_approval)→ 判断是否需审批闸门 needs_gate = ( action.risk_level not in AUTOMATION_AUTO_EXECUTABLE_RISKS ) or (session.mode == "plan_only") if needs_gate: channel = (action.payload or {}).get("confirm_channel") or "agent" if session.mode == "plan_only": channel = "agent" # 方案预览阶段统一走坐席确认 ticket = await self.approval.ensure_ticket( action, channel=channel, reason=action.description ) action.status = "await_approval" session.status = "paused" session.current_action_id = action.id await self.db.flush() await publish_action_required(session.id, action, ticket) return # 暂停,等待审批/确认后续跑 # 可自动执行 await self._do_execute(session, action, mapping, clients) # 全部动作处理完毕 await self._finish_success(session) async def _do_execute( self, session: AutoSession, action: AutoAction, mapping: dict, clients: dict ) -> None: """执行单个动作,处理成功/失败分支。""" handler = get_handler(action.action_type) if handler is None: action.status = "failed" action.error = f"未注册的动作类型: {action.action_type}" await self.db.flush() await self._on_action_failed( session, action, AutomationException(AutomationErrorCode.CONFIG_ERROR, action.error), ) return ctx = ActionContext( db=self.db, session=session, action=action, mapping=mapping, clients=clients, employee_id=session.employee_id, audit=self.audit, ) try: result = await handler.execute(ctx) action.status = "success" action.result = result await self.db.flush() await publish_progress( session.id, "action_done", f"动作完成:{action.title}", action.id ) except BaseClientError as e: action.status = "failed" action.error = str(e) await self.db.flush() await self._on_action_failed( session, action, AutomationException(AutomationErrorCode.EXTERNAL_CALL_FAILED, str(e)), ) except AutomationException as e: action.status = "failed" action.error = e.message await self.db.flush() await self._on_action_failed(session, action, e) except Exception as e: # noqa: BLE001 action.status = "failed" action.error = str(e) await self.db.flush() await self._on_action_failed( session, action, AutomationException(AutomationErrorCode.EXTERNAL_CALL_FAILED, str(e)), ) async def _on_action_failed( self, session: AutoSession, action: AutoAction, exc: AutomationException ) -> None: """动作失败:回滚已执行动作 + 转人工接管。""" try: await self.rollback.compensate(session, action) except Exception as e: # noqa: BLE001 logger.warning(f"回滚补偿异常 session={session.id}: {e}") session.status = "handoff" session.closed_by = "system(auto)" await self.db.flush() await publish_error(session.id, exc.code, exc.message) await publish_takeover(session.id, f"动作失败转人工:{exc.message}") async def _finish_success(self, session: AutoSession) -> None: """全部动作成功:置 resolved + 调度静默关单。""" session.status = "resolved" session.resolved_at = datetime.now(timezone.utc) session.current_action_id = None await self.db.flush() await publish_resolved(session.id, "处置已完成,等待您确认") schedule_silent_close(session.id, SILENT_CLOSE_TTL, on_expire=self._auto_close) # -------------------------------------------------------------------------- # 审批/确认后续跑 # -------------------------------------------------------------------------- async def resume( self, session_id: str, action_id: str, decision: str, note: Optional[str], approver_id: Optional[str], ) -> None: """审批/确认结果回来后,更新审批单并续跑计划。""" session, actions = await self._load(session_id) if session is None: return if session.status in AUTOMATION_SESSION_TERMINAL_STATES: return ticket = await self.approval.get_pending_for_action(action_id) if ticket is not None: await self.approval.decide(ticket.id, decision, note, approver_id) action = next((a for a in actions if a.id == action_id), None) if action is None: return if decision == "approve": action.status = "approved" action.approved_by = approver_id action.approved_at = datetime.now(timezone.utc) await self.db.flush() # 续跑主循环(会继续执行已批准动作及后续动作) await self.run(session_id) else: action.status = "rejected" session.status = "handoff" session.closed_by = approver_id await self.db.flush() cancel_silent_close(session_id) await publish_takeover( session.id, f"审批驳回转人工:{note or action.title}" ) # -------------------------------------------------------------------------- # 静默关单回调(独立会话,避免长事务) # -------------------------------------------------------------------------- async def _auto_close(self, session_id: str) -> None: """静默关单:仅在会话仍为 resolved 时自动关单。""" factory = _get_session_factory() async with factory() as db: stmt = select(AutoSession).where(AutoSession.id == session_id) session = (await db.execute(stmt)).scalar_one_or_none() if session is None: return if session.status != "resolved": return # 已被接管/关单/反馈 session.status = "closed" session.closed_by = "system(auto)" await db.commit() await publish_progress(session_id, "auto_closed", "已静默关单")