"""知识摘要服务 — 管理会计OS持久记忆 仿 OpenCode SummarizeProvider 的持久记忆机制,做三处改进: 1. 分层压缩:日→周→月→全部,不是一次性全量压缩 2. 结构化存储:MySQL 关系表,不是 SQLite JSON 消息 3. 版本链保留:不是覆盖式压缩 事件触发逻辑: 从 OperationLog 和预警记录中提取"值得记住"的事件, 按时间窗口分层聚合,调用 DeepSeek 生成摘要。 """ import logging import os import json import httpx from datetime import datetime, date, timedelta from typing import Optional from sqlalchemy.orm import Session from sqlalchemy import func, and_ from app.models import KnowledgeEvent, KnowledgeSummary, OperationLog, KPIAlert, KPIDefinition, KPIValue logger = logging.getLogger("cma.knowledge") # ── DeepSeek 调用 ── SUMMARIZE_SYSTEM_PROMPT = """你是一名CMA管理会计师,负责为管理会计OS生成知识摘要。 你的工作是:审核一组经营事件记录,提炼出"必须记住"的核心信息。 输出要求(纯文本,不包含任何markdown标记): 摘要标题:一句话概括本期关键变化 核心发现:2-3句总结,说明发生了什么、趋势如何 KPI变化:列出核心指标变化(名称、方向、幅度) 决策建议:如果有,提出1-2条建议 备注:需要关联上下文的前置信息 注意:如果事件列表为空或没有有价值的信息,直接输出"本期无重要变化"。""" async def _call_deepseek(prompt: str, timeout: int = 30) -> str: """调用DeepSeek API生成摘要""" api_key = os.getenv("DEEPSEEK_API_KEY", "sk-8e2...c2e8") try: async with httpx.AsyncClient(timeout=timeout) 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": SUMMARIZE_SYSTEM_PROMPT}, {"role": "user", "content": prompt}, ], "stream": False, "temperature": 0.2, "max_tokens": 1024, }, ) data = resp.json() return data.get("choices", [{}])[0].get("message", {}).get("content", "") except Exception as e: logger.error(f"DeepSeek摘要调用失败: {e}") return "" # ── 事件抽取 ── def extract_events(db: Session, since: datetime, until: Optional[datetime] = None) -> list[dict]: """从 OperationLog + KPIAlert 中抽取关键事件 返回 dict 列表,用于喂给摘要 prompt """ now = until or datetime.utcnow() # 1. 操作日志 → 事件 logs = ( db.query(OperationLog) .filter(OperationLog.created_at >= since, OperationLog.created_at <= now) .order_by(OperationLog.created_at.asc()) .all() ) events = [] for log in logs: detail_str = "" if log.detail: # 裁短 detail 避免 token 浪费 d = json.dumps(log.detail, ensure_ascii=False) detail_str = d[:300] + ("..." if len(d) > 300 else "") events.append({ "type": "action", "time": log.created_at.isoformat() if log.created_at else "", "action": log.action, "target": f"{log.target_type}#{log.target_id}", "detail": detail_str, }) # 2. 预警记录 → 事件 alerts = ( db.query(KPIAlert, KPIDefinition.kpi_name, KPIDefinition.kpi_code) .join(KPIDefinition, KPIAlert.kpi_id == KPIDefinition.id) .filter(KPIAlert.created_at >= since, KPIAlert.created_at <= now) .order_by(KPIAlert.created_at.asc()) .all() ) for alert, kpi_name, kpi_code in alerts: events.append({ "type": "alert", "time": alert.created_at.isoformat() if alert.created_at else "", "level": alert.alert_level, "kpi": f"{kpi_code} ({kpi_name})", "message": (alert.alert_message or "")[:200], "status": alert.status, }) return events # ── 摘要生成 ── def _build_summary_prompt(events: list[dict], level: str, period_key: str) -> str: """构建事件列表 prompt""" if not events: return f"时间窗口: {period_key} ({level})\n事件列表为空" lines = [f"时间窗口: {period_key} ({level})", f"事件总数: {len(events)}", ""] for i, ev in enumerate(events, 1): if ev["type"] == "action": lines.append(f"{i}. [操作] {ev['time']} {ev['action']} on {ev['target']} | {ev['detail']}") elif ev["type"] == "alert": lines.append(f"{i}. [预警] {ev['time']} [{ev['level']}] {ev['kpi']} | {ev['message']} (状态: {ev['status']})") else: lines.append(f"{i}. [其他] {ev['time']} {ev}") return "\n".join(lines) async def _kpi_snapshot(db: Session) -> list[dict]: """当前KPI快照 — 每个活跃KPI的最新实际值""" kpis = db.query(KPIDefinition).filter(KPIDefinition.status == "active").all() snapshot = [] for k in kpis: latest = ( db.query(KPIValue) .filter(KPIValue.kpi_id == k.id) .order_by(KPIValue.period.desc()) .first() ) if latest: snapshot.append({ "code": k.kpi_code, "name": k.kpi_name, "value": latest.actual_value, "period": latest.period, "unit": k.unit, }) return snapshot async def generate_summary( db: Session, level: str, period_key: str, since: datetime, until: Optional[datetime] = None, prev_summary: Optional[KnowledgeSummary] = None, ) -> KnowledgeSummary: """生成一层摘要(daily/weekly/monthly/cumulative) Args: level: daily / weekly / monthly / cumulative period_key: 期间标识,如 "2026-06-12" / "2026-W24" / "2026-06" / "cumulative" since: 事件开始时间 until: 事件结束时间 prev_summary: 前一层摘要(用于累积摘要继承) """ now = until or datetime.utcnow() # 1. 抽取事件 events = extract_events(db, since, now) # 2. 构建 prompt prompt = _build_summary_prompt(events, level, period_key) # 3. 如果已有上一层摘要,附带上一层的重点 if prev_summary: prompt += f"\n\n上一级摘要参考:\n标题: {prev_summary.title}\n内容: {prev_summary.content[:500]}\n" # 4. 调用 DeepSeek result_text = await _call_deepseek(prompt) # 5. 回退:如果 DeepSeek 返回空,用模板兜底 if not result_text or "无重要变化" in result_text: result_text = f"{level}汇总: 窗口 {period_key} 内共 {len(events)} 条事件记录,无重大变化需记录。" # 6. 提取 KPI 快照 kpi_snapshot = await _kpi_snapshot(db) # 7. 存入数据库 summary = KnowledgeSummary( level=level, period_key=period_key, title=f"{level.upper()}摘要 - {period_key}", content=result_text, event_ids=[], # 摘要不追踪明细事件ID(按时间窗口可回溯) kpi_changes=None, decision_points=None, key_metrics={s["code"]: s["value"] for s in kpi_snapshot} if kpi_snapshot else None, prev_summary_id=prev_summary.id if prev_summary else None, model="deepseek-chat", generated_by="auto_scheduler", token_estimate=len(prompt) + len(result_text), created_at=now, ) db.add(summary) db.commit() db.refresh(summary) logger.info(f"知识摘要已生成: level={level} period={period_key} id={summary.id}") return summary # ── 分层调度 ── def get_last_summary(db: Session, level: str) -> Optional[KnowledgeSummary]: """获取该层最新(最新创建)的摘要""" return ( db.query(KnowledgeSummary) .filter(KnowledgeSummary.level == level, KnowledgeSummary.is_stale == 0) .order_by(KnowledgeSummary.id.desc()) .first() ) async def run_daily_summary(db: Session) -> KnowledgeSummary: """运行每日摘要""" today = date.today() period_key = today.isoformat() since = datetime(today.year, today.month, today.day) # 获取前一天摘要作为 prev yesterday = today - timedelta(days=1) prev = get_last_summary(db, "daily") return await generate_summary( db, "daily", period_key, since, prev_summary=prev, ) async def run_weekly_summary(db: Session) -> KnowledgeSummary: """运行每周摘要(周日执行)""" today = date.today() # ISO 周算法: 本周一到今天 iso_week = today.isocalendar() period_key = f"{iso_week[0]}-W{iso_week[1]:02d}" since = today - timedelta(days=today.weekday()) # 本周一 since_dt = datetime(since.year, since.month, since.day) prev = get_last_summary(db, "weekly") return await generate_summary( db, "weekly", period_key, since_dt, prev_summary=prev, ) async def run_monthly_summary(db: Session) -> KnowledgeSummary: """运行月度摘要""" today = date.today() period_key = today.strftime("%Y-%m") since = datetime(today.year, today.month, 1) # 附属前一个月的 daily 和 weekly 摘要 prev_month = today.replace(day=1) - timedelta(days=1) prev = get_last_summary(db, "monthly") return await generate_summary( db, "monthly", period_key, since, prev_summary=prev, ) # ── 对外接口(同步包装,供手动触发使用) ── import asyncio def generate_summary_sync( db: Session, level: str, period_key: str, since: datetime, until: Optional[datetime] = None, ) -> dict: """同步包装,用于 API 手动触发""" loop = asyncio.new_event_loop() try: summary = loop.run_until_complete( generate_summary(db, level, period_key, since, until) ) return {"id": summary.id, "level": summary.level, "period_key": summary.period_key} finally: loop.close() def generate_daily_sync(db: Session) -> dict: loop = asyncio.new_event_loop() try: summary = loop.run_until_complete(run_daily_summary(db)) return {"id": summary.id, "level": summary.level, "period_key": summary.period_key} finally: loop.close() def generate_weekly_sync(db: Session) -> dict: loop = asyncio.new_event_loop() try: summary = loop.run_until_complete(run_weekly_summary(db)) return {"id": summary.id, "level": summary.level, "period_key": summary.period_key} finally: loop.close() def generate_monthly_sync(db: Session) -> dict: loop = asyncio.new_event_loop() try: summary = loop.run_until_complete(run_monthly_summary(db)) return {"id": summary.id, "level": summary.level, "period_key": summary.period_key} finally: loop.close()