feat: bot-bridge+verify合并 — API+数据库+验证引擎

This commit is contained in:
Hermes CI Fix
2026-07-22 17:43:59 +08:00
parent e8ffa0d550
commit 2d0abe3411
5 changed files with 561 additions and 2 deletions
+501
View File
@@ -0,0 +1,501 @@
"""
Bot-Bridge V2 — 合并改造:Bot桥接 + 自动验证引擎
数据流:
bot-bridge(财务Bot→CMA)→ 写入KPI
verifyCMA→校验)→ 读取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.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 <op> 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] <op> 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("/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,
}