diff --git a/app/api/simulate.py b/app/api/simulate.py index fb27d9c..b369259 100644 --- a/app/api/simulate.py +++ b/app/api/simulate.py @@ -46,7 +46,7 @@ def submit_trade_api(req: TradeRequest, auth: AuthContext = Depends(get_auth_con agent_type="platform", # 网关越权与放行审计同口径(复审 P3) ) try: - return submit_trade(req.model_dump()) + return submit_trade(req.model_dump(), actor_id=auth.actor_id) except UnsupportedTradeType as exc: raise ApiError(400, "BAD_REQUEST", str(exc)) from exc except LookupError as exc: # NotFoundError 亦为其子类;已统一(B6 评审 P3-5) diff --git a/app/gateway/trade_gateway.py b/app/gateway/trade_gateway.py index 8be96d2..1410a18 100644 --- a/app/gateway/trade_gateway.py +++ b/app/gateway/trade_gateway.py @@ -55,6 +55,7 @@ def _audit( req: dict[str, Any], rule_id: str | None = None, detail: dict[str, Any] | None = None, + actor_id: str | None = None, ) -> None: from app.utils.trace import current_trace, new_trace @@ -63,7 +64,7 @@ def _audit( "trace_id": current_trace() or new_trace(), "event_type": "trade_request", "agent_type": "platform", - "actor_id": "SYSTEM", # B6 接入 AuthContext 后透传 actor_id + "actor_id": actor_id or "SYSTEM", # C6 透传:代理人发起交易归属发起人;缺省 SYSTEM 保持现有测试/脚本零改动 "customer_id": req.get("customer_id"), "rule_id": rule_id, "input_summary": { @@ -90,6 +91,7 @@ def submit_trade( thresholds: RiskThresholds | None = None, now: datetime | None = None, trade_id_factory: Callable[[datetime], str] | None = None, + actor_id: str | None = None, ) -> dict[str, Any]: """处理一笔模拟交易请求(PRD FR-1 流程 ①~⑤)。 @@ -188,6 +190,7 @@ def submit_trade( decision="risk_engine_error", trade_id=trade_id, req=req, + actor_id=actor_id, detail={"error_stage": "process_trade_event"}, ) return { @@ -203,6 +206,7 @@ def submit_trade( decision="trade_accepted", trade_id=trade_id, req=req, + actor_id=actor_id, detail=dict(engine_result), # 全量输出(评审 P2-3) ) return {"blocked": False, "trade_id": trade_id, **engine_result} diff --git a/app/repository/risk_repository.py b/app/repository/risk_repository.py index 5382744..397a108 100644 --- a/app/repository/risk_repository.py +++ b/app/repository/risk_repository.py @@ -184,7 +184,7 @@ class RiskRepository: """ ), { - "payload": json.dumps(payload, ensure_ascii=False), + "payload": json.dumps(payload, ensure_ascii=False, default=str), "rules": json.dumps(merged_rules), "score": new_score, "atype": alert_type, @@ -391,6 +391,115 @@ class RiskRepository: ) return res.rowcount == 1 + # ---------- 代理人行为链(C6 · RISK-008 数据源) ---------- + + def list_audit_events( + self, + event_type: str, + since: datetime, + actor_id: str | None = None, + decision: str | None = None, + limit: int = 2000, + ) -> list[dict[str, Any]]: + """audit_log 滑窗查询(RISK-008 数据源;input_summary 解析回 dict)。 + + WHERE event_type=:et AND created_at>=:since [AND actor_id=:aid] [AND decision=:d] + ORDER BY created_at ASC。条件 B/C 用 actor_id+decision='forbidden'; + 条件 A 用 event_type='trade_request' 全量(actor 过滤在 Python 端做,排除 SYSTEM)。 + """ + where = ["event_type = :et", "created_at >= :since"] + params: dict[str, Any] = {"et": event_type, "since": since} + if actor_id: + where.append("actor_id = :aid") + params["aid"] = actor_id + if decision: + where.append("decision = :dec") + params["dec"] = decision + params["lim"] = limit + sql = text( + f"SELECT * FROM audit_log WHERE {' AND '.join(where)} ORDER BY created_at ASC LIMIT :lim" + ) + with self._engine.connect() as conn: + rows = [dict(r) for r in conn.execute(sql, params).mappings()] + for row in rows: + if isinstance(row.get("input_summary"), str): + try: + row["input_summary"] = json.loads(row["input_summary"]) + except (ValueError, TypeError): + row["input_summary"] = {} + if isinstance(row.get("created_at"), str): + try: + row["created_at"] = datetime.fromisoformat(row["created_at"]) + except ValueError: + pass + return rows + + def find_agent_behavior_alert(self, actor_id: str, day_start: datetime) -> dict | None: + """同代理人当日 pending 行为链单(去重锚点)。 + + SQL 取最近 pending pattern 单,Python 过滤 payload.alert_subtype 含 + agent_behavior 且 payload.actor_id==actor_id(不用 payload LIKE:JSON + 键值歧义 + 序列化空格差异两坑,挂账 #9 登记性能,一期 Python 过滤最稳)。 + """ + sql = text( + """ + SELECT * FROM risk_alert + WHERE status = 'pending_review' AND alert_type = 'pattern' + AND created_at >= :day_start + ORDER BY created_at DESC + """ + ) + with self._engine.connect() as conn: + rows = [self._parse_alert(dict(r)) for r in conn.execute(sql, {"day_start": day_start}).mappings()] + for row in rows: + payload = row.get("payload") or {} + sub = payload.get("alert_subtype") + if isinstance(sub, str): + sub = [sub] + if "agent_behavior" in (sub or []) and payload.get("actor_id") == actor_id: + return row + return None + + def merge_agent_behavior_payload( + self, + alert_id: str, + *, + subtypes_hit: list[str], + customers: list[str], + evidence_trace_ids: list[str], + timeline: dict[str, Any], + ) -> bool: + """agent_behavior 单追加:subtypes_hit/customers/evidence_trace_ids/timeline 并集(读改写)。 + + 与 append_alert_event 互补——后者只合并 payload.events + alert_subtype; + 本方法合并行为链专属的标量列表字段,保证"同日多子条件只一张单"语义(A-12 断言)。 + """ + with self._engine.begin() as conn: + row = conn.execute( + text("SELECT payload FROM risk_alert WHERE alert_id = :aid"), + {"aid": alert_id}, + ).mappings().first() + if row is None: + return False + payload = json.loads(row["payload"]) if isinstance(row["payload"], str) else row["payload"] + payload["subtypes_hit"] = sorted(set(payload.get("subtypes_hit") or []) | set(subtypes_hit)) + payload["customers"] = sorted(set(payload.get("customers") or []) | set(customers)) + payload["evidence_trace_ids"] = sorted( + set(payload.get("evidence_trace_ids") or []) | set(evidence_trace_ids) + ) + old_timeline = payload.get("timeline") or {} + for k, v in (timeline or {}).items(): + old_timeline.setdefault(k, []).extend(v or []) + payload["timeline"] = old_timeline + conn.execute( + text("UPDATE risk_alert SET payload = :payload WHERE alert_id = :aid"), + { + "payload": json.dumps(payload, ensure_ascii=False, default=str), + "aid": alert_id, + }, + ) + return True + @staticmethod def _dump_alert(alert: dict[str, Any]) -> dict[str, Any]: out = dict(alert) diff --git a/app/service/risk/agent_behavior_service.py b/app/service/risk/agent_behavior_service.py new file mode 100644 index 0000000..a72dd7f --- /dev/null +++ b/app/service/risk/agent_behavior_service.py @@ -0,0 +1,383 @@ +"""代理人异常行为链识别(C6 · PRD FR-10 / 规则表 v1.1 RISK-008 补充约束)。 + +数据源(定案口径,input_guard_log 不作为数据源): + 条件 A:audit_log(event_type='trade_request') 中 actor_id 为代理人(排除 SYSTEM + 与 actor_id==customer_id 的本人交易)的交易请求,按行内 trade_id 从 + core_trade 取 traded_at/trade_type/product_id 组装时间线;同 (actor,customer) + 组内,对每个 redeem 在 2h 内找不同产品的 subscribe → 计一次诱导调仓。 + 条件 B:audit_log(event_type='authz', decision='forbidden', input_summary.code= + 'AUTH_403_SCOPE') 越权试探,按 actor 滑动窗口计数。 + 条件 C:同上族 code IN ('AUTH_403_NOT_OWNER','AUTH_403_NOT_ASSIGNED') 越权查询。 + +出单:代理人维度独立 pattern 单(payload.alert_subtype=['agent_behavior']), +与客户维度事件单互不并入(find_pending_event_alert 已排除 agent_behavior 单,P0-1)。 +L3 不写(代理人画像本期仅 payload 承载,挂账 #3)。 +审计 event_type='agent_behavior_detected'(每 actor 每次命中一条;仅 INSERT)。 +红线:不改表 / 审计只 INSERT / 新 Tool 只读 / 无 LLM。 +""" + +from __future__ import annotations + +import logging +from collections import defaultdict +from dataclasses import dataclass +from datetime import datetime, timedelta +from typing import Any + +from app.config.settings import settings +from app.repository.risk_repository import RiskRepository +from app.service.risk.alert_service import _publish_alert +from app.utils.trace import current_trace, new_trace + +logger = logging.getLogger(__name__) + +# 条件 A:赎回后 2h 内申购不同产品 → 诱导调仓(消费制配对窗口) +_INDUCE_WINDOW_SECONDS = 2 * 3600 + + +@dataclass(frozen=True) +class BehaviorThresholds: + """行为链判定阈值(冻结 dataclass,from_settings 读配置)。""" + + a_window_hours: int = 24 + a_count: int = 3 + b_window_hours: int = 72 + b_count: int = 5 + c_window_hours: int = 24 + c_count: int = 10 + + @classmethod + def from_settings(cls) -> "BehaviorThresholds": + return cls( + a_window_hours=settings.risk_agent_behavior_a_window_hours, + a_count=settings.risk_agent_behavior_a_count, + b_window_hours=settings.risk_agent_behavior_b_window_hours, + b_count=settings.risk_agent_behavior_b_count, + c_window_hours=settings.risk_agent_behavior_c_window_hours, + c_count=settings.risk_agent_behavior_c_count, + ) + + +def _new_alert_id() -> str: + from uuid import uuid4 + + return f"ALT-{datetime.now():%Y%m%d}-{uuid4().hex[:8].upper()}" + + +def _parse_dt(value: Any) -> datetime | None: + if isinstance(value, datetime): + return value + if isinstance(value, str): + try: + return datetime.fromisoformat(value) + except ValueError: + return None + return None + + +def _count_induce(events: list[dict[str, Any]]) -> int: + """同 (actor,customer) 组内诱导调仓计数:每个 redeem 找 2h 内不同产品的 subscribe。 + + events:[{trade_type, product_id, traded_at}];消费制(每个 subscribe 至多配对一次)。 + 窗口按 traded_at(事件时点,非扫描墙钟,与 RISK-004 rebuild 口径一致)。 + """ + redeems = sorted( + (e for e in events if e.get("trade_type") == "redeem"), + key=lambda e: e["traded_at"], + ) + subscribes = sorted( + (e for e in events if e.get("trade_type") == "subscribe"), + key=lambda e: e["traded_at"], + ) + used: set[int] = set() + count = 0 + for r in redeems: + for i, s in enumerate(subscribes): + if i in used: + continue + if s.get("product_id") == r.get("product_id"): + continue + delta = (s["traded_at"] - r["traded_at"]).total_seconds() + if 0 <= delta <= _INDUCE_WINDOW_SECONDS: + used.add(i) + count += 1 + break + return count + + +def detect_hits( + now: datetime, + core_ro: Any, + risk_repo: RiskRepository, + th: BehaviorThresholds, +) -> dict[str, dict[str, Any]]: + """按代理人分组聚合三条件证据(纯查询,无副作用,单测友好)。 + + 返回 {actor_id: {"A": [...], "B": [...], "C": [...], "customers": [...], "trace_ids": [...]}}。 + 每条证据含 {trace_id, at, customer_id, detail}。仅保留至少命中一条条件的 actor。 + """ + since_a = now - timedelta(hours=th.a_window_hours) + since_b = now - timedelta(hours=th.b_window_hours) + since_c = now - timedelta(hours=th.c_window_hours) + + # ---- 条件 A:代理人发起的交易请求(排除 SYSTEM 与本人交易) ---- + a_rows = risk_repo.list_audit_events(event_type="trade_request", since=since_a) + a_timeline: list[dict[str, Any]] = [] + for row in a_rows: + actor = row.get("actor_id") + customer = row.get("customer_id") + if not actor or actor == "SYSTEM" or not customer or actor == customer: + continue + summary = row.get("input_summary") or {} + trade_id = summary.get("trade_id") + if not trade_id: + continue + trade = core_ro.get_trade_by_id(trade_id) + if not trade: + continue + traded_at = _parse_dt(trade.get("traded_at")) + if traded_at is None: + continue + a_timeline.append( + { + "actor_id": actor, + "customer_id": customer, + "trade_type": trade.get("trade_type"), + "product_id": trade.get("product_id"), + "traded_at": traded_at, + "trace_id": row.get("trace_id"), + } + ) + + # ---- 条件 B / C:越权拒绝(authz forbidden) ---- + b_rows = risk_repo.list_audit_events(event_type="authz", since=since_b, decision="forbidden") + c_rows = risk_repo.list_audit_events(event_type="authz", since=since_c, decision="forbidden") + + by_actor: dict[str, dict[str, Any]] = defaultdict( + lambda: {"A": [], "B": [], "C": [], "customers": [], "trace_ids": []} + ) + + # 条件 A:按 (actor, customer) 分组计诱导 + a_by_ac: dict[tuple[str, str], list[dict[str, Any]]] = defaultdict(list) + for ev in a_timeline: + a_by_ac[(ev["actor_id"], ev["customer_id"])].append(ev) + for (actor, customer), evts in a_by_ac.items(): + n = _count_induce(evts) + if n <= 0: + continue + bucket = by_actor[actor] + bucket["A"].append( + { + "trace_id": evts[0].get("trace_id"), + "at": max((e["traded_at"] for e in evts), default=now), + "customer_id": customer, + "detail": {"induce_count": n, "trade_count": len(evts)}, + } + ) + if customer not in bucket["customers"]: + bucket["customers"].append(customer) + if evts[0].get("trace_id"): + bucket["trace_ids"].append(evts[0]["trace_id"]) + + # 条件 B:越权试探(AUTH_403_SCOPE) + for row in b_rows: + actor = row.get("actor_id") + code = (row.get("input_summary") or {}).get("code") + if not actor or actor == "SYSTEM": + continue + if code != "AUTH_403_SCOPE": + continue + bucket = by_actor[actor] + bucket["B"].append( + { + "trace_id": row.get("trace_id"), + "at": _parse_dt(row.get("created_at")) or now, + "customer_id": row.get("customer_id"), + "detail": {"code": code}, + } + ) + if row.get("customer_id") and row["customer_id"] not in bucket["customers"]: + bucket["customers"].append(row["customer_id"]) + if row.get("trace_id"): + bucket["trace_ids"].append(row["trace_id"]) + + # 条件 C:越权查询(NOT_OWNER / NOT_ASSIGNED) + for row in c_rows: + actor = row.get("actor_id") + code = (row.get("input_summary") or {}).get("code") + if not actor or actor == "SYSTEM": + continue + if code not in ("AUTH_403_NOT_OWNER", "AUTH_403_NOT_ASSIGNED"): + continue + bucket = by_actor[actor] + bucket["C"].append( + { + "trace_id": row.get("trace_id"), + "at": _parse_dt(row.get("created_at")) or now, + "customer_id": row.get("customer_id"), + "detail": {"code": code}, + } + ) + if row.get("customer_id") and row["customer_id"] not in bucket["customers"]: + bucket["customers"].append(row["customer_id"]) + if row.get("trace_id"): + bucket["trace_ids"].append(row["trace_id"]) + + # 阈值过滤:仅保留达到阈值的条件;至少一条条件达阈值的 actor 才进入结果 + # (detect_hits 返回的是"已触发"证据,阈值判定在此完成,scan_and_alert 仅消费) + result: dict[str, dict[str, Any]] = {} + for actor, bucket in by_actor.items(): + a_total = sum(e["detail"]["induce_count"] for e in bucket["A"]) + keep_a = bool(bucket["A"]) and a_total >= th.a_count + keep_b = len(bucket["B"]) >= th.b_count + keep_c = len(bucket["C"]) >= th.c_count + if not (keep_a or keep_b or keep_c): + continue + result[actor] = { + "A": bucket["A"] if keep_a else [], + "B": bucket["B"] if keep_b else [], + "C": bucket["C"] if keep_c else [], + "customers": bucket["customers"], + "trace_ids": bucket["trace_ids"], + } + return result + + +def _most_common(ids: list[str]) -> str | None: + if not ids: + return None + counts: dict[str, int] = defaultdict(int) + for i in ids: + counts[i] += 1 + return max(counts.items(), key=lambda kv: kv[1])[0] + + +def _audit_detected( + repo: RiskRepository, + trace_id: str, + actor_id: str, + customer_id: str, + subtypes_hit: list[str], + evidence: dict[str, Any], +) -> None: + """agent_behavior_detected 审计(仅 INSERT;写库失败降级不阻塞扫描,红线 5 口径)。""" + try: + repo.insert_audit_log( + { + "trace_id": trace_id, + "event_type": "agent_behavior_detected", + "agent_type": "risk", + "actor_id": actor_id, + "customer_id": customer_id or None, + "rule_id": "RISK-008", + "input_summary": { + "subtypes_hit": subtypes_hit, + "evidence_counts": {k: len(v) for k, v in evidence.items()}, + "note": "代理人维度行为链预警(与客户维度事件单独立出单)", + }, + "decision": "detected", + "risk_score": 75, + "handler_id": None, + "handler_result": None, + "handler_comment": None, + } + ) + except Exception: + logger.exception("agent_behavior detect audit failed (degraded)") + + +def scan_and_alert( + now: datetime | None = None, + core_ro: Any = None, + risk_repo: RiskRepository | None = None, + thresholds: BehaviorThresholds | None = None, +) -> dict[str, Any]: + """扫描入口(脚本/测试共用;模式对齐 escalation_service)。 + + 返回 {"scanned_windows", "hits", "created", "appended"}。 + hits 为逐 actor 命中摘要(含 subtypes_hit);created/appended 为出单/并入记录。 + """ + from app.repository.core_ro import CoreReadOnlyRepository + + repo = risk_repo or RiskRepository() + core = core_ro or CoreReadOnlyRepository() + th = thresholds or BehaviorThresholds.from_settings() + now = now or datetime.now() + trace_id = current_trace() or new_trace() + + hits = detect_hits(now, core, repo, th) + created: list[dict[str, Any]] = [] + appended: list[dict[str, Any]] = [] + hit_list: list[dict[str, Any]] = [] + day_start = datetime(now.year, now.month, now.day) + + for actor_id, bucket in hits.items(): + # detect_hits 已按阈值过滤:桶非空即代表该条件达阈值(A 为诱导总数、B/C 为事件数) + a_hit = bool(bucket["A"]) + b_hit = bool(bucket["B"]) + c_hit = bool(bucket["C"]) + subtypes_hit = [s for s, ok in (("A", a_hit), ("B", b_hit), ("C", c_hit)) if ok] + if not subtypes_hit: + continue + hit_list.append({"actor_id": actor_id, "subtypes_hit": subtypes_hit}) + + customer_id = _most_common(bucket["customers"]) or "" + evidence = {"A": bucket["A"], "B": bucket["B"], "C": bucket["C"]} + existing = repo.find_agent_behavior_alert(actor_id, day_start) + if existing: + # 同日同 actor 仅一张单:证据并集(events + alert_subtype 由 append 合并, + # 行为链专属标量字段由 merge 合并) + repo.append_alert_event( + existing["alert_id"], evidence, ["RISK-008"], 75, "pattern", + extra_subtypes=["agent_behavior"], + ) + repo.merge_agent_behavior_payload( + existing["alert_id"], + subtypes_hit=subtypes_hit, + customers=bucket["customers"], + evidence_trace_ids=bucket["trace_ids"], + timeline=evidence, + ) + _audit_detected(repo, trace_id, actor_id, customer_id, subtypes_hit, evidence) + appended.append( + {"alert_id": existing["alert_id"], "actor_id": actor_id, "subtypes_hit": subtypes_hit} + ) + continue + + alert = { + "alert_id": _new_alert_id(), + "trace_id": trace_id, + "customer_id": customer_id, + "trade_id": None, + "alert_type": "pattern", + "triggered_rules": ["RISK-008"], + "risk_score": 75, + "status": "pending_review", + "payload": { + "alert_subtype": ["agent_behavior"], + "actor_id": actor_id, + "actor_type": "agent", + "subtypes_hit": subtypes_hit, + "timeline": evidence, + "customers": bucket["customers"], + "evidence_trace_ids": bucket["trace_ids"], + "note": "customer_id 为涉及客户(多客户取众数)非归属客户", + }, + } + repo.insert_alert(alert) + _audit_detected(repo, trace_id, actor_id, customer_id, subtypes_hit, evidence) + _publish_alert( + alert, + notify_role=["risk_officer", "risk_manager"], + extra={"alert_subtype": "agent_behavior"}, + ) + created.append( + {"alert_id": alert["alert_id"], "actor_id": actor_id, "subtypes_hit": subtypes_hit} + ) + + return { + "scanned_windows": len(hits), + "hits": hit_list, + "created": created, + "appended": appended, + } diff --git a/app/service/risk/chat_tools.py b/app/service/risk/chat_tools.py index 666f077..2ec6b14 100644 --- a/app/service/risk/chat_tools.py +++ b/app/service/risk/chat_tools.py @@ -30,6 +30,7 @@ from app.repository.core_ro import CoreReadOnlyRepository from app.repository.risk_repository import RiskRepository from app.service import suitability from app.tool.core_tools import _jsonable +from app.utils.desensitize import mask_name class RiskToolSpec(dict): @@ -262,6 +263,68 @@ def query_overdue_alerts(customer_id: str, core_ro: CoreReadOnlyRepository | Non # ---------- 注册表(供 tool_service.get_registered_tool 分发) ---------- +def _mask_customer(customer_id: str, core_ro: CoreReadOnlyRepository) -> dict[str, Any] | None: + """涉及客户脱敏展示(PRD 系统动作 6):customer_id 内部键不脱敏,display_name 走 mask_name。""" + ro = core_ro or CoreReadOnlyRepository() + l0 = ro.get_customer_l0(customer_id) + if not l0: + return {"customer_id": customer_id, "display_name": None} + name = l0.get("display_name") + return {"customer_id": customer_id, "display_name": mask_name(name) if name else None} + + +def query_agent_behavior(customer_id: str, core_ro: CoreReadOnlyRepository | None = None, + risk_repo: RiskRepository | None = None, **params: Any) -> dict[str, Any]: + """代理人行为链查询(只读;FR-10 行为链预警查询入口)。 + + agent_id = params.get("agent_id")(可缺省=全部)。过滤 payload.alert_subtype + 含 agent_behavior 的 pattern 单;agent_id 传入时再按 payload.actor_id 精确匹配。 + customers 经 utils/desensitize.mask_name 二次脱敏(PRD 系统动作 6)。 + 可见性:对话线准入已限 risk_officer(manager 走 HTTP 台账),Tool 层天然 fail-closed。 + """ + repo = risk_repo or RiskRepository() + ro = core_ro or CoreReadOnlyRepository() + agent_id = params.get("agent_id") + rows, _ = repo.list_alerts(alert_type="pattern", status=None, page_size=100) + out: list[dict[str, Any]] = [] + for a in rows: + payload = a.get("payload") or {} + sub = payload.get("alert_subtype") + if isinstance(sub, str): + sub = [sub] + if "agent_behavior" not in (sub or []): + continue + # customer_id 为涉及客户众数(scan_and_alert 写入),对话线按客户维度收敛 + if customer_id and a.get("customer_id") != customer_id: + continue + if agent_id and payload.get("actor_id") != agent_id: + continue + customers_raw = payload.get("customers") or [] + masked = [_mask_customer(cid, ro) for cid in customers_raw] + out.append( + { + "alert_id": a.get("alert_id"), + "actor_id": payload.get("actor_id"), + "subtypes_hit": payload.get("subtypes_hit"), + "customers": masked, + "created_at": a.get("created_at").isoformat() + if isinstance(a.get("created_at"), (_dt.datetime, _dt.date)) + else a.get("created_at"), + "risk_score": a.get("risk_score"), + "status": a.get("status"), + "evidence_trace_ids": (payload.get("evidence_trace_ids") or [])[:5], + } + ) + return _jsonable( + { + "scope": customer_id or None, + "agent_id": agent_id, + "count": len(out), + "items": out, + } + ) + + RISK_TOOL_REGISTRY: dict[str, RiskToolSpec] = { "query_overdue_alerts": RiskToolSpec( func=query_overdue_alerts, diff --git a/app/service/tool_service.py b/app/service/tool_service.py index ff0370d..4a4b62f 100644 --- a/app/service/tool_service.py +++ b/app/service/tool_service.py @@ -86,6 +86,8 @@ _INTENT_KEYWORDS: dict[str, list[tuple[str, tuple[str, ...]]]] = { "risk": [ # FR-9 时效升级查询:置于 alert_query 之前("超时/超期预警"话术不得误命中台账) ("query_overdue_alerts", ("超期", "超时", "逾期", "多久没处理", "处置时效")), + # FR-10 行为链查询:置于 alert_query 之前("代理人预警"应命中行为链而非台账) + ("query_agent_behavior", ("代理人", "行为链", "异常行为", "诱导", "越权记录")), ("alert_query", ("预警", "待审", "预警台账", "待审预警")), ("customer_context", ("客户上下文", "监测信息", "风险画像", "客户监测")), ("suitability_check", ("适当性", "能不能买", "可购", "适合买", "购买资格")), @@ -392,6 +394,15 @@ def summarize(record: dict[str, Any]) -> str: f"(客户 {data.get('customer_id')} 待审预警 {pending} 条," f"今日 {today} 条;仅供参考,处置须经风控专员人工完成)" ) + if name == "query_agent_behavior": + items = data.get("items") or [] + if not items: + return "(代理人行为链查询:无命中记录)" + lines = [f"(代理人行为链:共 {len(items)} 条记录)"] + for it in items[:5]: + subs = "/".join(it.get("subtypes_hit") or []) + lines.append(f"- 代理人 {it.get('actor_id')}:命中 {subs}({it.get('status')})") + return "\n".join(lines) if name == "customer_context": if not data.get("found"): return f"(客户风控上下文:未找到客户 {data.get('customer_id')})" diff --git a/scripts/cron/agent_behavior_scan.py b/scripts/cron/agent_behavior_scan.py new file mode 100644 index 0000000..1929929 --- /dev/null +++ b/scripts/cron/agent_behavior_scan.py @@ -0,0 +1,32 @@ +"""RISK-008 代理人行为链扫描(建议 30min 周期;RISK_AGENT_BEHAVIOR_SCAN_MINUTES 可配)。 + +初期独立脚本 + 系统 cron;结构对齐 scripts/cron/escalation_scan.py: +sys.path 引导 → new_trace → 扫描 → 打印 JSON 摘要(供 cron 日志 / 演示走查)。 +审计写库失败由 service 内降级留底,不阻塞下次扫描(红线 5 口径)。 +""" + +from __future__ import annotations + +import json +import sys +from pathlib import Path + +# ① sys.path 引导项目根(escalation_scan 先例) +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) + +from app.service.risk.agent_behavior_service import scan_and_alert # noqa: E402 +from app.utils.trace import new_trace # noqa: E402 + + +def main() -> int: + # ② 显式生成 trace(无 HTTP 上下文,service 内 current_trace 复用) + new_trace() + # ③ 扫描 → 打印 JSON 摘要 + summary = scan_and_alert() + print(json.dumps(summary, ensure_ascii=False, default=str, indent=2)) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tests/conftest.py b/tests/conftest.py index 51671d5..6fae557 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -108,6 +108,64 @@ def backdated_alert(sqlite_engine): ) +@pytest.fixture() +def backdated_audit_event(sqlite_engine): + """注入 created_at 回拨的 audit_log 行(C6);teardown 按 trace_id(TEST-TRACE- 前缀)精确删除。 + + 用法: + tid = backdated_audit_event(event_type="authz", actor_id="STAFF-X", + customer_id="CUST-Y", + input_summary={"code": "AUTH_403_SCOPE"}, + decision="forbidden", hours_ago=5) + 注意:注入与读取必须共用同一个 sqlite_engine;回拨行 teardown 按 trace_id 前缀 + 精确删除,不污染正常 trace 空间(评审必查点,_cleanup_test_rows 时间窗清不掉回拨行)。 + """ + import json + from uuid import uuid4 + + created_trace_ids: list[str] = [] + + def _make( + event_type: str, + actor_id: str, + customer_id: str | None = None, + input_summary: dict | None = None, + decision: str = "forbidden", + hours_ago: float = 5.0, + agent_type: str = "risk", + trace_id: str | None = None, + ) -> str: + tid = trace_id or f"TEST-TRACE-{uuid4().hex[:12].upper()}" + created_at = datetime.now() - timedelta(hours=hours_ago) + with sqlite_engine.begin() as conn: + conn.execute( + text( + "INSERT INTO audit_log (trace_id, event_type, agent_type, actor_id," + " customer_id, rule_id, input_summary, decision, risk_score," + " handler_id, handler_result, handler_comment, created_at)" + " VALUES (:tid, :et, :agent_type, :aid, :cid, NULL, :summary," + " :decision, NULL, NULL, NULL, NULL, :created_at)" + ), + { + "tid": tid, + "et": event_type, + "agent_type": agent_type, + "aid": actor_id, + "cid": customer_id, + "summary": json.dumps(input_summary or {}, ensure_ascii=False), + "decision": decision, + "created_at": created_at, + }, + ) + created_trace_ids.append(tid) + return tid + + yield _make + with sqlite_engine.begin() as conn: + for tid in created_trace_ids: + conn.execute(text("DELETE FROM audit_log WHERE trace_id = :tid"), {"tid": tid}) + + # ---------- 真库集成环境(B8 集成测试专用) ---------- diff --git a/tests/test_agent_behavior_service.py b/tests/test_agent_behavior_service.py new file mode 100644 index 0000000..fade1bb --- /dev/null +++ b/tests/test_agent_behavior_service.py @@ -0,0 +1,254 @@ +"""C6 · RISK-008 代理人行为链单测(A-12 验收覆盖)。 + +纯函数 _count_induce 边界 + detect_hits 三条件 + scan_and_alert 出单/同日去重 + +payload.actor_id 归属 + query_agent_behavior Tool 过滤/脱敏。 +数据源:单 sqlite 引擎同时承载 Core 与 Agent 表(_ddl 已含 core_trade/audit_log/ +risk_alert),单测共用 RiskRepository(engine) + CoreReadOnlyRepository(engine)。 +""" + +from __future__ import annotations + +from datetime import datetime, timedelta + +from sqlalchemy import text + +from app.repository.core_ro import CoreReadOnlyRepository +from app.repository.risk_repository import RiskRepository +from app.service.risk.agent_behavior_service import ( + BehaviorThresholds, + _count_induce, + detect_hits, + scan_and_alert, +) +from app.service.risk.chat_tools import query_agent_behavior + + +def _seed_trade(engine, trade_id, customer_id, trade_type, product_id, traded_at) -> None: + with engine.begin() as conn: + conn.execute( + text( + "INSERT INTO core_trade (trade_id, customer_id, product_id, trade_type," + " amount, trade_status, traded_at) VALUES (:tid, :cid, :pid, :tt, 1000," + " 'confirmed', :ta)" + ), + { + "tid": trade_id, "cid": customer_id, "pid": product_id, + "tt": trade_type, "ta": traded_at, + }, + ) + + +def _seed_trade_request(make, trade_id, actor_id, customer_id, hours_ago=1.0) -> None: + make( + event_type="trade_request", actor_id=actor_id, customer_id=customer_id, + input_summary={"trade_id": trade_id}, decision="trade_accepted", hours_ago=hours_ago, + ) + + +# ---------- _count_induce 纯函数边界 ---------- + +def test_count_induce_one_pair(): + evts = [ + {"trade_type": "redeem", "product_id": "P1", "traded_at": datetime(2026, 1, 1, 0, 0)}, + {"trade_type": "subscribe", "product_id": "P2", "traded_at": datetime(2026, 1, 1, 0, 30)}, + {"trade_type": "redeem", "product_id": "P3", "traded_at": datetime(2026, 1, 1, 1, 0)}, + ] + assert _count_induce(evts) == 1 # 仅第一对配对,第二个 redeem 无 subscribe + + +def test_count_induce_three_pairs(): + base = datetime(2026, 1, 1, 0, 0) + evts = [] + for i in range(3): + r = base + timedelta(minutes=10 * i) + s = r + timedelta(minutes=30) + evts.append({"trade_type": "redeem", "product_id": "P1", "traded_at": r}) + evts.append({"trade_type": "subscribe", "product_id": "P2", "traded_at": s}) + assert _count_induce(evts) == 3 + + +def test_count_induce_same_product_not_count(): + evts = [ + {"trade_type": "redeem", "product_id": "P1", "traded_at": datetime(2026, 1, 1, 0, 0)}, + {"trade_type": "subscribe", "product_id": "P1", "traded_at": datetime(2026, 1, 1, 0, 30)}, + ] + assert _count_induce(evts) == 0 + + +def test_count_induce_window_boundary(): + r = datetime(2026, 1, 1, 0, 0) + s_in = r + timedelta(seconds=7200) # 恰好 2h 不同产品 → 计入(含边界) + s_out = r + timedelta(seconds=7201) # 超 2h 1s → 不计 + evts_in = [ + {"trade_type": "redeem", "product_id": "P1", "traded_at": r}, + {"trade_type": "subscribe", "product_id": "P2", "traded_at": s_in}, + ] + evts_out = [ + {"trade_type": "redeem", "product_id": "P1", "traded_at": r}, + {"trade_type": "subscribe", "product_id": "P2", "traded_at": s_out}, + ] + assert _count_induce(evts_in) == 1 + assert _count_induce(evts_out) == 0 + + +# ---------- 条件 A:本人 / SYSTEM 排除 + 3 次触发 ---------- + +def test_condition_a_self_and_system_excluded(sqlite_engine, backdated_audit_event): + now = datetime.now() + # 本人交易:actor_id == customer_id + _seed_trade(sqlite_engine, "TRD-TEST-SELF-R", "CUST-1", "redeem", "P1", now - timedelta(hours=1)) + _seed_trade(sqlite_engine, "TRD-TEST-SELF-S", "CUST-1", "subscribe", "P2", now - timedelta(minutes=30)) + backdated_audit_event( + event_type="trade_request", actor_id="CUST-1", customer_id="CUST-1", + input_summary={"trade_id": "TRD-TEST-SELF-R"}, hours_ago=1, + ) + # SYSTEM 无归属 + _seed_trade(sqlite_engine, "TRD-TEST-SYS-R", "CUST-2", "redeem", "P1", now - timedelta(hours=1)) + _seed_trade(sqlite_engine, "TRD-TEST-SYS-S", "CUST-2", "subscribe", "P2", now - timedelta(minutes=30)) + backdated_audit_event( + event_type="trade_request", actor_id="SYSTEM", customer_id="CUST-2", + input_summary={"trade_id": "TRD-TEST-SYS-R"}, hours_ago=1, + ) + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + assert detect_hits(now, core, repo, BehaviorThresholds()) == {} + + +def test_condition_a_three_triggers(sqlite_engine, backdated_audit_event): + now = datetime.now() + base = now - timedelta(hours=1) + for i in range(3): + r_t = base + timedelta(minutes=10 * i) + s_t = r_t + timedelta(minutes=30) + _seed_trade(sqlite_engine, f"TRD-TEST-A{i}-R", "CUST-1", "redeem", "P1", r_t) + _seed_trade(sqlite_engine, f"TRD-TEST-A{i}-S", "CUST-1", "subscribe", "P2", s_t) + # 代理人发起赎回 + 申购两笔请求(a_timeline 需同时含 redeem/subscribe 才能配对诱导) + _seed_trade_request(backdated_audit_event, f"TRD-TEST-A{i}-R", "STAFF-A", "CUST-1", hours_ago=1) + _seed_trade_request(backdated_audit_event, f"TRD-TEST-A{i}-S", "STAFF-A", "CUST-1", hours_ago=1) + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + hits = detect_hits(now, core, repo, BehaviorThresholds()) + assert "STAFF-A" in hits + assert hits["STAFF-A"]["A"][0]["detail"]["induce_count"] == 3 + + +# ---------- 条件 B / C 边界 ---------- + +def test_condition_b_boundary(sqlite_engine, backdated_audit_event): + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + now = datetime.now() + for _ in range(4): # 4 次不触发 + backdated_audit_event( + event_type="authz", actor_id="STAFF-B", customer_id="CUST-1", + input_summary={"code": "AUTH_403_SCOPE"}, decision="forbidden", hours_ago=1, + ) + assert "STAFF-B" not in detect_hits(now, core, repo, BehaviorThresholds()) + backdated_audit_event( # 第 5 次触发 + event_type="authz", actor_id="STAFF-B", customer_id="CUST-1", + input_summary={"code": "AUTH_403_SCOPE"}, decision="forbidden", hours_ago=1, + ) + hits = detect_hits(now, core, repo, BehaviorThresholds()) + assert "STAFF-B" in hits + assert len(hits["STAFF-B"]["B"]) == 5 + + +def test_condition_c_boundary_and_union(sqlite_engine, backdated_audit_event): + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + now = datetime.now() + for i in range(9): # 9 次(双 code 并集)不触发 + code = "AUTH_403_NOT_OWNER" if i % 2 == 0 else "AUTH_403_NOT_ASSIGNED" + backdated_audit_event( + event_type="authz", actor_id="STAFF-C", customer_id="CUST-1", + input_summary={"code": code}, decision="forbidden", hours_ago=1, + ) + assert "STAFF-C" not in detect_hits(now, core, repo, BehaviorThresholds()) + backdated_audit_event( # 第 10 次触发 + event_type="authz", actor_id="STAFF-C", customer_id="CUST-1", + input_summary={"code": "AUTH_403_NOT_OWNER"}, decision="forbidden", hours_ago=1, + ) + hits = detect_hits(now, core, repo, BehaviorThresholds()) + assert "STAFF-C" in hits + assert len(hits["STAFF-C"]["C"]) == 10 + + +def test_condition_b_wrong_code_excluded(sqlite_engine, backdated_audit_event): + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + now = datetime.now() + for _ in range(5): + backdated_audit_event( + event_type="authz", actor_id="STAFF-X", customer_id="CUST-1", + input_summary={"code": "AUTH_403_OTHER"}, decision="forbidden", hours_ago=1, + ) + assert "STAFF-X" not in detect_hits(now, core, repo, BehaviorThresholds()) + + +# ---------- scan_and_alert 出单 / 同日去重 / payload ---------- + +def _seed_agent_behavior_scenario(engine, make, actor_id="STAFF-A", customer_id="CUST-1") -> None: + now = datetime.now() + base = now - timedelta(hours=1) + for i in range(3): + r_t = base + timedelta(minutes=10 * i) + s_t = r_t + timedelta(minutes=30) + _seed_trade(engine, f"TRD-TEST-Z{i}-R", customer_id, "redeem", "P1", r_t) + _seed_trade(engine, f"TRD-TEST-Z{i}-S", customer_id, "subscribe", "P2", s_t) + # 代理人发起赎回 + 申购两笔请求 + _seed_trade_request(make, f"TRD-TEST-Z{i}-R", actor_id, customer_id, hours_ago=1) + _seed_trade_request(make, f"TRD-TEST-Z{i}-S", actor_id, customer_id, hours_ago=1) + + +def test_scan_creates_one_alert_per_actor_per_day(sqlite_engine, backdated_audit_event): + _seed_agent_behavior_scenario(sqlite_engine, backdated_audit_event) + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + r1 = scan_and_alert(core_ro=core, risk_repo=repo) + assert len(r1["created"]) == 1 + assert len(r1["appended"]) == 0 + alert_id = r1["created"][0]["alert_id"] + # 同日二次扫描:并入不新建 + r2 = scan_and_alert(core_ro=core, risk_repo=repo) + assert len(r2["created"]) == 0 + assert len(r2["appended"]) == 1 + alert = repo.get_alert(alert_id) + assert alert["payload"]["alert_subtype"] == ["agent_behavior"] + assert alert["payload"]["actor_id"] == "STAFF-A" + assert alert["payload"]["subtypes_hit"] == ["A"] + + +def test_payload_actor_id_points_to_agent(sqlite_engine, backdated_audit_event): + _seed_agent_behavior_scenario( + sqlite_engine, backdated_audit_event, actor_id="STAFF-AGENT-1", customer_id="CUST-9" + ) + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + r = scan_and_alert(core_ro=core, risk_repo=repo) + alert = repo.get_alert(r["created"][0]["alert_id"]) + assert alert["payload"]["actor_id"] == "STAFF-AGENT-1" + assert alert["customer_id"] == "CUST-9" # 涉及客户众数(仅一个) + assert alert["risk_score"] == 75 + assert alert["triggered_rules"] == ["RISK-008"] + + +# ---------- query_agent_behavior Tool ---------- + +def test_query_agent_behavior_tool_filters_and_masks(sqlite_engine, backdated_audit_event): + _seed_agent_behavior_scenario( + sqlite_engine, backdated_audit_event, actor_id="STAFF-Q", customer_id="CUST-7" + ) + repo = RiskRepository(engine=sqlite_engine) + core = CoreReadOnlyRepository(engine=sqlite_engine) + scan_and_alert(core_ro=core, risk_repo=repo) + # 不带 agent_id:返回该单 + out_all = query_agent_behavior("CUST-7", core_ro=core, risk_repo=repo) + assert out_all["count"] == 1 + assert out_all["items"][0]["actor_id"] == "STAFF-Q" + assert out_all["items"][0]["customers"][0]["customer_id"] == "CUST-7" + # 错误 agent_id:过滤为空 + out_none = query_agent_behavior("CUST-7", core_ro=core, risk_repo=repo, agent_id="STAFF-OTHER") + assert out_none["count"] == 0 + # 正确 agent_id:命中 + out_hit = query_agent_behavior("CUST-7", core_ro=core, risk_repo=repo, agent_id="STAFF-Q") + assert out_hit["count"] == 1