Files
cma-management/backend/app/api/ai_analysis.py
T
Hermes CI Fix ad68471b29 feat: 路线图R1决策建议一键落地+R2机会推送+R5预算闭环
R1(P0): AI建议一键应用到KPI/预算/行动方案
- 新表 ai_suggestions + AISuggestion 模型(init_db自动建)
- /api/cma/ai/suggestions CRUD + /{id}/apply(复用kpis/budget/action_plans) + dismiss
- 应用写 OperationLog(action=ai_suggestion_apply, detail含suggestion_id/before/after)
- 规则驱动建议生成 generate_rule_suggestions(低执行率/高执行率/预算超支/pending预警)
- 幂等: 同entity+type+target_id+title+unapplied不重复建; applied后拒绝重复应用
- 前端: Dashboard AI面板建议卡(应用到/忽略) + 建议中心页 /ai-suggestions

R2(P1): 数据找人扩大-机会类推送
- scripts/opportunity_detector.py: KPI向好(执行率>110%)/预算余量(<70%且actual>0)/预测上行
- scripts/daily_push.py: 异常+机会 每日9:15推企微(8800/send, --dry-run调试)
- crontab: 15 9 * * * (alert_generator 9:00之后)

R5(P0): 预算闭环加固
- auto-decompose批量幂等: 只取年度行(period=YYYY-00)+同KPI多版本取一行
- scripts/closed_loop_check.py: 预算执行率异常→检查现金流/行动同步→缺失提示+报告
- scripts/verify_decompose_idempotent.py: 幂等验证脚本

测试: test_ai_suggestions(10例)+test_roadmap_r2r5(14例); 修test_budget幂等契约适配年度行
全量: 673 passed
2026-08-30 12:05:36 +08:00

517 lines
21 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""AI分析引擎 — 侧边栏智能分析"""
from fastapi import APIRouter, Depends, HTTPException, Query, Request
from fastapi.responses import StreamingResponse
from sqlalchemy.orm import Session
from sqlalchemy import func, text as sa_text, or_
from app.database import get_db
from app.deps import get_entity_id
from app.auth_middleware import require_auth, require_role
from app.models import KPIDefinition, KPIValue, KPIAlert, StrategicMap, User, ActionPlan, BudgetPlan, AISuggestion
from app.utils.cache import get as cache_get, set as cache_set
import json, hashlib, httpx, os
from datetime import datetime, date
router = APIRouter(prefix="/api/cma/ai", tags=["AI分析"],
dependencies=[Depends(require_role("ceo", "finance", "business", "it"))],
)
# ============================================================
# R1 决策建议生成(规则驱动,稳定可复现,落库 ai_suggestions
# ============================================================
def _sug_dict(s: AISuggestion) -> dict:
return {
"id": s.id,
"entity_id": s.entity_id,
"source": s.source,
"suggestion_type": s.suggestion_type,
"target_type": s.target_type,
"target_id": s.target_id,
"title": s.title,
"content": s.content,
"suggestion_data": s.suggestion_data or {},
"status": s.status,
"applied_by": s.applied_by,
"applied_at": s.applied_at.isoformat() if s.applied_at else None,
"apply_detail": s.apply_detail or [],
"created_at": s.created_at.isoformat() if s.created_at else None,
}
def _existing_unapplied(db: Session, entity_id: int, suggestion_type: str,
target_id: int, title: str) -> bool:
"""幂等:同entity+类型+目标+标题的未应用建议存在则跳过"""
return db.query(AISuggestion).filter(
AISuggestion.entity_id == entity_id,
AISuggestion.suggestion_type == suggestion_type,
AISuggestion.target_id == target_id,
AISuggestion.title == title,
AISuggestion.status == "unapplied",
).first() is not None
def generate_rule_suggestions(db: Session, entity_id: int,
source: str = "dashboard", user_id: int = None,
kpi_id: int = None) -> list:
"""从数据规则生成决策建议并落库(R1,路线图2026-08-30
规则:
1. KPI执行率<70% → 建议建行动方案(异常类)
2. KPI执行率>110% → 建议上调KPI目标(机会类)
3. 预算执行率>110% → 建议调预算(预算类)
4. 有pending预警 → 建议建行动方案处理预警
幂等:同 entity+type+target_id+title+status=unapplied 不重复建。
"""
now = datetime.now()
period = now.strftime("%Y-%m")
created = []
def _add(suggestion_type: str, target_type: str, tid: int,
title: str, content: str, suggestion_data: dict):
nonlocal created
if _existing_unapplied(db, entity_id, suggestion_type, tid, title):
return
sug = AISuggestion(
entity_id=entity_id,
user_id=user_id,
source=source,
suggestion_type=suggestion_type,
target_type=target_type,
target_id=tid,
title=title,
content=content,
suggestion_data=suggestion_data,
status="unapplied",
)
db.add(sug)
created.append(sug)
# 查询KPI(可按kpi_id过滤)
q = db.query(KPIDefinition).filter(KPIDefinition.entity_id == entity_id,
KPIDefinition.status == "active")
if kpi_id:
q = q.filter(KPIDefinition.id == kpi_id)
kpis = q.all()
for k in kpis:
latest = db.query(KPIValue).filter(
KPIValue.kpi_id == k.id,
or_(
KPIValue.entity_id == entity_id,
KPIValue.entity_id.is_(None),
),
).order_by(KPIValue.period.desc()).first()
if not latest or latest.actual_value is None:
continue
actual = latest.actual_value
target = k.target_value
ratio = (actual / target) if target else None
# 1. 异常:执行率<70% → 建行动方案
if ratio is not None and ratio < 0.7:
title = f"提升 {k.kpi_name}:达成率仅{ratio*100:.0f}%"
content = (f"KPI[{k.kpi_name}] 最新期间{latest.period}实际值{actual:g}"
f"目标{target:g},达成率{ratio*100:.1f}%,低于70%预警线。"
f"建议制定专项改善行动方案。")
_add("action_plan", "kpi", k.id, title, content, {
"kpi_id": k.id, "priority": "high",
"title": f"改善: {k.kpi_name}达成率提升",
})
# 2. 机会:执行率>110% → 上调KPI目标
elif ratio is not None and ratio > 1.1:
new_target = round(actual * 1.05, 2)
title = f"上调 {k.kpi_name} 目标:达成率{ratio*100:.0f}%超预期"
content = (f"KPI[{k.kpi_name}] 达成率{ratio*100:.1f}%超过110%"
f"建议将目标从{target:g}上调至{new_target:g},保持牵引力。")
_add("kpi_target", "kpi", k.id, title, content, {
"kpi_id": k.id, "target_value": new_target,
})
# 3. 预算执行率>110% → 调预算
budget_rows = db.query(BudgetPlan).filter(
BudgetPlan.entity_id == entity_id,
BudgetPlan.status == "active",
BudgetPlan.period == period,
).all()
for b in budget_rows:
actual = db.query(func.max(KPIValue.actual_value)).filter(
KPIValue.kpi_id == b.kpi_id,
KPIValue.period == b.period,
).scalar()
if actual is None or b.budget_value is None or b.budget_value <= 0:
continue
exec_ratio = actual / b.budget_value
if exec_ratio > 1.1:
kpi_name = "KPI"
k = db.query(KPIDefinition).filter(KPIDefinition.id == b.kpi_id).first()
if k:
kpi_name = k.kpi_name
title = f"调整 {kpi_name} 预算:执行率{exec_ratio*100:.0f}%超预算"
content = (f"预算[{kpi_name}] {period}预算值{b.budget_value:g}"
f"实际{actual:g},执行率{exec_ratio*100:.1f}%超过110%。"
f"建议同步调整预算/现金流/行动方案。")
_add("budget_adjust", "budget", b.kpi_id, title, content, {
"kpi_id": b.kpi_id, "period": period, "budget_value": round(actual, 2),
})
# 4. pending预警 → 建行动方案
alerts = db.query(KPIAlert).filter(KPIAlert.status == "pending").all()
for a in alerts:
k = db.query(KPIDefinition).filter(KPIDefinition.id == a.kpi_id).first()
kpi_name = k.kpi_name if k else f"KPI#{a.kpi_id}"
title = f"处理预警:{kpi_name} {a.alert_message[:30]}"
content = f"存在待处理预警({a.alert_level}级):{a.alert_message}。建议建立行动方案跟进。"
_add("action_plan", "alert", a.id, title, content, {
"kpi_id": a.kpi_id, "priority": "high" if a.alert_level == "red" else "medium",
"alert_id": a.id,
"title": f"处理预警: {kpi_name}",
})
if created:
db.commit()
for s in created:
db.refresh(s)
return created
def _unapplied_suggestions(db: Session, entity_id: int, limit: int = 20) -> list:
items = db.query(AISuggestion).filter(
AISuggestion.entity_id == entity_id,
AISuggestion.status == "unapplied",
).order_by(AISuggestion.created_at.desc()).limit(limit).all()
return [_sug_dict(s) for s in items]
async def _call_deepseek(prompt: str) -> str:
"""调用DeepSeek API"""
api_key = os.getenv("DEEPSEEK_API_KEY", "sk-8e24e6eb87f2475e96ea0980002dc2e8")
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(
"https://api.deepseek.com/v1/chat/completions",
headers={"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"},
json={
"model": "deepseek-chat",
"messages": [
{"role": "system", "content": "你是一名CMA管理会计师,擅长用数据驱动的方式分析企业经营状况,给出专业的财务分析和管理建议。回答要简洁、专业、有数据支撑。"},
{"role": "user", "content": prompt}
],
"stream": False,
"temperature": 0.3,
}
)
data = resp.json()
return data.get("choices", [{}])[0].get("message", {}).get("content", "")
@router.get("/dashboard-analysis")
async def dashboard_analysis(role: str = Query("ceo"), db: Session = Depends(get_db),
entity_id: int = Depends(get_entity_id)):
"""AI分析驾驶舱数据"""
# 尝试缓存
cache_key = f"dashboard_analysis:{role}:{entity_id}"
cached = cache_get("ai", cache_key)
if cached:
# 缓存命中(LLM文本10分钟内不重复调用),但轻量规则建议仍执行(幂等)
try:
generate_rule_suggestions(db, entity_id, source="dashboard")
except Exception:
pass
cached["suggestions"] = _unapplied_suggestions(db, entity_id)
return cached
# 获取当前KPI数据
kpis = db.query(KPIDefinition).filter(KPIDefinition.entity_id == entity_id,
KPIDefinition.status == "active").all()
kpi_summary = []
for k in kpis:
latest = db.query(KPIValue).filter(KPIValue.kpi_id == k.id).order_by(KPIValue.period.desc()).first()
kpi_summary.append({
"name": k.kpi_name,
"code": k.kpi_code,
"dimension": k.dimension,
"target": k.target_value,
"actual": latest.actual_value if latest else None,
"period": latest.period if latest else None,
"unit": k.unit,
})
# 获取预警
alerts = db.query(KPIAlert).filter(KPIAlert.status == "pending").count()
# 构建分析prompt
kpi_text = "\n".join([f"- {k['name']}({k['code']}): 目标={k['target']}, 实际={k['actual']}({k['period']}), 维度={k['dimension']}" for k in kpi_summary if k['actual'] is not None])
prompt = f"""我是一家公司的管理层,以下是当前管理会计系统的KPI数据和系统状态,请给出专业的分析和管理建议:
当前KPI数据:
{kpi_text}
待处理预警数:{alerts}
请从以下三个方面分析:
1. **核心发现**:当前数据反映的最关键问题是什么?
2. **深入解读**:从CMA管理会计角度,这些数据意味着什么?
3. **行动建议**:基于数据,财务和业务部门应该采取什么具体行动?
注意:角色视角为{"CEO(总经理)" if role == "ceo" else "财务部" if role == "finance" else "业务部"}。"""
try:
analysis = await _call_deepseek(prompt)
except Exception as e:
analysis = f"AI分析暂时不可用: {str(e)}"
# R1: 规则驱动生成可落地决策建议(幂等落库)
try:
generate_rule_suggestions(db, entity_id, source="dashboard")
except Exception as e:
pass
result = {"analysis": analysis, "kpi_count": len(kpi_summary), "alert_count": alerts,
"suggestions": _unapplied_suggestions(db, entity_id)}
# 缓存10分钟
cache_set("ai", cache_key, result, ttl_seconds=600)
return result
@router.get("/kpi-analysis/{kpi_id}")
async def kpi_analysis(kpi_id: int, db: Session = Depends(get_db),
entity_id: int = Depends(get_entity_id)):
"""AI分析单个KPI"""
# 尝试缓存
cache_key = f"kpi_analysis:{kpi_id}:{entity_id}"
cached = cache_get("ai", cache_key)
if cached:
return cached
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == kpi_id).first()
if not kpi:
raise HTTPException(404, "KPI不存在")
values = db.query(KPIValue).filter(KPIValue.kpi_id == kpi_id).order_by(KPIValue.period.asc()).all()
trend_data = []
for v in values:
trend_data.append({"period": v.period, "value": v.actual_value})
prompt = f"""请分析以下KPI指标:
KPI名称:{kpi.kpi_name}
维度:{kpi.dimension}
计算公式:{kpi.formula}
目标值:{kpi.target_value}
单位:{kpi.unit}
负责部门:{kpi.responsible_dept}
历史数据趋势:
{json.dumps(trend_data, ensure_ascii=False, indent=2)}
请分析:
1. 当前表现如何,是否达到目标
2. 趋势走势是否健康(上升/下降/波动)
3. 存在什么风险
4. 建议采取什么管理行动"""
try:
analysis = await _call_deepseek(prompt)
except Exception as e:
analysis = f"分析暂时不可用: {str(e)}"
# R1: 生成该KPI的可落地建议
try:
generate_rule_suggestions(db, entity_id, source="kpi", kpi_id=kpi_id)
except Exception as e:
pass
result = {"kpi_name": kpi.kpi_name, "analysis": analysis,
"suggestions": _unapplied_suggestions(db, entity_id)}
cache_set("ai", cache_key, result, ttl_seconds=600)
return result
async def _stream_analysis(prompt: str):
"""流式调用DeepSeek并生成SSE事件"""
async with httpx.AsyncClient(timeout=60) as client:
async with client.stream(
"POST",
"https://api.deepseek.com/v1/chat/completions",
headers={
"Authorization": f"Bearer {os.getenv('DEEPSEEK_API_KEY', 'sk-8e24e6eb87f2475e96ea0980002dc2e8')}",
"Content-Type": "application/json",
},
json={
"model": "deepseek-chat",
"messages": [
{"role": "system", "content": "你是一名CMA管理会计师,擅长用数据驱动的方式分析企业经营状况,给出专业的财务分析和管理建议。"},
{"role": "user", "content": prompt},
],
"stream": True,
"temperature": 0.3,
}
) as response:
async for line in response.aiter_lines():
if not line or line.startswith(":"):
continue
if line.startswith("data: "):
data_str = line[6:]
if data_str.strip() == "[DONE]":
break
try:
chunk = json.loads(data_str)
delta = chunk.get("choices", [{}])[0].get("delta", {}).get("content", "")
if delta:
yield f"data: {json.dumps({'text': delta})}\n\n"
except json.JSONDecodeError:
continue
yield "data: {\"text\": \"[DONE]\"}\n\n"
@router.get("/dashboard-analysis-stream")
async def dashboard_analysis_stream(role: str = Query("ceo"), db: Session = Depends(get_db)):
"""AI分析驾驶舱数据 — SSE流式输出"""
cache_key = f"dashboard_analysis:{role}"
cached = cache_get("ai", cache_key)
if cached:
# 缓存存在,直接以流的形式一次性返回
full_text = cached.get("analysis", "")
async def cached_stream():
yield f"data: {json.dumps({'text': full_text})}\n\n"
yield "data: {\"text\": \"[DONE]\"}\n\n"
return StreamingResponse(cached_stream(), media_type="text/event-stream")
kpis = db.query(KPIDefinition).filter(KPIDefinition.status == "active").all()
kpi_summary = []
for k in kpis:
latest = db.query(KPIValue).filter(KPIValue.kpi_id == k.id).order_by(KPIValue.period.desc()).first()
kpi_summary.append({
"name": k.kpi_name, "code": k.kpi_code, "dimension": k.dimension,
"target": k.target_value, "actual": latest.actual_value if latest else None,
"period": latest.period if latest else None, "unit": k.unit,
})
alerts = db.query(KPIAlert).filter(KPIAlert.status == "pending").count()
kpi_text = "\n".join([f"- {k['name']}({k['code']}): 目标={k['target']}, 实际={k['actual']}({k['period']}), 维度={k['dimension']}" for k in kpi_summary if k['actual'] is not None])
role_label = {"ceo": "CEO(总经理)", "finance": "财务部", "business": "业务部"}.get(role, "管理层")
prompt = f"""我是一家公司的管理层,以下是当前管理会计系统的KPI数据和系统状态,请给出专业的分析和管理建议:
当前KPI数据:
{kpi_text}
待处理预警数:{alerts}
请从以下三个方面分析:
1. **核心发现**:当前数据反映的最关键问题是什么?
2. **深入解读**:从CMA管理会计角度,这些数据意味着什么?
3. **行动建议**:基于数据,财务和业务部门应该采取什么具体行动?
注意:角色视角为{role_label}。"""
return StreamingResponse(_stream_analysis(prompt), media_type="text/event-stream")
@router.post("/ask")
async def ask_question(
request: Request,
db: Session = Depends(get_db),
current_user: User = Depends(require_auth),
):
"""自然语言查询 — CEO问企业经营问题"""
body = await request.json()
question = body.get("question", "").strip()
if not question:
raise HTTPException(400, "请输入问题")
# 收集系统数据作为上下文
kpis = db.query(KPIDefinition).filter(KPIDefinition.status == "active").all()
kpi_context = []
for k in kpis:
latest = db.query(KPIValue).filter(KPIValue.kpi_id == k.id).order_by(KPIValue.period.desc()).first()
alert = db.query(KPIAlert).filter(KPIAlert.kpi_id == k.id, KPIAlert.status == "pending").first()
kpi_context.append(
f"{k.kpi_name}({k.kpi_code}): 当前值={latest.actual_value if latest else '无'}"
f"{' ⚠️' + alert.alert_level if alert else ''}"
)
# 获取改善计划
plans = db.query(ActionPlan).order_by(ActionPlan.created_at.desc()).limit(10).all()
plan_context = [f"- {p.title}({p.assignee}, {p.status}, {p.progress}%)" for p in plans]
# 获取预警
red_alerts = db.query(KPIAlert).filter(KPIAlert.status == "pending", KPIAlert.alert_level == "red").count()
yellow_alerts = db.query(KPIAlert).filter(KPIAlert.status == "pending", KPIAlert.alert_level == "yellow").count()
system_context = f"""你是管理会计OS的AI助手,基于以下企业数据回答管理层问题。
时间:{datetime.now().strftime('%Y-%m-%d %H:%M')}
当前用户:{current_user.name} ({current_user.role})
## KPI数据
{chr(10).join(kpi_context)}
## 预警概况
红色(紧急): {red_alerts}条 | 黄色(预警): {yellow_alerts}
## 改善计划
{chr(10).join(plan_context) if plan_context else '暂无'}
请基于以上数据回答问题。如果问题需要具体数据但上下文中没有,可以根据KPI编码名称推断。回答要简洁、有数据支撑。"""
prompt = f"{system_context}\n\n用户问题:{question}"
try:
analysis = await _call_deepseek(prompt)
except Exception as e:
analysis = f"查询失败: {str(e)}"
return {"question": question, "answer": analysis, "timestamp": datetime.now().isoformat()}
@router.post("/review-plans")
async def review_plans(
db: Session = Depends(get_db),
current_user: User = Depends(require_auth),
):
"""AI复盘改善行动计划执行效果"""
plans = db.query(ActionPlan).order_by(ActionPlan.created_at.asc()).all()
if not plans:
return {"analysis": "暂无改善行动计划,无法复盘"}
plan_text = []
for p in plans:
kpi = db.query(KPIDefinition).filter(KPIDefinition.id == p.kpi_id).first()
kpi_name = kpi.kpi_name if kpi else "未知"
latest = db.query(KPIValue).filter(KPIValue.kpi_id == p.kpi_id).order_by(KPIValue.period.desc()).first()
plan_text.append(
f"- {p.title}\n"
f" 关联KPI: {kpi_name}(当前值: {latest.actual_value if latest else '无'})\n"
f" 负责人: {p.assignee} | 状态: {p.status} | 进度: {p.progress}%\n"
f" 描述: {p.description}\n"
f" 截止日: {p.due_date.strftime('%Y-%m-%d') if p.due_date else '无'}"
)
completed = sum(1 for p in plans if p.status == "completed")
in_progress = sum(1 for p in plans if p.status == "in_progress")
pending = sum(1 for p in plans if p.status == "pending")
prompt = f"""请复盘以下改善行动计划的执行情况:
## 改善计划概览
总数: {len(plans)} | 已完成: {completed} | 进行中: {in_progress} | 待开始: {pending}
## 各计划详情
{chr(10).join(plan_text)}
请分析:
1. **执行概况**:整体执行到位吗?哪些计划需要重点关注?
2. **效果评估**:已完成的计划是否真正改善了关联KPI?
3. **风险提示**:哪些计划存在延期或执行不力的风险?
4. **改进建议**:接下来应该调整或优先推进哪些计划?"""
try:
analysis = await _call_deepseek(prompt)
except Exception as e:
analysis = f"复盘失败: {str(e)}"
return {
"analysis": analysis,
"stats": {"total": len(plans), "completed": completed, "in_progress": in_progress, "pending": pending},
"timestamp": datetime.now().isoformat(),
}