""" Bot-Bridge V2 — 合并改造:Bot桥接 + 自动验证引擎 数据流: bot-bridge(财务Bot→CMA)→ 写入KPI verify(CMA→校验)→ 读取KPI → 通过/失败 API: POST /api/cma/bot-bridge/mpm-result — 接收MPM分析结果 POST /api/cma/bot-bridge/verify/{action_plan_id} — 验证ActionPlan执行结果 """ import json, logging, re from datetime import datetime from typing import Optional from fastapi import APIRouter, Depends, HTTPException, Header from sqlalchemy.orm import Session from sqlalchemy import func from app.database import get_db from app.models import ( MpmResult, BotBridgeConfig, KPIDefinition, KPIValue, KPIAlert, ActionPlan, Entity, ) logger = logging.getLogger("cma.bot_bridge_v2") router = APIRouter(prefix="/api/cma/bot-bridge", tags=["Bot桥接V2"]) # ═══════════════════════════════════════════════ # 鉴权 # ═══════════════════════════════════════════════ def verify_bridge_token( x_bridge_token: str = Header(None, alias="X-BRIDGE-TOKEN"), db: Session = Depends(get_db), ): """从bot_bridge_config表校验Token""" if not x_bridge_token: raise HTTPException(401, "缺少X-BRIDGE-TOKEN请求头") config = db.query(BotBridgeConfig).filter( BotBridgeConfig.token == x_bridge_token, BotBridgeConfig.is_active == True, ).first() if not config: raise HTTPException(401, "Token无效或已停用") return config.bot_name # ═══════════════════════════════════════════════ # 验证引擎 # ═══════════════════════════════════════════════ def _get_kpi_current_value(db: Session, kpi_id: int, entity_id: int = None) -> Optional[dict]: """获取KPI当前最新值及目标值""" kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first() if not kpi: return None latest = db.query(KPIValue).filter( KPIValue.kpi_id == kpi_id, KPIValue.data_status == "verified", ).order_by(KPIValue.calculated_at.desc(), KPIValue.period.desc()).first() return { "kpi_id": kpi.id, "kpi_code": kpi.kpi_code, "kpi_name": kpi.kpi_name, "value": latest.actual_value if latest else None, "target": kpi.target_value, "period": latest.period if latest else None, "unit": kpi.unit, } def _parse_condition(condition: str) -> list: """ 解析condition表达式为可执行结构 支持格式: "value > target" — 当前值 > 目标值 "value >= 80" — 当前值 >= 80 (绝对值) "value / target * 100 > 80" — 完成率 > 80% "value < baseline" — 当前值 < 基线值(目标值=基线) "value == 0" — 精确等于 "value > target AND value < 100" — 多条件AND "value > 50 OR value == -1" — 多条件OR """ if not condition: return [] # 分割AND/OR parts = re.split(r'\s+(AND|OR)\s+', condition, flags=re.IGNORECASE) clauses = [] for i in range(0, len(parts), 2): expr = parts[i].strip() # 检查是否有运算符连接 match = re.match( r'(value|target|baseline)\s*([><=!]+)\s*(value|target|baseline|-?\d+\.?\d*)(?:\s*([+\-*/])\s*(\d+\.?\d*))?', expr ) if not match: # 尝试 value/target*100 > 80 这种复合表达式 match = re.match( r'(value)\s*(/|\\*)\s*(target|baseline)\s*(\*?\s*\d+)?\s*([><=!]+)\s*(-?\d+\.?\d*)', expr ) if not match: logger.warning(f"无法解析condition表达式: {expr}") continue logical_op = parts[i + 1].upper() if i + 1 < len(parts) else "AND" clauses.append({"expr": expr, "logical_op": logical_op, "match": match.groups() if match else ()}) return clauses def _evaluate_condition(condition: str, kpi_data: dict) -> bool: """ 评估condition表达式 Args: condition: 条件表达式 e.g. "value > target" kpi_data: {value, target, ...} Returns: True=通过, False=不通过 """ if not condition: return False value = kpi_data.get("value") target = kpi_data.get("target") baseline = kpi_data.get("baseline", target) if value is None: logger.warning(f"KPI {kpi_data.get('kpi_code')} 当前值为空,无法验证") return False result = True current_logical = "AND" # 按AND/OR分割 parts = re.split(r'\s+(AND|OR)\s+', condition, flags=re.IGNORECASE) for i in range(0, len(parts), 2): expr = parts[i].strip() if i + 1 < len(parts): current_logical = parts[i + 1].upper() clause_passed = _eval_single_expr(expr, value, target, baseline) if current_logical == "AND": result = result and clause_passed elif current_logical == "OR": result = result or clause_passed # 短路优化 if current_logical == "AND" and not result: break if current_logical == "OR" and result: break return result def _eval_single_expr(expr: str, value: float, target: Optional[float], baseline: Optional[float]) -> bool: """评估单条条件表达式""" # 模式1: value target/baseline/number (如 value > target, value >= 80) m = re.match( r'(value|target|baseline)\s*([><=!]+)\s*(value|target|baseline|-?\d+\.?\d*)', expr ) if m: left = m.group(1) op = m.group(2) right_raw = m.group(3) left_val = _resolve_var(left, value, target, baseline) right_val = _resolve_var(right_raw, value, target, baseline) if left_val is None or right_val is None: return False return _apply_op(left_val, op, right_val) # 模式2: value [*/] target/baseline [* number] number (如 value/target*100 > 80) m = re.match( r'(value)\s*([/*])\s*(target|baseline)(?:\s*\*\s*(\d+))?\s*([><=!]+)\s*(-?\d+\.?\d*)', expr ) if m: left = m.group(1) op1 = m.group(2) # / or * var2 = m.group(3) # target or baseline multiplier = float(m.group(4)) if m.group(4) else 100 op2 = m.group(5) # > < >= <= == != right_num = float(m.group(6)) left_val = _resolve_var(left, value, target, baseline) right_var = _resolve_var(var2, value, target, baseline) if left_val is None or right_var is None or right_var == 0: return False if op1 == '/': computed = (left_val / right_var) * multiplier else: # * computed = left_val * right_var return _apply_op(computed, op2, right_num) # 模式3: 纯数字比较 (提供兼容) logger.warning(f"无法解析表达式: {expr}") return False def _resolve_var(token: str, value: float, target: Optional[float], baseline: Optional[float]) -> Optional[float]: """将变量名解析为数值""" token = token.strip() if token == "value": return value elif token == "target": return target elif token == "baseline": return baseline else: try: return float(token) except (ValueError, TypeError): return None def _apply_op(left: float, op: str, right: float) -> bool: """应用比较运算符""" try: if op == ">": return left > right elif op == ">=": return left >= right elif op == "<": return left < right elif op == "<=": return left <= right elif op in ("==", "="): return abs(left - right) < 0.0001 elif op == "!=": return abs(left - right) >= 0.0001 else: logger.warning(f"未知运算符: {op}") return False except (TypeError, ValueError): return False # ═══════════════════════════════════════════════ # KPI映射表:MPM结果字段 → KPI编码 # ═══════════════════════════════════════════════ MPM_TO_KPI_MAP = { "revenue": None, # 不做KPI映射,保留在raw_data "cost": None, "gross_margin_standard": None, "gross_margin_adjusted": "F_GROSS_MARGIN", "net_profit_standard": None, "net_profit_adjusted": "F_NET_PROFIT", "channel_rebate_rate": "C_REBATE_RATE", "mgmt_expense_ratio": "F_COST_RATIO", } # 预警阈值配置 ALERT_THRESHOLDS = { "channel_rebate_rate": {"threshold": 75, "operator": ">", "kpi_code": "C_REBATE_RATE", "kpi_name": "渠补率"}, "mgmt_expense_ratio": {"threshold": 50, "operator": ">", "kpi_code": "F_COST_RATIO", "kpi_name": "管理费/净收入"}, } # ═══════════════════════════════════════════════ # API端点 # ═══════════════════════════════════════════════ @router.post("/mpm-result") def receive_mpm_result( data: dict, bridge_bot: str = Depends(verify_bridge_token), db: Session = Depends(get_db), ): """ 接收财务Bot/其他Bot的MPM分析结果 - 鉴权(校验X-BRIDGE-TOKEN) - 写入mpm_results表 - 更新KPI当前值(kpi_values表) - 触发预警(如渠补率超75%) """ source = data.get("source", bridge_bot) entity_id = data.get("entity_id") period = data.get("period") results = data.get("results", {}) if not entity_id or not period: raise HTTPException(400, "缺少必填字段: entity_id, period") # 验证entity存在 entity = db.query(Entity).filter(Entity.id == entity_id).first() if not entity: raise HTTPException(404, f"实体entity_id={entity_id}不存在") # ── 1. 写入MPM结果 ── record = MpmResult( entity_id=entity_id, period=period, source=source, raw_data=results, ) db.add(record) db.flush() # ── 2. 更新KPI当前值 ── kpi_updates = [] for field, kpi_code in MPM_TO_KPI_MAP.items(): if kpi_code is None: continue field_value = results.get(field) if field_value is None: continue kpi = db.query(KPIDefinition).filter( KPIDefinition.kpi_code == kpi_code, KPIDefinition.entity_id == entity_id, ).first() if not kpi: logger.warning(f"KPI编码 {kpi_code} 未找到 (entity_id={entity_id})") continue # 写入最新值 (upsert: 存在同period则更新,否则插入) existing = db.query(KPIValue).filter( KPIValue.kpi_id == kpi.id, KPIValue.period == period, ).first() if existing: existing.actual_value = field_value existing.source_type = "bot-bridge-v2" existing.data_status = "verified" else: kpi_val = KPIValue( kpi_id=kpi.id, period=period, actual_value=field_value, source_type="bot-bridge-v2", data_status="verified", ) db.add(kpi_val) kpi_updates.append(kpi_code) # ── 3. 触发预警 ── alerts = [] for field, config in ALERT_THRESHOLDS.items(): field_value = results.get(field) if field_value is None: continue threshold = config["threshold"] operator = config["operator"] kpi_code = config.get("kpi_code") kpi_name = config["kpi_name"] if (operator == ">" and field_value > threshold) or \ (operator == ">=" and field_value >= threshold): # 查找关联KPI kpi_id = None if kpi_code: kpi_def = db.query(KPIDefinition).filter( KPIDefinition.kpi_code == kpi_code, KPIDefinition.entity_id == entity_id, ).first() if kpi_def: kpi_id = kpi_def.id # 创建预警 alert_msg = f"⚠️ {kpi_name}异常: {field_value:.1f}% (阈值: {threshold}%) [entity={entity.short_name or entity.name}, period={period}]" alert = KPIAlert( kpi_id=kpi_id or 1, # fallback to first KPI if not found alert_level="red", alert_message=alert_msg, alert_type="actual", status="pending", ) db.add(alert) alerts.append(alert_msg) logger.warning(alert_msg) db.commit() db.refresh(record) return { "success": True, "kpi_updated": len(kpi_updates), "alerts_triggered": len(alerts), "mpm_record_id": record.id, "details": { "kpi_codes": kpi_updates, "alerts": alerts, }, } @router.post("/kpi-result") def push_kpi_result( data: dict, bridge_bot: str = Depends(verify_bridge_token), db: Session = Depends(get_db), ): """ Bot分析结果回填KPI值 + 触发预警(PRD第七部分 bot-bridge数据通道) 入参: entity_id (int, 默认1): 企业实体ID kpi_code (str, 必填): KPI编码 period (str, 必填): 期间 YYYY-MM value (num, 必填): 实际值 source (str, 默认finance-bot): 来源Bot标识 remark (str, 可选): 备注 行为: 1. 按(entity_id, kpi_code)查KPI → 不存在返回错误 2. 写入KPIValue (source_type=bot, data_status=verified) — 与auto-verify引擎兼容 3. 调用 run_alert_check 触发阈值预警 → 返回 new_alerts """ entity_id = data.get("entity_id", 1) kpi_code = data.get("kpi_code") period = data.get("period") value = data.get("value") source = data.get("source", "finance-bot") remark = data.get("remark") if not kpi_code: raise HTTPException(400, "缺少必填字段: kpi_code") if not period: raise HTTPException(400, "缺少必填字段: period") if value is None or value == "": raise HTTPException(400, "缺少必填字段: value") try: actual_value = float(value) except (TypeError, ValueError): raise HTTPException(400, f"value不是有效数值: {value!r}") # 1. 按(entity_id, kpi_code)查KPI kpi = db.query(KPIDefinition).filter( KPIDefinition.entity_id == entity_id, KPIDefinition.kpi_code == kpi_code, ).first() if not kpi: raise HTTPException(404, f"KPI {kpi_code} 不存在 (entity_id={entity_id})") # 2. 写入KPIValue(同period已存在则更新,幂等upsert) existing = db.query(KPIValue).filter( KPIValue.kpi_id == kpi.id, KPIValue.period == period, ).first() if existing: existing.actual_value = actual_value existing.source_type = "bot" existing.source_batch = source existing.data_status = "verified" existing.calculated_at = datetime.now() if remark: existing.remark = remark val = existing else: val = KPIValue( kpi_id=kpi.id, period=period, actual_value=actual_value, source_type="bot", source_batch=source, data_status="verified", calculated_at=datetime.now(), remark=remark, ) db.add(val) db.commit() db.refresh(val) # 3. 触发预警检查(复用alert_generator引擎,yellow/red自动联动行动计划) from scripts.alert_generator import run_alert_check new_count = run_alert_check(db, period) # 收集本次写入值直接触发的预警(kpi_value_id关联) new_alerts = [] if new_count: triggered = db.query(KPIAlert).filter( KPIAlert.kpi_value_id == val.id, ).order_by(KPIAlert.created_at.desc()).all() new_alerts = [ { "id": a.id, "kpi_id": a.kpi_id, "kpi_code": kpi.kpi_code, "level": a.alert_level, "message": a.alert_message, "action_plan_id": a.action_plan_linked_id, "created_at": a.created_at.isoformat() if a.created_at else None, } for a in triggered ] logger.info( f"[bot-bridge] KPI回填: {kpi.kpi_code}@{period}={actual_value} " f"source={source} entity={entity_id} alerts={new_count}" ) return { "status": "ok", "kpi_value_id": val.id, "kpi_code": kpi.kpi_code, "kpi_name": kpi.kpi_name, "period": period, "value": actual_value, "source": source, "new_alerts_count": len(new_alerts), "new_alerts": new_alerts, } @router.post("/verify/{action_plan_id}") def verify_action_plan( action_plan_id: int, bridge_bot: str = Depends(verify_bridge_token), db: Session = Depends(get_db), ): """ 验证ActionPlan的执行结果 1. 读取ActionPlan的auto_verify_rule 2. 读取关联KPI的当前值 3. 按condition校验 4. 返回通过/失败 + 详细数据 """ plan = db.query(ActionPlan).filter(ActionPlan.id == action_plan_id).first() if not plan: raise HTTPException(404, f"ActionPlan {action_plan_id} 不存在") if not plan.auto_verify_rule: raise HTTPException(400, "该ActionPlan未配置auto_verify_rule验证规则") # 解析验证规则 if isinstance(plan.auto_verify_rule, str): rule = json.loads(plan.auto_verify_rule) else: rule = plan.auto_verify_rule condition = rule.get("condition", "") description = rule.get("description", "") # 获取KPI当前值 kpi_data = _get_kpi_current_value(db, plan.kpi_id) if not kpi_data: return { "success": False, "action_plan_id": action_plan_id, "title": plan.title, "verify_result": "fail", "reason": "关联KPI不存在", } # 执行验证 passed = _evaluate_condition(condition, kpi_data) # 记录验证日志 verify_log_entry = { "timestamp": datetime.now().isoformat(), "condition": condition, "kpi_data": kpi_data, "passed": passed, } existing_logs = plan.verify_log or [] if isinstance(existing_logs, list): existing_logs.append(verify_log_entry) else: existing_logs = [verify_log_entry] plan.verify_result = "pass" if passed else "fail" plan.verify_log = existing_logs db.commit() return { "success": True, "action_plan_id": action_plan_id, "title": plan.title, "verify_result": plan.verify_result, "condition": condition, "condition_description": description, "kpi_data": kpi_data, "verification_detail": { "value": kpi_data.get("value"), "target": kpi_data.get("target"), "condition": condition, "passed": passed, }, } @router.get("/verify/{action_plan_id}/history") def verify_history( action_plan_id: int, bridge_bot: str = Depends(verify_bridge_token), db: Session = Depends(get_db), ): """获取ActionPlan的验证历史""" plan = db.query(ActionPlan).filter(ActionPlan.id == action_plan_id).first() if not plan: raise HTTPException(404, f"ActionPlan {action_plan_id} 不存在") return { "action_plan_id": action_plan_id, "title": plan.title, "verify_result": plan.verify_result, "verify_log": plan.verify_log or [], "auto_verify_rule": plan.auto_verify_rule, }