- 新增 POST /api/cma/bot-bridge/kpi-result: (entity_id,kpi_code)查KPI, 写入KPIValue(source_type=bot, data_status=verified), 复用run_alert_check 触发预警(自动联动行动计划), 返回new_alerts; KPI不存在/缺字段返回错误 - 修复verify引擎_get_kpi_current_value排序: 按calculated_at取最新值, bot回填值可被auto-verify读到 (修复2026H1字符串排序遮蔽月值问题) - 新增scripts/bot_bridge_push.py: 财务Bot MPM结果自动回填(MPM字段→KPI 映射), 支持--mpm/--mpm-file/单KPI直推/--ping
619 lines
20 KiB
Python
619 lines
20 KiB
Python
"""
|
||
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 <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("/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,
|
||
}
|