接 b3da1b6。上一条只改了凌晨规则与日报日界,剩下几处一并收掉: 1. risk_query_service.py:77-78:REST 的 start_time/end_time 是**裸 datetime**, 原先原样透传去比库内 UTC 列,而 Agent 路径本来就带时区 (risk_natural_language.py:117)——同一条筛选条件在界面与对话里会查出不同结果。 timeutil 新增 from_local:裸值按**北京时间**解释(面向中国客户的业务系统, 填表人的预期就是本地时间),带时区的按其自身时区处理。它与 to_utc_naive 的区别 正在裸值上:取库里的值用后者,接客户端输入用这个。 2. risk_scan_service.py 与 risk_judgement_service.py 的 _age:一处用 UTC 日期、 一处用服务器 date.today(),生日边界上同一客户会差一岁、65 岁阈值可能翻面。 统一走 local_date(北京时间)。 3. risk_daily_report_service.py:122 的 report_date 与 :244 的 created_today: 原先取 UTC 日期,北京 08:00 之前会把"今天新增的预警"算成昨天。 ruff / mypy(135 文件) / 603 unit+contract 全绿。
427 lines
17 KiB
Python
427 lines
17 KiB
Python
"""风控九段式日报 Service。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import logging
|
||
from collections import Counter
|
||
from collections.abc import AsyncIterator
|
||
from datetime import UTC, datetime
|
||
from typing import Any
|
||
from zoneinfo import ZoneInfo
|
||
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.core.config import get_settings
|
||
from app.core.contracts import RequestContext
|
||
from app.core.timeutil import local_date, local_day_bounds
|
||
from app.model.audit import InteractionAudit
|
||
from app.repository.fund_query_repository import FundRecord
|
||
from app.repository.risk_repository import RiskRepository
|
||
from app.service.authorization_service import AuthorizationService
|
||
from app.service.model_gateway import (
|
||
DatabaseModelEndpointResolver,
|
||
ModelGenerationService,
|
||
)
|
||
from app.service.risk_query_service import scope_from_context
|
||
|
||
OPEN_STATUSES = ("待处理", "调查中")
|
||
LEVEL_ORDER = ("高", "中", "低")
|
||
MAX_SUGGESTION_CHARS = 3000
|
||
FORBIDDEN_CLAIMS = ("已修改规则", "已调整规则", "已关闭预警", "已升级预警", "已完成处置")
|
||
PROMPT_VERSION = "risk-daily-report-v2"
|
||
SYSTEM_PROMPT = (
|
||
"你是公募基金风控日报辅助工具。只能根据确定性统计提出建议,"
|
||
"不得补造数据,不得声称已经修改规则、关闭预警或完成处置。"
|
||
)
|
||
USER_INSTRUCTION = "请生成最多 5 条简体中文风控日报优化建议,只输出建议正文。"
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
class RiskDailyReportService:
|
||
def __init__(
|
||
self,
|
||
session: AsyncSession,
|
||
*,
|
||
repository: RiskRepository | None = None,
|
||
model_service: ModelGenerationService | None = None,
|
||
endpoint_resolver: DatabaseModelEndpointResolver | None = None,
|
||
) -> None:
|
||
self.session = session
|
||
self.repository = repository
|
||
self.model_service = model_service
|
||
self.endpoint_resolver = endpoint_resolver
|
||
|
||
async def generate(
|
||
self,
|
||
context: RequestContext,
|
||
now: datetime | None = None,
|
||
) -> dict[str, Any]:
|
||
await AuthorizationService.require(context, "risk:alert:read")
|
||
generated_at = _utc_naive(now)
|
||
report = await self._build(context, generated_at)
|
||
suggestions, source = await self._suggestions(report)
|
||
return await self._complete(report, suggestions, source)
|
||
|
||
async def stream(
|
||
self,
|
||
context: RequestContext,
|
||
now: datetime | None = None,
|
||
) -> AsyncIterator[dict[str, Any]]:
|
||
await AuthorizationService.require(context, "risk:alert:read")
|
||
generated_at = _utc_naive(now)
|
||
yield {"type": "start", "generated_at": _local_datetime_text(generated_at)}
|
||
yield {"type": "progress", "stage": "statistics", "message": "正在统计日报数据"}
|
||
report = await self._build(context, generated_at)
|
||
report["optimization_suggestions"] = ""
|
||
yield {"type": "replace", "content": self.render_text(report)}
|
||
yield {"type": "progress", "stage": "suggestions", "message": "正在生成优化建议"}
|
||
suggestions, source = await self._suggestions(report)
|
||
report["optimization_suggestions"] = suggestions
|
||
yield {"type": "replace", "content": self.render_text(report)}
|
||
yield {"type": "done", "report": await self._complete(report, suggestions, source)}
|
||
|
||
async def _build(
|
||
self,
|
||
context: RequestContext,
|
||
generated_at: datetime,
|
||
) -> dict[str, Any]:
|
||
# 日界按**北京时间**切,再换算回库内 UTC 格式去查询。原先直接用 UTC 日期切日,
|
||
# 北京 08:00 前生成的日报,统计窗口会变成"前一日 08:00–当日 08:00"、跨了零点
|
||
# (docs/25 P1 #6)。查询列是 UTC,所以 local_day_bounds 返回的也是 UTC naive。
|
||
day_start, next_day = local_day_bounds(generated_at)
|
||
repository = self.repository or RiskRepository(
|
||
self.session,
|
||
scope=scope_from_context(context),
|
||
)
|
||
snapshot = await repository.daily_report_snapshot(day_start, next_day)
|
||
daily = [_record_values(item) for item in snapshot.daily]
|
||
unresolved = [_record_values(item) for item in snapshot.unresolved]
|
||
false_positive = [_record_values(item) for item in snapshot.false_positive]
|
||
dispositions = [_record_values(item) for item in snapshot.dispositions]
|
||
key_alerts = {
|
||
item["alert_no"]: item
|
||
for item in [*daily, *unresolved]
|
||
if item.get("risk_level") == "高"
|
||
}
|
||
unresolved_items = [
|
||
self._alert_item(item, generated_at)
|
||
for item in unresolved
|
||
]
|
||
key_items = [
|
||
self._alert_item(key_alerts[key], generated_at)
|
||
for key in sorted(key_alerts)
|
||
]
|
||
false_positive_items = [
|
||
self._alert_item(item, generated_at)
|
||
for item in false_positive
|
||
]
|
||
return {
|
||
"type": "风控日报",
|
||
"report_date": local_date(generated_at).isoformat(),
|
||
"generated_at": _local_datetime_text(generated_at),
|
||
"daily_alert_count": len(daily),
|
||
"level_distribution": _distribution(daily, "risk_level", LEVEL_ORDER),
|
||
"key_risk_events": key_items,
|
||
"unresolved_items": {
|
||
"total": len(unresolved_items),
|
||
"new_today": sum(item["created_today"] for item in unresolved_items),
|
||
"historical": sum(not item["created_today"] for item in unresolved_items),
|
||
"overdue": sum(item["is_overdue"] for item in unresolved_items),
|
||
"items": unresolved_items,
|
||
},
|
||
"false_positive_statistics": {
|
||
"total": len(false_positive_items),
|
||
"reasons": [
|
||
{
|
||
"alert_id": item["alert_id"],
|
||
"reason": item.get("close_reason") or "未填写",
|
||
}
|
||
for item in false_positive_items
|
||
],
|
||
},
|
||
"type_distribution": _distribution(daily, "alert_type"),
|
||
"disposition_results": {
|
||
"acknowledged": sum(
|
||
_in_day(item.get("ack_at"), day_start, next_day)
|
||
for item in dispositions
|
||
),
|
||
"false_positive_closed": len(false_positive_items),
|
||
"escalated": sum(
|
||
_in_day(item.get("escalated_at"), day_start, next_day)
|
||
for item in dispositions
|
||
),
|
||
"investigating": sum(
|
||
item.get("status") == "调查中"
|
||
for item in dispositions
|
||
),
|
||
},
|
||
"rule_effectiveness": _rule_effectiveness(
|
||
daily,
|
||
unresolved,
|
||
false_positive,
|
||
),
|
||
"optimization_suggestions": "",
|
||
"source": "",
|
||
"prompt_version": PROMPT_VERSION,
|
||
}
|
||
|
||
async def _suggestions(self, report: dict[str, Any]) -> tuple[str, str]:
|
||
model_service = self.model_service
|
||
endpoint_resolver = self.endpoint_resolver
|
||
if model_service is None or endpoint_resolver is None:
|
||
from app.service.agent.bootstrap import get_model_service
|
||
from app.service.model_gateway import DatabaseModelEndpointResolver
|
||
|
||
model_service = model_service or get_model_service()
|
||
endpoint_resolver = endpoint_resolver or DatabaseModelEndpointResolver()
|
||
try:
|
||
endpoints = await endpoint_resolver.resolve(
|
||
agent_type="risk",
|
||
task_type="daily_report_suggestion",
|
||
)
|
||
execution = await model_service.generate(
|
||
endpoints,
|
||
self._suggestion_prompt(report),
|
||
max_attempts=2,
|
||
)
|
||
text = execution.text.strip()
|
||
if (
|
||
not text
|
||
or len(text) > MAX_SUGGESTION_CHARS
|
||
or any(claim in text for claim in FORBIDDEN_CLAIMS)
|
||
):
|
||
raise ValueError("日报建议无效")
|
||
return text, "模型"
|
||
except Exception:
|
||
logger.exception("日报建议生成失败,已使用规则化建议")
|
||
return self._fallback(report), "规则化模板"
|
||
|
||
async def _complete(
|
||
self,
|
||
report: dict[str, Any],
|
||
suggestions: str,
|
||
source: str,
|
||
) -> dict[str, Any]:
|
||
report["optimization_suggestions"] = suggestions
|
||
report["source"] = source
|
||
report["content"] = self.render_text(report)
|
||
self.session.add(InteractionAudit(
|
||
actor_type="system",
|
||
portal="api",
|
||
action_type="risk_daily_report_generated",
|
||
detail={
|
||
"report_date": report["report_date"],
|
||
"daily_alert_count": report["daily_alert_count"],
|
||
"unresolved_total": report["unresolved_items"]["total"],
|
||
"source": source,
|
||
"prompt_version": PROMPT_VERSION,
|
||
},
|
||
created_at=datetime.now(UTC).replace(tzinfo=None),
|
||
))
|
||
await self.session.commit()
|
||
return report
|
||
|
||
@staticmethod
|
||
def _alert_item(item: dict[str, Any], generated_at: datetime) -> dict[str, Any]:
|
||
created_at = _as_datetime(item.get("created_at"))
|
||
due_at = _as_datetime(item.get("due_at"))
|
||
return {
|
||
"alert_id": item.get("alert_no"),
|
||
"alert_level": f"{item.get('risk_level')}风险",
|
||
"alert_type": item.get("alert_type"),
|
||
"triggered_rules": list(item.get("rule_codes") or []),
|
||
"evidence_summary": item.get("evidence_summary"),
|
||
"status": item.get("status"),
|
||
"ack_status": item.get("ack_status"),
|
||
"handler_id": item.get("handler_id"),
|
||
"created_at": _local_datetime_text(created_at) if created_at else None,
|
||
"due_time": _local_datetime_text(due_at) if due_at else None,
|
||
"is_overdue": bool(due_at and due_at <= generated_at),
|
||
"is_escalated": bool(item.get("is_escalated")),
|
||
"close_reason": item.get("close_reason"),
|
||
# 两边都按**北京时间**取日期:原先比的是 UTC 日期,北京 08:00 之前
|
||
# 会把"今天新增的预警"算成昨天(docs/25 P1 #6)。
|
||
"created_today": bool(
|
||
created_at and local_date(created_at) == local_date(generated_at)
|
||
),
|
||
}
|
||
|
||
@staticmethod
|
||
def _suggestion_prompt(report: dict[str, Any]) -> str:
|
||
summary = {
|
||
"report_date": report["report_date"],
|
||
"daily_alert_count": report["daily_alert_count"],
|
||
"level_distribution": report["level_distribution"],
|
||
"unresolved_summary": {
|
||
key: value
|
||
for key, value in report["unresolved_items"].items()
|
||
if key != "items"
|
||
},
|
||
"false_positive_total": report["false_positive_statistics"]["total"],
|
||
"type_distribution": report["type_distribution"],
|
||
"disposition_results": report["disposition_results"],
|
||
"rule_effectiveness": report["rule_effectiveness"],
|
||
}
|
||
statistics = json.dumps(summary, ensure_ascii=False)
|
||
return f"{SYSTEM_PROMPT}\n{USER_INSTRUCTION}\n统计数据:{statistics}"
|
||
|
||
@staticmethod
|
||
def _fallback(report: dict[str, Any]) -> str:
|
||
unresolved = report["unresolved_items"]
|
||
suggestions: list[str] = []
|
||
if unresolved["overdue"]:
|
||
suggestions.append(
|
||
f"优先复核 {unresolved['overdue']} 条已超时未闭环预警,"
|
||
"并补充处理留痕。"
|
||
)
|
||
if unresolved["historical"]:
|
||
suggestions.append(
|
||
f"逐条清理 {unresolved['historical']} 条历史遗留未闭环预警,明确责任人和完成时限。"
|
||
)
|
||
false_positive_rules = [
|
||
item["rule_code"]
|
||
for item in report["rule_effectiveness"]
|
||
if item["false_positives"] > 0
|
||
]
|
||
if false_positive_rules:
|
||
suggestions.append(
|
||
f"对规则 {','.join(false_positive_rules)} 的误报样本人工复盘,"
|
||
"完成回测和审批后再调整。"
|
||
)
|
||
if not suggestions:
|
||
suggestions.append("继续监测预警变化,保持规则命中证据和人工处置记录完整。")
|
||
return "\n".join(f"{index}. {item}" for index, item in enumerate(suggestions, start=1))
|
||
|
||
@classmethod
|
||
def render_text(cls, report: dict[str, Any]) -> str:
|
||
unresolved = report["unresolved_items"]
|
||
lines = [
|
||
f"风控预警日报({report['report_date']})",
|
||
"",
|
||
f"1. 当日预警数量:{report['daily_alert_count']}",
|
||
f"2. 等级分布:{_distribution_text(report['level_distribution'])}",
|
||
"3. 重点风险事件:",
|
||
*_alert_lines(report["key_risk_events"]),
|
||
(
|
||
"4. 未闭环事项:"
|
||
f"共 {unresolved['total']} 条,当日新增 {unresolved['new_today']} 条,"
|
||
f"历史遗留 {unresolved['historical']} 条,已超时 {unresolved['overdue']} 条"
|
||
),
|
||
*_alert_lines(unresolved["items"]),
|
||
f"5. 误报统计:当日关闭误报 {report['false_positive_statistics']['total']} 条",
|
||
*[
|
||
f"- {item['alert_id']}:{item['reason']}"
|
||
for item in report["false_positive_statistics"]["reasons"]
|
||
],
|
||
f"6. 类型分布:{_distribution_text(report['type_distribution'])}",
|
||
(
|
||
"7. 处置结果:"
|
||
f"确认接收 {report['disposition_results']['acknowledged']} 条,"
|
||
f"关闭误报 {report['disposition_results']['false_positive_closed']} 条,"
|
||
f"调查中 {report['disposition_results']['investigating']} 条,"
|
||
f"升级 {report['disposition_results']['escalated']} 条"
|
||
),
|
||
"8. 规则效果:",
|
||
*[
|
||
(
|
||
f"- {item['rule_code']}:当日命中 {item['daily_hits']} 条,"
|
||
f"未闭环 {item['unresolved']} 条,当日误报 {item['false_positives']} 条"
|
||
)
|
||
for item in report["rule_effectiveness"]
|
||
],
|
||
"9. 建议优化方向:",
|
||
report["optimization_suggestions"],
|
||
]
|
||
return "\n".join(lines)
|
||
|
||
|
||
def _record_values(record: FundRecord) -> dict[str, Any]:
|
||
return record.to_dict()
|
||
|
||
|
||
def _distribution(
|
||
items: list[dict[str, Any]],
|
||
field: str,
|
||
preferred_order: tuple[str, ...] = (),
|
||
) -> list[dict[str, Any]]:
|
||
counts = Counter(item.get(field) for item in items if item.get(field) is not None)
|
||
keys = [key for key in preferred_order if key in counts]
|
||
keys.extend(sorted(str(key) for key in counts if key not in preferred_order))
|
||
return [
|
||
{
|
||
"name": f"{key}风险" if field == "risk_level" else key,
|
||
"count": counts[key],
|
||
}
|
||
for key in keys
|
||
]
|
||
|
||
|
||
def _rule_effectiveness(
|
||
daily: list[dict[str, Any]],
|
||
unresolved: list[dict[str, Any]],
|
||
false_positive: list[dict[str, Any]],
|
||
) -> list[dict[str, Any]]:
|
||
counters = [
|
||
Counter(code for item in group for code in (item.get("rule_codes") or []))
|
||
for group in (daily, unresolved, false_positive)
|
||
]
|
||
return [
|
||
{
|
||
"rule_code": code,
|
||
"daily_hits": counters[0][code],
|
||
"unresolved": counters[1][code],
|
||
"false_positives": counters[2][code],
|
||
}
|
||
for code in sorted(set().union(*(set(counter) for counter in counters)))
|
||
]
|
||
|
||
|
||
def _in_day(value: Any, day_start: datetime, next_day: datetime) -> bool:
|
||
parsed = _as_datetime(value)
|
||
return bool(parsed and day_start <= parsed < next_day)
|
||
|
||
|
||
def _as_datetime(value: Any) -> datetime | None:
|
||
if isinstance(value, datetime):
|
||
return value
|
||
if isinstance(value, str):
|
||
try:
|
||
return datetime.fromisoformat(value.replace("Z", "+00:00")).replace(tzinfo=None)
|
||
except ValueError:
|
||
return None
|
||
return None
|
||
|
||
|
||
def _distribution_text(items: list[dict[str, Any]]) -> str:
|
||
return ",".join(f"{item['name']} {item['count']} 条" for item in items) or "无"
|
||
|
||
|
||
def _alert_lines(items: list[dict[str, Any]]) -> list[str]:
|
||
if not items:
|
||
return ["- 无"]
|
||
return [
|
||
(
|
||
f"- {item['alert_id']}|{item['alert_level']}|{item['alert_type']}|"
|
||
f"状态 {item['status']}|规则 {','.join(item['triggered_rules']) or '-'}|"
|
||
f"期限 {item['due_time'] or '-'}|{'已超时' if item['is_overdue'] else '未超时'}"
|
||
)
|
||
for item in items
|
||
]
|
||
|
||
|
||
def _utc_naive(value: datetime | None) -> datetime:
|
||
if value is None:
|
||
return datetime.now(UTC).replace(tzinfo=None)
|
||
if value.tzinfo is not None:
|
||
return value.astimezone(UTC).replace(tzinfo=None)
|
||
return value
|
||
|
||
|
||
def _local_datetime_text(value: datetime) -> str:
|
||
"""把 UTC 时间转换为配置时区的常用展示格式。"""
|
||
aware = value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
|
||
local_value = aware.astimezone(ZoneInfo(get_settings().timezone))
|
||
return local_value.strftime("%Y-%m-%d %H:%M:%S")
|