Compare commits

...

6 Commits

Author SHA1 Message Date
guke 1e7f6024c7 Merge main into feat/aliyun-sms-verify
冲突仅 app/api/v1/auth.py::sms_login,合并两侧改动:
  - 保留(本分支)阿里云 provider 的 verify_code try/except —— 校验降级抛 SmsError 时透出其状态码(503),不误报 400。
  - 保留(main)风控失败事件 record_behavior_event —— 但仅在「验证码错误」(ok 为 False)时记;
    provider 降级 503 已在 try/except 提前 raise,不计入「验证失败」风控。
wechat_bind_phone_sms 的 try/except 自动合并、无冲突(main 未在该处加风控)。

验证:全量 pytest 552 passed / 8 failed(8 个均为合并前已存在、与本次无关)。
2026-07-26 00:42:31 +08:00
guke 7a2b7cb8ed 功能:阿里云号码认证短信 provider + sms.py 拆为可切换 provider 包
现有 sms.py(极光自管码)升级为 app/integrations/sms/ 包:
  - base   : SmsError + provider 无关的 mock_verify
  - jiguang: 原极光自管码逻辑逐字迁入,行为零改动(git 识别为 sms.py 的 rename)
  - aliyun : 新增阿里云 dypns 号码认证(Mode A:阿里云生成+下发+校验,核验免费)
  - __init__: 按 SMS_PROVIDER 每次调用路由的分派器(默认 jiguang,可秒切回退)

关键决策:
  - Mode A:发码 SendSmsVerifyCode(##code## 占位)、校验 CheckSmsVerifyCode(PASS/UNKNOWN);
    本服务不再存码 → 消除极光路径「内存存码、多 worker 不共享」技术债。
  - 防爆破与极光一致:aliyun 保留 per-phone 失败计数(SMS_MAX_VERIFY_ATTEMPTS),达上限本地作废,
    避免两 provider 行为不同致排查困惑(此为 aliyun 路径唯一本地态)。
  - 校验降级:阿里云接口异常 → verify_code 抛 SmsError(503),auth 两处 try/except 透出 503(非误报 400)。
  - 官方 SDK alibabacloud_dypnsapi20170525;SDK 交互隔离在 _call_send/_call_check(惰性 import + 惰性建
    client),单测 monkeypatch 不触真网络。

测试:test_sms_aliyun(17)+ test_sms_dispatch(3)全绿;test_auth 内部访问 retarget 到 jiguang.*。
配置:SMS_PROVIDER + ALIYUN_SMS_*(见 .env.example);文档 docs/integrations/sms.md + aliyun 接口参考。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-25 23:51:30 +08:00
guke deb7730773 docs: 阿里云短信验证服务设计spec
Mode A(阿里云托管码) + 可切换 provider(SMS_PROVIDER) + 官方 SDK;
防爆破保留与极光一致的单码失败计数,校验降级返回 503。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-25 17:58:44 +08:00
linkeyu 7bf04f2655 功能:新增风控监控与处置能力 (#174)
## 修改内容
- 新增短信、一键登录、比价三类风控事件与规则聚合
- 新增管理员忽略、封禁、解封、重置报警和阈值配置接口
- 风险列表返回当前有效限制的 restriction_id,供后台已封禁视图解除
- 在短信/一键登录、比价、任务领奖、提现链路接入风险记录与限制
- 新增通用行为流水、风险事件、主体限制模型及 Alembic 迁移
- 新增本地演示数据脚本与风控测试

## 验证
- ruff check:通过
- Alembic 全新 SQLite upgrade head / downgrade -1:通过
- 风控测试:通过,覆盖封禁列表 restriction_id 与解除链路
- 全量测试:532 passed,8 failed;其中 7 项在干净 origin/main 独立复现,另 1 项独立复跑通过,未发现本分支新增回归

---------

Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #174
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-07-25 17:56:23 +08:00
marco 5f6593eaf2 fix(alembic): 缩短 withdraw 迁移 revision id 至 varchar(32) 内
drop_withdraw_active_unique_index (33 字符) 超过 alembic_version.version_num
的 varchar(32) 上限,生产 alembic 写版本号时 StringDataRightTruncation 报错、
部署中断。改短为 drop_withdraw_active_uniq_idx (29 字符):迁移逻辑(drop_index)
一字不变、无下游引用、仍为唯一 head。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-24 20:32:38 +08:00
guke e11f506e1e docs: 提现允许在途时继续提交申请 设计spec (#173)
## 摘要
取消「同一用户同时仅一笔在途提现」限制:已有 reviewing/pending 提现单时可继续发起新申请。

- 删应用层在途单检查(WithdrawTooFrequentError)+ 删 DB 分区唯一索引 ux_withdraw_order_user_active(含迁移)
- 清理失效死代码;IntegrityError 兜底瘦身为仅处理 out_bill_no 幂等
- 既有约束不变:建单先扣款(防超提)、coin_cash 每日档位次数、out_bill_no 幂等、解绑退款、admin 审核/对账均按单号维度

## 测试
- 新增:多笔在途并存放行(coin_cash & invite_cash)、第二笔仅受余额约束(409 现金余额不足)
- 迁移 upgrade→downgrade→upgrade 回环验证
- 提现域全绿(test_withdraw / test_invite_cash_withdraw / test_withdraw_ledger_check)

## 注意
- 客户端:每次提交需生成新的 out_bill_no;未开免确认时多笔 pending 各返回一个微信确认页,App 需能处理多笔待确认
- 无并发硬上限(产品拍板):coin_cash 由每日档位次数天然封顶,invite_cash 仅受余额约束

---------

Co-authored-by: guke <guke@autohome.com.cn>
Reviewed-on: #173
2026-07-24 16:40:57 +08:00
42 changed files with 5039 additions and 132 deletions
+18 -2
View File
@@ -81,12 +81,28 @@ HEARTBEAT_TIMEOUT_MINUTES=60
HEARTBEAT_SCAN_INTERVAL_SEC=60
# ===== 短信 (mock 模式) =====
# mock = true 时,任意 6 位数字均通过,且 /sms/send 不真发短信(只 log)。
# 后续接阿里云/腾讯云短信时,改成 false 并填供应商相关 key。
# mock = true 时,任意 6 位数字均通过,且 /sms/send 不真发短信(只 log)。生产改 false。
SMS_MOCK=true
SMS_CODE_TTL_SEC=300
SMS_SEND_INTERVAL_SEC=60
# ===== 短信提供商(可切换:jiguang 默认 / aliyun 阿里云号码认证)=====
# jiguang:本服务生成验证码,极光 REST 只负责下发,本地内存校验(复用上面极光 JG_* 凭证)。
# aliyun :阿里云 dypns 号码认证,阿里云生成+下发+校验(Mode A,核验免费);缺凭证时 /sms/* 返 503。
# 需在阿里云号码认证控制台开通「融合认证」,并使用系统赠送签名 + 赠送模板。
SMS_PROVIDER=jiguang
ALIYUN_SMS_ACCESS_KEY_ID=
ALIYUN_SMS_ACCESS_KEY_SECRET=
ALIYUN_SMS_SIGN_NAME=
ALIYUN_SMS_TEMPLATE_CODE=
# 方案名:留空=默认方案;若填,发码与校验须一致(本服务已共用同一配置项,不会不匹配)。
ALIYUN_SMS_SCHEME_NAME=
ALIYUN_SMS_ENDPOINT=dypnsapi.aliyuncs.com
ALIYUN_SMS_CODE_LENGTH=6
ALIYUN_SMS_VALID_TIME_SEC=300
ALIYUN_SMS_INTERVAL_SEC=60
ALIYUN_SMS_TIMEOUT_SEC=15
# ===== 测试账号(release 包全流程联调用)=====
# 配一个固定测试手机号,专供无 SIM 卡 / 不走一键登录时打通全流程:该号登录【免短信验证码】
# (real 模式下也跳过校验)、每次登录【都重走新手引导】,并有【每日登录上限】防被人猜到号后脚本刷。
@@ -0,0 +1,37 @@
"""drop withdraw active-order partial unique index (allow multiple in-flight withdrawals)
Revision ID: drop_withdraw_active_uniq_idx
Revises: d8dd2106e438
Create Date: 2026-07-24 00:00:00.000000
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
# revision identifiers, used by Alembic.
revision: str = "drop_withdraw_active_uniq_idx"
down_revision: Union[str, Sequence[str], None] = "d8dd2106e438"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# 取消「同一用户同一时刻仅一笔在途提现」:允许 reviewing/pending 并存。
# 仅删本索引;姊妹索引 ux_cash_transaction_withdraw_refund_ref(退款幂等)保持不动。
op.drop_index("ux_withdraw_order_user_active", table_name="withdraw_order")
def downgrade() -> None:
# 回滚重建分区唯一索引。注意:若届时某用户已有 ≥2 张在途单,重建会因唯一冲突失败——
# 属预期的回滚代价(取消限制后本就允许多单),需先人工收敛在途单再回滚。
op.create_index(
"ux_withdraw_order_user_active",
"withdraw_order",
["user_id"],
unique=True,
sqlite_where=sa.text("status IN ('reviewing', 'pending')"),
postgresql_where=sa.text("status IN ('reviewing', 'pending')"),
)
@@ -0,0 +1,210 @@
"""add generic behavior, risk incident and subject restriction tables
Revision ID: risk_monitor_generic
Revises: drop_withdraw_active_uniq_idx
"""
from __future__ import annotations
from collections.abc import Sequence
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
from alembic import op
revision: str = "risk_monitor_generic"
down_revision: str | None = "drop_withdraw_active_uniq_idx"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
_RISK_MONITOR_PAGE = "risk-monitor"
_DEFAULT_RISK_MONITOR_ROLES = ("operator", "tech")
def _role_table() -> sa.TableClause:
return sa.table(
"admin_role",
sa.column("name", sa.String),
sa.column("pages", _JSON),
)
def _add_default_role_permissions() -> None:
"""Grant the new page without replacing any existing role customisation."""
role = _role_table()
conn = op.get_bind()
rows = conn.execute(
sa.select(role.c.name, role.c.pages).where(
role.c.name.in_(_DEFAULT_RISK_MONITOR_ROLES)
)
).all()
for name, pages in rows:
current_pages = list(pages or [])
if _RISK_MONITOR_PAGE not in current_pages:
conn.execute(
role.update()
.where(role.c.name == name)
.values(pages=[*current_pages, _RISK_MONITOR_PAGE])
)
def _remove_default_role_permissions() -> None:
role = _role_table()
conn = op.get_bind()
rows = conn.execute(
sa.select(role.c.name, role.c.pages).where(
role.c.name.in_(_DEFAULT_RISK_MONITOR_ROLES)
)
).all()
for name, pages in rows:
current_pages = list(pages or [])
if _RISK_MONITOR_PAGE in current_pages:
conn.execute(
role.update()
.where(role.c.name == name)
.values(
pages=[
page
for page in current_pages
if page != _RISK_MONITOR_PAGE
]
)
)
def upgrade() -> None:
op.create_table(
"behavior_event",
sa.Column("id", sa.Integer(), autoincrement=True, nullable=False),
sa.Column("event_type", sa.String(length=64), nullable=False),
sa.Column("subject_type", sa.String(length=32), nullable=False),
sa.Column("subject_id", sa.String(length=128), nullable=False),
sa.Column("user_id", sa.Integer(), nullable=True),
sa.Column("device_id", sa.String(length=128), nullable=True),
sa.Column("device_model", sa.String(length=128), nullable=True),
sa.Column("phone", sa.String(length=20), nullable=True),
sa.Column("client_ip", sa.String(length=64), nullable=True),
sa.Column("outcome", sa.String(length=24), nullable=False),
sa.Column("reason", sa.String(length=256), nullable=True),
sa.Column("details", sa.JSON(), nullable=True),
sa.Column("occurred_at", sa.DateTime(timezone=True), nullable=False),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("(CURRENT_TIMESTAMP)"),
nullable=False,
),
sa.PrimaryKeyConstraint("id"),
)
op.create_index("ix_behavior_event_occurred_at", "behavior_event", ["occurred_at"])
op.create_index(
"ix_behavior_event_type_time", "behavior_event", ["event_type", "occurred_at"]
)
op.create_index(
"ix_behavior_event_subject_time",
"behavior_event",
["subject_type", "subject_id", "occurred_at"],
)
op.create_index(
"ix_behavior_event_user_time", "behavior_event", ["user_id", "occurred_at"]
)
op.create_table(
"risk_incident",
sa.Column("id", sa.Integer(), autoincrement=True, nullable=False),
sa.Column("rule_code", sa.String(length=64), nullable=False),
sa.Column("event_type", sa.String(length=64), nullable=False),
sa.Column("subject_type", sa.String(length=32), nullable=False),
sa.Column("subject_id", sa.String(length=128), nullable=False),
sa.Column("window_key", sa.String(length=64), nullable=False),
sa.Column("window_start", sa.DateTime(timezone=True), nullable=False),
sa.Column("window_end", sa.DateTime(timezone=True), nullable=False),
sa.Column("first_event_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("triggered_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("last_event_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("event_count", sa.Integer(), nullable=False),
sa.Column("status", sa.String(length=24), nullable=False),
sa.Column("action_reason", sa.String(length=256), nullable=True),
sa.Column("handled_by", sa.Integer(), nullable=True),
sa.Column("handled_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("details", sa.JSON(), nullable=True),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("(CURRENT_TIMESTAMP)"),
nullable=False,
),
sa.Column(
"updated_at",
sa.DateTime(timezone=True),
server_default=sa.text("(CURRENT_TIMESTAMP)"),
nullable=False,
),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint(
"rule_code",
"subject_type",
"subject_id",
"window_key",
name="uq_risk_incident_rule_subject_window",
),
)
op.create_index(
"ix_risk_incident_rule_status",
"risk_incident",
["rule_code", "status", "triggered_at"],
)
op.create_index(
"ix_risk_incident_subject",
"risk_incident",
["subject_type", "subject_id", "triggered_at"],
)
op.create_table(
"subject_restriction",
sa.Column("id", sa.Integer(), autoincrement=True, nullable=False),
sa.Column("subject_type", sa.String(length=32), nullable=False),
sa.Column("subject_id", sa.String(length=128), nullable=False),
sa.Column("scope", sa.String(length=32), nullable=False),
sa.Column("active", sa.Boolean(), server_default=sa.true(), nullable=False),
sa.Column("reason", sa.String(length=256), nullable=True),
sa.Column("incident_id", sa.Integer(), nullable=True),
sa.Column("created_by", sa.Integer(), nullable=True),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.text("(CURRENT_TIMESTAMP)"),
nullable=False,
),
sa.Column("revoked_by", sa.Integer(), nullable=True),
sa.Column("revoked_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("details", sa.JSON(), nullable=True),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint(
"subject_type",
"subject_id",
"scope",
name="uq_subject_restriction_subject_scope",
),
)
op.create_index(
"ix_subject_restriction_lookup",
"subject_restriction",
["subject_type", "subject_id", "scope", "active"],
)
_add_default_role_permissions()
def downgrade() -> None:
_remove_default_role_permissions()
op.drop_index("ix_subject_restriction_lookup", table_name="subject_restriction")
op.drop_table("subject_restriction")
op.drop_index("ix_risk_incident_subject", table_name="risk_incident")
op.drop_index("ix_risk_incident_rule_status", table_name="risk_incident")
op.drop_table("risk_incident")
op.drop_index("ix_behavior_event_user_time", table_name="behavior_event")
op.drop_index("ix_behavior_event_subject_time", table_name="behavior_event")
op.drop_index("ix_behavior_event_type_time", table_name="behavior_event")
op.drop_index("ix_behavior_event_occurred_at", table_name="behavior_event")
op.drop_table("behavior_event")
+4 -2
View File
@@ -17,6 +17,7 @@ from app.admin.routers.ad_audit import router as ad_audit_router
from app.admin.routers.ad_config import router as ad_config_router
from app.admin.routers.ad_revenue import router as ad_revenue_router
from app.admin.routers.admins import router as admins_router
from app.admin.routers.analytics_health import router as analytics_health_router
from app.admin.routers.audit import router as audit_router
from app.admin.routers.auth import router as auth_router
from app.admin.routers.comparison import router as comparison_router
@@ -25,8 +26,6 @@ from app.admin.routers.coupon_data import router as coupon_data_router
from app.admin.routers.cps import router as cps_router
from app.admin.routers.dashboard import router as dashboard_router
from app.admin.routers.device_liveness import router as device_liveness_router
from app.admin.routers.ops_stat_config import router as ops_stat_config_router
from app.admin.routers.analytics_health import router as analytics_health_router
from app.admin.routers.event_logs import router as event_logs_router
from app.admin.routers.feedback import router as feedback_router
from app.admin.routers.feedback_qr import router as feedback_qr_router
@@ -34,7 +33,9 @@ from app.admin.routers.guide_video import router as guide_video_router
from app.admin.routers.huawei_review import router as huawei_review_router
from app.admin.routers.onboarding import router as onboarding_router
from app.admin.routers.ops_marquee_seed import router as ops_marquee_seed_router
from app.admin.routers.ops_stat_config import router as ops_stat_config_router
from app.admin.routers.price_report import router as price_report_router
from app.admin.routers.risk_monitor import router as risk_monitor_router
from app.admin.routers.roles import router as roles_router
from app.admin.routers.users import router as users_router
from app.admin.routers.wallet import router as wallet_router
@@ -101,6 +102,7 @@ admin_app.include_router(onboarding_router)
admin_app.include_router(wallet_router)
admin_app.include_router(withdraw_router)
admin_app.include_router(price_report_router)
admin_app.include_router(risk_monitor_router)
admin_app.include_router(feedback_router)
admin_app.include_router(event_logs_router)
admin_app.include_router(analytics_health_router)
+3 -2
View File
@@ -34,6 +34,7 @@ PERMISSION_CATALOG: list[dict] = [
{"key": "users", "label": "用户管理"},
]},
{"group": "监控审计", "pages": [
{"key": "risk-monitor", "label": "风控监控"},
{"key": "device-liveness", "label": "设备存活"},
{"key": "analytics-health", "label": "埋点成功率"},
{"key": "event-logs", "label": "埋点日志"},
@@ -55,13 +56,13 @@ BUILTIN_ROLES: list[dict] = [
{"name": SUPER_ADMIN_ROLE, "label": "管理员", "pages": []},
{"name": "operator", "label": "运营", "pages": [
"dashboard", "coupon-data", "ad-revenue-report", "comparison-records",
"cps", "device-liveness", "price-reports", "feedbacks", "huawei-review",
"cps", "risk-monitor", "device-liveness", "price-reports", "feedbacks", "huawei-review",
]},
{"name": "finance", "label": "财务", "pages": [
"dashboard", "ad-revenue-report", "cps", "withdraws",
]},
{"name": "tech", "label": "技术", "pages": [
"dashboard", "device-liveness", "analytics-health", "config", "ad-revenue", "huawei-review",
"dashboard", "risk-monitor", "device-liveness", "analytics-health", "config", "ad-revenue", "huawei-review",
"event-logs", "audit-logs",
]},
]
+545
View File
@@ -0,0 +1,545 @@
"""风控监控读模型:聚合通用行为流水、风险事件与比价记录。"""
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from sqlalchemy import desc, func, select
from sqlalchemy.orm import Session
from app.core.config import settings
from app.models.comparison import ComparisonRecord
from app.models.device import DeviceLiveness
from app.models.risk import BehaviorEvent, RiskIncident
from app.models.user import User
from app.repositories import risk as risk_repo
KIND_TO_RULE = {
"sms": risk_repo.RULE_SMS_HOURLY,
"oneclick": risk_repo.RULE_ONECLICK_DAILY,
"compare": risk_repo.RULE_COMPARE_DAILY,
}
RULE_TO_KIND = {rule: kind for kind, rule in KIND_TO_RULE.items()}
def _day_bounds(now: datetime | None = None) -> tuple[str, datetime, datetime]:
current = (now or datetime.now(UTC)).astimezone(risk_repo.CN_TZ)
start = current.replace(hour=0, minute=0, second=0, microsecond=0)
end = start + timedelta(days=1)
return (
start.strftime("%Y-%m-%d"),
start.astimezone(UTC),
end.astimezone(UTC),
)
def _compare_day_bounds(now: datetime | None = None) -> tuple[datetime, datetime]:
current = (now or datetime.now(UTC)).astimezone(risk_repo.CN_TZ).replace(
tzinfo=None
)
start = current.replace(hour=0, minute=0, second=0, microsecond=0)
return start, start + timedelta(days=1)
def sync_today_compare_incidents(db: Session, now: datetime | None = None) -> None:
risk_repo.reconcile_compare_rule(db, at=now)
def reset_open_alerts(
db: Session, *, admin_id: int, reset_at: datetime
) -> dict[str, int]:
"""将三类待处理报警归零,并从 reset_at 开始重新累计。"""
rule_codes = tuple(KIND_TO_RULE.values())
incidents = list(
db.scalars(
select(RiskIncident).where(
RiskIncident.rule_code.in_(rule_codes),
RiskIncident.status == "open",
)
).all()
)
counts = {kind: 0 for kind in KIND_TO_RULE}
for incident in incidents:
incident.status = "resolved"
incident.action_reason = risk_repo.MANUAL_RESET_REASON
incident.handled_by = admin_id
incident.handled_at = reset_at
counts[RULE_TO_KIND[incident.rule_code]] += 1
risk_repo.reset_rule_baselines(
db,
rule_codes=rule_codes,
reset_at=reset_at,
admin_id=admin_id,
)
return counts
def summary(db: Session, now: datetime | None = None) -> dict:
date_key, start, end = _day_bounds(now)
compare_start, compare_end = _compare_day_bounds(now)
current = now or datetime.now(UTC)
risk_repo.reconcile_behavior_rule(
db, rule_code=risk_repo.RULE_SMS_HOURLY, at=current, commit=False
)
risk_repo.reconcile_behavior_rule(
db, rule_code=risk_repo.RULE_ONECLICK_DAILY, at=current, commit=False
)
risk_repo.reconcile_compare_rule(db, at=current, commit=False)
db.commit()
sms_threshold = risk_repo.get_rule_threshold(db, risk_repo.RULE_SMS_HOURLY)
oneclick_threshold = risk_repo.get_rule_threshold(
db, risk_repo.RULE_ONECLICK_DAILY
)
compare_threshold = risk_repo.get_rule_threshold(db, risk_repo.RULE_COMPARE_DAILY)
def _behavior_total(event_type: str, outcomes: tuple[str, ...]) -> int:
return int(
db.scalar(
select(func.count(BehaviorEvent.id)).where(
BehaviorEvent.event_type == event_type,
BehaviorEvent.outcome.in_(outcomes),
BehaviorEvent.occurred_at >= start,
BehaviorEvent.occurred_at < end,
)
)
or 0
)
compare_total = int(
db.scalar(
select(func.count(ComparisonRecord.id)).where(
ComparisonRecord.created_at >= compare_start,
ComparisonRecord.created_at < compare_end,
)
)
or 0
)
def _open_count(rule_code: str) -> int:
return int(
db.scalar(
select(func.count(RiskIncident.id)).where(
RiskIncident.rule_code == rule_code,
RiskIncident.status == "open",
RiskIncident.window_key.like(f"{date_key}%"),
)
)
or 0
)
return {
"date": date_key,
"updated_at": now or datetime.now(UTC),
"cards": [
{
"kind": "sms",
"alert_subject_count": _open_count(risk_repo.RULE_SMS_HOURLY),
"today_total": _behavior_total(risk_repo.EVENT_SMS_SEND, ("success",)),
"threshold": sms_threshold,
"rule_text": f"单设备 1 小时 ≥ {sms_threshold}",
},
{
"kind": "oneclick",
"alert_subject_count": _open_count(risk_repo.RULE_ONECLICK_DAILY),
"today_total": _behavior_total(
risk_repo.EVENT_ONECLICK_LOGIN, ("success", "failed")
),
"threshold": oneclick_threshold,
"rule_text": f"单设备当日 ≥ {oneclick_threshold}",
},
{
"kind": "compare",
"alert_subject_count": _open_count(risk_repo.RULE_COMPARE_DAILY),
"today_total": compare_total,
"threshold": compare_threshold,
"rule_text": f"单账户当日 ≥ {compare_threshold}",
},
],
}
def _latest_device_model(db: Session, subject_id: str) -> str | None:
return db.execute(
select(BehaviorEvent.device_model)
.where(
BehaviorEvent.subject_type == "device",
BehaviorEvent.subject_id == subject_id,
BehaviorEvent.device_model.is_not(None),
BehaviorEvent.device_model != "",
)
.order_by(BehaviorEvent.occurred_at.desc(), BehaviorEvent.id.desc())
.limit(1)
).scalar_one_or_none()
def _first_used_at(db: Session, subject_id: str) -> datetime | None:
behavior_first = db.scalar(
select(func.min(BehaviorEvent.occurred_at)).where(
BehaviorEvent.subject_type == "device",
BehaviorEvent.subject_id == subject_id,
)
)
registered_first = db.scalar(
select(func.min(DeviceLiveness.created_at)).where(
DeviceLiveness.device_id == subject_id
)
)
observed = [value for value in (behavior_first, registered_first) if value is not None]
if not observed:
return None
def _utc(value: datetime) -> datetime:
return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
return min(observed, key=_utc)
def _common_device(db: Session, user_id: int, start: datetime, end: datetime) -> str | None:
return db.execute(
select(ComparisonRecord.device_id)
.where(
ComparisonRecord.user_id == user_id,
ComparisonRecord.device_id.is_not(None),
ComparisonRecord.created_at >= start,
ComparisonRecord.created_at < end,
)
.group_by(ComparisonRecord.device_id)
.order_by(func.count(ComparisonRecord.id).desc(), ComparisonRecord.device_id.asc())
.limit(1)
).scalar_one_or_none()
def list_incidents(
db: Session,
*,
kind: str,
status: str = "open",
limit: int = 20,
cursor: int = 0,
) -> tuple[list[dict], int | None, int]:
rule_code = KIND_TO_RULE[kind]
stmt = select(RiskIncident).where(RiskIncident.rule_code == rule_code)
count_stmt = select(func.count(RiskIncident.id)).where(
RiskIncident.rule_code == rule_code
)
if status != "all":
stmt = stmt.where(RiskIncident.status == status)
count_stmt = count_stmt.where(RiskIncident.status == status)
total = int(db.scalar(count_stmt) or 0)
incidents = list(
db.scalars(
stmt.order_by(desc(RiskIncident.triggered_at), desc(RiskIncident.id))
.offset(cursor)
.limit(limit)
).all()
)
user_ids = [int(i.subject_id) for i in incidents if i.subject_type == "user"]
users = {
u.id: u
for u in db.scalars(select(User).where(User.id.in_(user_ids))).all()
} if user_ids else {}
items: list[dict] = []
for incident in incidents:
item = {
"incident_id": incident.id,
"kind": kind,
"subject_type": incident.subject_type,
"subject_id": incident.subject_id,
"window_start": incident.window_start,
"window_end": incident.window_end,
"first_event_at": incident.first_event_at,
"triggered_at": incident.triggered_at,
"last_event_at": incident.last_event_at,
"event_count": incident.event_count,
"status": incident.status,
}
if incident.subject_type == "device":
restriction = risk_repo.get_active_restriction(
db,
subject_type="device",
subject_id=incident.subject_id,
scope=risk_repo.SCOPE_AUTH_DEVICE,
)
item["device_model"] = _latest_device_model(db, incident.subject_id)
item["first_used_at"] = _first_used_at(db, incident.subject_id)
item["restricted"] = restriction is not None
item["restriction_id"] = restriction.id if restriction else None
else:
uid = int(incident.subject_id)
user = users.get(uid)
restriction = risk_repo.get_active_restriction(
db,
subject_type="user",
subject_id=incident.subject_id,
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
)
item.update(
{
"user_id": uid,
"phone": user.phone if user else None,
"registered_at": user.created_at if user else None,
"common_device_id": _common_device(
db, uid, incident.window_start, incident.window_end
),
"restricted": restriction is not None,
"restriction_id": restriction.id if restriction else None,
}
)
items.append(item)
next_cursor = cursor + len(items) if cursor + len(items) < total else None
return items, next_cursor, total
def _page_events(
db: Session,
incident: RiskIncident,
*,
event_type: str,
outcomes: tuple[str, ...] | None,
limit: int,
cursor: int,
) -> tuple[list[BehaviorEvent], int]:
filters = [
BehaviorEvent.event_type == event_type,
BehaviorEvent.subject_type == incident.subject_type,
BehaviorEvent.subject_id == incident.subject_id,
BehaviorEvent.occurred_at >= incident.window_start,
BehaviorEvent.occurred_at < incident.window_end,
]
if outcomes:
filters.append(BehaviorEvent.outcome.in_(outcomes))
total = int(db.scalar(select(func.count(BehaviorEvent.id)).where(*filters)) or 0)
rows = list(
db.scalars(
select(BehaviorEvent)
.where(*filters)
.order_by(BehaviorEvent.occurred_at.desc(), BehaviorEvent.id.desc())
.offset(cursor)
.limit(limit)
).all()
)
return rows, total
def _user_map(db: Session, user_ids: set[int]) -> dict[int, User]:
if not user_ids:
return {}
return {u.id: u for u in db.scalars(select(User).where(User.id.in_(user_ids))).all()}
def _distinct_event_users(
db: Session,
incident: RiskIncident,
*,
event_type: str,
outcomes: tuple[str, ...],
) -> int:
return int(
db.scalar(
select(func.count(func.distinct(BehaviorEvent.user_id))).where(
BehaviorEvent.event_type == event_type,
BehaviorEvent.subject_type == incident.subject_type,
BehaviorEvent.subject_id == incident.subject_id,
BehaviorEvent.outcome.in_(outcomes),
BehaviorEvent.user_id.is_not(None),
BehaviorEvent.occurred_at >= incident.window_start,
BehaviorEvent.occurred_at < incident.window_end,
)
)
or 0
)
def _sms_distinct_accounts(db: Session, incident: RiskIncident) -> int:
"""统计整个告警窗口内,确实完成本设备短信验证的不同账号数。"""
send_rows = db.execute(
select(BehaviorEvent.phone, BehaviorEvent.occurred_at).where(
BehaviorEvent.event_type == risk_repo.EVENT_SMS_SEND,
BehaviorEvent.subject_type == incident.subject_type,
BehaviorEvent.subject_id == incident.subject_id,
BehaviorEvent.outcome == "success",
BehaviorEvent.phone.is_not(None),
BehaviorEvent.occurred_at >= incident.window_start,
BehaviorEvent.occurred_at < incident.window_end,
)
).all()
phones = {phone for phone, _ in send_rows}
if not phones:
return 0
login_rows = db.execute(
select(
BehaviorEvent.user_id,
BehaviorEvent.phone,
BehaviorEvent.occurred_at,
)
.where(
BehaviorEvent.event_type == risk_repo.EVENT_SMS_LOGIN,
BehaviorEvent.subject_id == incident.subject_id,
BehaviorEvent.phone.in_(phones),
BehaviorEvent.user_id.is_not(None),
BehaviorEvent.outcome == "success",
BehaviorEvent.occurred_at >= incident.window_start,
BehaviorEvent.occurred_at
< incident.window_end + timedelta(seconds=settings.SMS_CODE_TTL_SEC),
)
.order_by(BehaviorEvent.occurred_at.asc(), BehaviorEvent.id.asc())
).all()
matched_users = {
user_id
for user_id, phone, login_at in login_rows
if any(
sent_phone == phone
and sent_at <= login_at
and login_at
<= sent_at + timedelta(seconds=settings.SMS_CODE_TTL_SEC)
for sent_phone, sent_at in send_rows
)
}
return len(matched_users)
def incident_details(
db: Session,
*,
incident: RiskIncident,
kind: str,
limit: int,
cursor: int,
) -> dict:
if kind == "sms":
rows, total = _page_events(
db,
incident,
event_type=risk_repo.EVENT_SMS_SEND,
outcomes=("success",),
limit=limit,
cursor=cursor,
)
phone_values = {r.phone for r in rows if r.phone}
login_rows = list(
db.scalars(
select(BehaviorEvent)
.where(
BehaviorEvent.event_type == risk_repo.EVENT_SMS_LOGIN,
BehaviorEvent.subject_id == incident.subject_id,
BehaviorEvent.phone.in_(phone_values),
BehaviorEvent.outcome == "success",
BehaviorEvent.occurred_at >= incident.window_start,
BehaviorEvent.occurred_at
< incident.window_end + timedelta(seconds=settings.SMS_CODE_TTL_SEC),
)
.order_by(BehaviorEvent.occurred_at.asc(), BehaviorEvent.id.asc())
).all()
) if phone_values else []
users = _user_map(db, {r.user_id for r in login_rows if r.user_id is not None})
items = []
for row in rows:
match = next(
(
login
for login in login_rows
if login.phone == row.phone
and login.occurred_at >= row.occurred_at
and login.occurred_at
<= row.occurred_at + timedelta(seconds=settings.SMS_CODE_TTL_SEC)
),
None,
)
user = users.get(match.user_id) if match and match.user_id else None
items.append(
{
"id": row.id,
"occurred_at": row.occurred_at,
"outcome": row.outcome,
"phone": row.phone,
"user_id": user.id if user else None,
"username": user.username if user else None,
"account_registered_at": user.created_at if user else None,
"verified": match is not None,
}
)
distinct_accounts = _sms_distinct_accounts(db, incident)
elif kind == "oneclick":
rows, total = _page_events(
db,
incident,
event_type=risk_repo.EVENT_ONECLICK_LOGIN,
outcomes=("success", "failed"),
limit=limit,
cursor=cursor,
)
users = _user_map(db, {r.user_id for r in rows if r.user_id is not None})
items = []
for row in rows:
user = users.get(row.user_id) if row.user_id else None
items.append(
{
"id": row.id,
"occurred_at": row.occurred_at,
"outcome": row.outcome,
"phone": row.phone,
"user_id": row.user_id,
"username": user.username if user else None,
"account_registered_at": user.created_at if user else None,
"reason": row.reason,
}
)
distinct_accounts = _distinct_event_users(
db,
incident,
event_type=risk_repo.EVENT_ONECLICK_LOGIN,
outcomes=("success", "failed"),
)
else:
user_id = int(incident.subject_id)
filters = (
ComparisonRecord.user_id == user_id,
ComparisonRecord.created_at >= incident.window_start,
ComparisonRecord.created_at < incident.window_end,
)
total = int(db.scalar(select(func.count(ComparisonRecord.id)).where(*filters)) or 0)
rows = list(
db.scalars(
select(ComparisonRecord)
.where(*filters)
.order_by(ComparisonRecord.created_at.desc(), ComparisonRecord.id.desc())
.offset(cursor)
.limit(limit)
).all()
)
items = []
for row in rows:
prices = {
str(result.get("platform_name") or result.get("platform_id") or "未知平台"): (
float(result["price"]) if result.get("price") is not None else None
)
for result in (row.comparison_results or [])
}
items.append(
{
"id": row.id,
"occurred_at": row.created_at,
"outcome": row.status,
"store_name": row.store_name,
"product_names": row.product_names,
"prices": prices,
"saved_amount_yuan": (
row.saved_amount_cents / 100
if row.saved_amount_cents is not None
else None
),
"status": row.status,
}
)
distinct_accounts = None
next_cursor = cursor + len(items) if cursor + len(items) < total else None
return {
"kind": kind,
"incident_id": incident.id,
"event_count": incident.event_count,
"distinct_accounts": distinct_accounts,
"items": items,
"next_cursor": next_cursor,
"total": total,
}
+6
View File
@@ -30,6 +30,12 @@ def _validate(key: str, value: Any) -> None:
if t == "int":
if not isinstance(value, int) or isinstance(value, bool) or value < 0:
raise ValueError("需为非负整数")
minimum = CONFIG_DEFS[key].get("min")
maximum = CONFIG_DEFS[key].get("max")
if minimum is not None and value < minimum:
raise ValueError(f"不能小于 {minimum}")
if maximum is not None and value > maximum:
raise ValueError(f"不能大于 {maximum}")
elif t == "int_list":
if not isinstance(value, list) or not value:
raise ValueError("需为非空整数列表")
+331
View File
@@ -0,0 +1,331 @@
"""风控监控:统计、风险事件列表/详情及忽略、封禁、解封。"""
from __future__ import annotations
from typing import Annotated, Literal
from fastapi import APIRouter, Depends, HTTPException, Query, Request, status
from app.admin.audit import write_audit
from app.admin.deps import (
AdminDb,
get_client_ip,
require_page,
require_role,
)
from app.admin.repositories import risk_monitor as repo
from app.admin.schemas.risk_monitor import (
RestrictionRevokeResponse,
RiskActionRequest,
RiskActionResponse,
RiskDetailPage,
RiskIncidentPage,
RiskMonitorSummary,
RiskResetResponse,
RiskRuleConfig,
)
from app.core.config_schema import (
RISK_COMPARE_DAILY_THRESHOLD_KEY,
RISK_ONECLICK_DAILY_THRESHOLD_KEY,
RISK_SMS_HOURLY_THRESHOLD_KEY,
)
from app.models.risk import RiskIncident, SubjectRestriction
from app.repositories import app_config
from app.repositories import risk as risk_repo
router = APIRouter(
prefix="/admin/api/risk-monitor",
tags=["admin-risk-monitor"],
dependencies=[Depends(require_page("risk-monitor"))],
)
RiskKind = Literal["sms", "oneclick", "compare"]
def _rule_config(db) -> RiskRuleConfig:
return RiskRuleConfig(
sms_hourly_threshold=risk_repo.get_rule_threshold(
db, risk_repo.RULE_SMS_HOURLY
),
oneclick_daily_threshold=risk_repo.get_rule_threshold(
db, risk_repo.RULE_ONECLICK_DAILY
),
compare_daily_threshold=risk_repo.get_rule_threshold(
db, risk_repo.RULE_COMPARE_DAILY
),
)
def _incident_or_404(db, incident_id: int, kind: RiskKind | None = None) -> RiskIncident:
incident = db.get(RiskIncident, incident_id)
if incident is None:
raise HTTPException(status_code=404, detail="风险事件不存在")
if kind and incident.rule_code != repo.KIND_TO_RULE[kind]:
raise HTTPException(status_code=404, detail="风险事件类型不匹配")
return incident
@router.get("/summary", response_model=RiskMonitorSummary, summary="风控监控顶部卡片")
def get_summary(db: AdminDb) -> RiskMonitorSummary:
return RiskMonitorSummary(**repo.summary(db))
@router.get("/rules", response_model=RiskRuleConfig, summary="读取风控规则阈值")
def get_rules(db: AdminDb) -> RiskRuleConfig:
return _rule_config(db)
@router.patch("/rules", response_model=RiskRuleConfig, summary="修改风控规则并即时重算")
def update_rules(
body: RiskRuleConfig,
request: Request,
admin: Annotated[object, Depends(require_role("operator", "tech"))],
db: AdminDb,
) -> RiskRuleConfig:
before = _rule_config(db).model_dump()
after = body.model_dump()
values = (
(RISK_SMS_HOURLY_THRESHOLD_KEY, body.sms_hourly_threshold),
(RISK_ONECLICK_DAILY_THRESHOLD_KEY, body.oneclick_daily_threshold),
(RISK_COMPARE_DAILY_THRESHOLD_KEY, body.compare_daily_threshold),
)
for key, value in values:
app_config.set_value(
db, key, value, admin_id=admin.id, commit=False
)
now = risk_repo.utcnow()
risk_repo.reconcile_behavior_rule(
db, rule_code=risk_repo.RULE_SMS_HOURLY, at=now, commit=False
)
risk_repo.reconcile_behavior_rule(
db, rule_code=risk_repo.RULE_ONECLICK_DAILY, at=now, commit=False
)
risk_repo.reconcile_compare_rule(db, at=now, commit=False)
write_audit(
db,
admin,
action="risk.rules.update",
target_type="risk_rule",
target_id="all",
detail={"before": before, "after": after},
ip=get_client_ip(request),
commit=False,
)
db.commit()
return _rule_config(db)
@router.post(
"/reset",
response_model=RiskResetResponse,
summary="重置三类待处理报警并从零重新累计",
)
def reset_alerts(
request: Request,
admin: Annotated[object, Depends(require_role("operator", "tech"))],
db: AdminDb,
) -> RiskResetResponse:
reset_at = risk_repo.utcnow()
counts = repo.reset_open_alerts(db, admin_id=admin.id, reset_at=reset_at)
total = sum(counts.values())
write_audit(
db,
admin,
action="risk.alerts.reset",
target_type="risk_monitor",
target_id="all",
detail={
"reset_at": reset_at.isoformat(),
"reset_incident_count": total,
"reset_counts": counts,
"restrictions_changed": False,
},
ip=get_client_ip(request),
commit=False,
)
db.commit()
return RiskResetResponse(
reset_at=reset_at,
reset_incident_count=total,
reset_counts=counts,
)
@router.get(
"/incidents/{kind}",
response_model=RiskIncidentPage,
summary="按短信/一键登录/比价读取风险事件",
)
def get_incidents(
kind: RiskKind,
db: AdminDb,
incident_status: Annotated[
Literal["open", "ignored", "blocked", "resolved", "all"], Query(alias="status")
] = "open",
limit: Annotated[int, Query(ge=1, le=100)] = 20,
cursor: Annotated[int, Query(ge=0)] = 0,
) -> RiskIncidentPage:
if kind == "compare":
repo.sync_today_compare_incidents(db)
items, next_cursor, total = repo.list_incidents(
db, kind=kind, status=incident_status, limit=limit, cursor=cursor
)
return RiskIncidentPage(items=items, next_cursor=next_cursor, total=total)
@router.get(
"/incidents/{kind}/{incident_id}/details",
response_model=RiskDetailPage,
summary="风险事件当期事实明细",
)
def get_incident_details(
kind: RiskKind,
incident_id: int,
db: AdminDb,
limit: Annotated[int, Query(ge=1, le=100)] = 20,
cursor: Annotated[int, Query(ge=0)] = 0,
) -> RiskDetailPage:
incident = _incident_or_404(db, incident_id, kind)
return RiskDetailPage(
**repo.incident_details(
db, incident=incident, kind=kind, limit=limit, cursor=cursor
)
)
@router.post(
"/incidents/{incident_id}/ignore",
response_model=RiskActionResponse,
summary="忽略风险事件",
)
def ignore_incident(
incident_id: int,
body: RiskActionRequest,
request: Request,
admin: Annotated[object, Depends(require_role("operator", "tech"))],
db: AdminDb,
) -> RiskActionResponse:
incident = _incident_or_404(db, incident_id)
if incident.status != "open":
raise HTTPException(status_code=409, detail="该风险事件已处理")
incident.status = "ignored"
incident.action_reason = body.reason.strip() or "管理员确认忽略"
incident.handled_by = admin.id
incident.handled_at = risk_repo.utcnow()
write_audit(
db,
admin,
action="risk.incident.ignore",
target_type="risk_incident",
target_id=incident.id,
detail={
"rule_code": incident.rule_code,
"subject_type": incident.subject_type,
"subject_id": incident.subject_id,
"reason": incident.action_reason,
},
ip=get_client_ip(request),
commit=False,
)
db.commit()
return RiskActionResponse(incident_id=incident.id, status=incident.status)
@router.post(
"/incidents/{incident_id}/block",
response_model=RiskActionResponse,
summary="封禁风险主体",
)
def block_incident(
incident_id: int,
body: RiskActionRequest,
request: Request,
admin: Annotated[object, Depends(require_role("operator", "tech"))],
db: AdminDb,
) -> RiskActionResponse:
incident = _incident_or_404(db, incident_id)
if incident.status not in ("open", "blocked"):
raise HTTPException(status_code=409, detail="该风险事件已处理")
if incident.subject_type == "device":
if incident.subject_id.startswith("legacy-ip:"):
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST,
detail="旧客户端未上报设备ID,不能执行设备封禁",
)
scope = risk_repo.SCOPE_AUTH_DEVICE
elif incident.subject_type == "user":
scope = risk_repo.SCOPE_ECONOMIC_ACCOUNT
else: # pragma: no cover - 当前三条规则只会生成 device/user
raise HTTPException(status_code=400, detail="该主体类型暂不支持封禁")
reason = body.reason.strip() or "风控规则命中"
restriction = risk_repo.activate_restriction(
db,
subject_type=incident.subject_type,
subject_id=incident.subject_id,
scope=scope,
reason=reason,
incident_id=incident.id,
admin_id=admin.id,
)
incident.status = "blocked"
incident.action_reason = reason
incident.handled_by = admin.id
incident.handled_at = risk_repo.utcnow()
write_audit(
db,
admin,
action="risk.subject.block",
target_type=incident.subject_type,
target_id=incident.subject_id,
detail={
"incident_id": incident.id,
"rule_code": incident.rule_code,
"scope": scope,
"reason": reason,
},
ip=get_client_ip(request),
commit=False,
)
db.commit()
return RiskActionResponse(
incident_id=incident.id,
status=incident.status,
restriction_id=restriction.id,
)
@router.post(
"/restrictions/{restriction_id}/revoke",
response_model=RestrictionRevokeResponse,
summary="解除主体限制",
)
def revoke_restriction(
restriction_id: int,
request: Request,
admin: Annotated[object, Depends(require_role("operator", "tech"))],
db: AdminDb,
) -> RestrictionRevokeResponse:
restriction = db.get(SubjectRestriction, restriction_id)
if restriction is None:
raise HTTPException(status_code=404, detail="限制记录不存在")
if restriction.active:
risk_repo.revoke_restriction(db, restriction=restriction, admin_id=admin.id)
if restriction.incident_id:
incident = db.get(RiskIncident, restriction.incident_id)
if incident and incident.status == "blocked":
incident.status = "resolved"
incident.handled_by = admin.id
incident.handled_at = risk_repo.utcnow()
write_audit(
db,
admin,
action="risk.subject.unblock",
target_type=restriction.subject_type,
target_id=restriction.subject_id,
detail={"restriction_id": restriction.id, "scope": restriction.scope},
ip=get_client_ip(request),
commit=False,
)
db.commit()
return RestrictionRevokeResponse(restriction_id=restriction.id)
+104
View File
@@ -0,0 +1,104 @@
"""风控监控后台接口 schema。"""
from __future__ import annotations
from datetime import datetime
from pydantic import BaseModel, Field
class RiskSummaryCard(BaseModel):
kind: str
alert_subject_count: int
today_total: int
threshold: int
rule_text: str
class RiskMonitorSummary(BaseModel):
date: str
updated_at: datetime
cards: list[RiskSummaryCard]
class RiskRuleConfig(BaseModel):
sms_hourly_threshold: int = Field(ge=1, le=5)
oneclick_daily_threshold: int = Field(ge=1, le=100_000)
compare_daily_threshold: int = Field(ge=1, le=100)
class RiskIncidentItem(BaseModel):
incident_id: int
kind: str
subject_type: str
subject_id: str
device_model: str | None = None
first_used_at: datetime | None = None
phone: str | None = None
user_id: int | None = None
registered_at: datetime | None = None
common_device_id: str | None = None
window_start: datetime
window_end: datetime
first_event_at: datetime
triggered_at: datetime
last_event_at: datetime
event_count: int
status: str
restricted: bool = False
restriction_id: int | None = None
class RiskIncidentPage(BaseModel):
items: list[RiskIncidentItem]
next_cursor: int | None = None
total: int
class RiskDetailItem(BaseModel):
id: int
occurred_at: datetime
outcome: str
phone: str | None = None
user_id: int | None = None
username: str | None = None
account_registered_at: datetime | None = None
verified: bool | None = None
reason: str | None = None
store_name: str | None = None
product_names: str | None = None
prices: dict[str, float | None] | None = None
saved_amount_yuan: float | None = None
status: str | None = None
class RiskDetailPage(BaseModel):
kind: str
incident_id: int
event_count: int
distinct_accounts: int | None = None
items: list[RiskDetailItem]
next_cursor: int | None = None
total: int
class RiskActionRequest(BaseModel):
reason: str = Field("", max_length=256)
class RiskActionResponse(BaseModel):
ok: bool = True
incident_id: int
status: str
restriction_id: int | None = None
class RiskResetResponse(BaseModel):
ok: bool = True
reset_at: datetime
reset_incident_count: int
reset_counts: dict[str, int]
class RestrictionRevokeResponse(BaseModel):
ok: bool = True
restriction_id: int
+147 -4
View File
@@ -37,6 +37,7 @@ from app.integrations.jiguang import JiguangError, mask_phone, verify_and_get_ph
from app.integrations.sms import SmsError, send_code, verify_code
from app.repositories import onboarding as onboarding_repo
from app.repositories import phone_rebind as rebind_repo
from app.repositories import risk as risk_repo
from app.repositories import user as user_repo
from app.schemas.auth import (
JverifyLoginRequest,
@@ -70,6 +71,18 @@ SMS_SEND_MAX_PER_HOUR_PER_DEVICE = 5 # 每小时上限
SMS_SEND_MAX_PER_DAY_PER_DEVICE = 20 # 每天上限(再叠一层日封顶,挡低频长时间轰炸)
def _client_ip(request: Request) -> str:
xff = request.headers.get("x-forwarded-for")
if xff:
return xff.split(",")[0].strip()
return request.client.host if request.client else ""
def _device_subject(device_id: str, request: Request) -> str:
# 兼容旧客户端空 device_id;新客户端均上传稳定硬件标识。
return device_id or f"legacy-ip:{_client_ip(request)}"
def _login_response(
user, *, onboarding_completed: bool, force_onboarding: bool = False
) -> TokenWithUser:
@@ -87,7 +100,15 @@ def _login_response(
# ===================== 极光一键登录 =====================
@router.post("/jverify-login", response_model=TokenWithUser, summary="极光一键登录")
def jverify_login(req: JverifyLoginRequest, db: DbSession) -> TokenWithUser:
def jverify_login(req: JverifyLoginRequest, request: Request, db: DbSession) -> TokenWithUser:
subject_id = _device_subject(req.device_id, request)
if req.device_id and risk_repo.is_restricted(
db,
subject_type="device",
subject_id=req.device_id,
scope=risk_repo.SCOPE_AUTH_DEVICE,
):
raise HTTPException(status_code=403, detail="当前设备环境异常,暂无法登录")
logger.info(
"jverify_login operator=%s token_len=%d",
req.operator or "-",
@@ -97,10 +118,37 @@ def jverify_login(req: JverifyLoginRequest, db: DbSession) -> TokenWithUser:
try:
phone = verify_and_get_phone(req.login_token)
except JiguangError as e:
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_ONECLICK_LOGIN,
subject_type="device",
subject_id=subject_id,
device_id=req.device_id or None,
device_model=req.device_model or None,
client_ip=_client_ip(request),
outcome="failed",
reason=str(e),
details={"operator": req.operator or None},
evaluate_rule=risk_repo.RULE_ONECLICK_DAILY,
)
logger.error("[JG] verify+decrypt failed: %s", e, exc_info=True)
raise HTTPException(status_code=502, detail=f"jiguang verify failed: {e}") from e
user = user_repo.upsert_user_for_login(db, phone=phone, register_channel="jverify")
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_ONECLICK_LOGIN,
subject_type="device",
subject_id=subject_id,
user_id=user.id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=phone,
client_ip=_client_ip(request),
outcome="success",
details={"operator": req.operator or None},
evaluate_rule=risk_repo.RULE_ONECLICK_DAILY,
)
if user.status != "active":
raise HTTPException(status_code=403, detail="account disabled")
@@ -116,12 +164,34 @@ def jverify_login(req: JverifyLoginRequest, db: DbSession) -> TokenWithUser:
response_model=SmsSendResponse,
summary="发送短信验证码",
)
def sms_send(req: SmsSendRequest, request: Request) -> SmsSendResponse:
def sms_send(req: SmsSendRequest, request: Request, db: DbSession) -> SmsSendResponse:
subject_id = _device_subject(req.device_id, request)
if req.device_id and risk_repo.is_restricted(
db,
subject_type="device",
subject_id=req.device_id,
scope=risk_repo.SCOPE_AUTH_DEVICE,
):
raise HTTPException(status_code=403, detail="当前设备环境异常,暂无法发送验证码")
# 测试账号:不真发短信(号码非真实手机号,真发会失败/浪费),直接放行让客户端进入填码界面。
# 真正的"免验证码"在 /sms/login 跳过校验;每日上限也在 login 处算(send 不耗额度)。
# (也不受下面设备发码限流约束:QA 联调要反复发码。)
if test_account.is_test_account(req.phone):
logger.info("test_account sms_send short-circuit (不真发)")
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=subject_id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=req.phone,
client_ip=_client_ip(request),
# 测试号不真发短信,也豁免正式频控;保留流水供排查,但不能污染
# “今日短信下发”与短信风控报警。
outcome="test",
details={"mock": True, "test_account": True},
)
return SmsSendResponse(sent=True, mock=True, cooldown_sec=0)
# 发码防刷:同一设备(device_id) + 同一 IP,每小时 / 每天两道闸,**均只按成功发码计数**。
@@ -139,6 +209,18 @@ def sms_send(req: SmsSendRequest, request: Request) -> SmsSendResponse:
try:
cooldown = send_code(req.phone)
except SmsError as e:
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=subject_id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=req.phone,
client_ip=_client_ip(request),
outcome="failed",
reason=str(e),
)
raise HTTPException(status_code=e.status_code, detail=str(e)) from e
# 发码成功 → 两道闸各 +1(被单号冷却挡下的重发走不到这里,故不占额度)
@@ -146,6 +228,20 @@ def sms_send(req: SmsSendRequest, request: Request) -> SmsSendResponse:
from app.core.config import settings # 局部 import 避免循环
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=subject_id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=req.phone,
client_ip=_client_ip(request),
outcome="success",
details={"mock": settings.SMS_MOCK},
evaluate_rule=risk_repo.RULE_SMS_HOURLY,
)
return SmsSendResponse(sent=True, mock=settings.SMS_MOCK, cooldown_sec=cooldown)
@@ -155,6 +251,7 @@ def sms_send(req: SmsSendRequest, request: Request) -> SmsSendResponse:
summary="手机号+验证码登录",
)
def sms_login(req: SmsLoginRequest, request: Request, db: DbSession) -> TokenWithUser:
subject_id = _device_subject(req.device_id, request)
# 测试账号:免验证码登录 + 每日上限 + 每次都走新手引导(详见 app/core/test_account.py)。
# 放在最前面:命中即不校验验证码;先扣当日额度,超限直接拒,挡住有人猜到号后脚本刷。
# 测试账号走自己的每日额度、不受下面 (设备+IP) 每小时限流约束(QA 需在一小时内反复登录联调)。
@@ -162,6 +259,19 @@ def sms_login(req: SmsLoginRequest, request: Request, db: DbSession) -> TokenWit
if not test_account.try_consume_quota():
raise HTTPException(status_code=429, detail="测试账号今日使用次数已达上限,请明天再试")
user = user_repo.upsert_user_for_login(db, phone=req.phone, register_channel="sms")
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_LOGIN,
subject_type="device",
subject_id=subject_id,
user_id=user.id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=req.phone,
client_ip=_client_ip(request),
outcome="success",
details={"test_account": True},
)
if user.status != "active":
raise HTTPException(status_code=403, detail="account disabled")
logger.info("sms_login test_account ok user_id=%d phone=%s", user.id, mask_phone(req.phone))
@@ -182,10 +292,39 @@ def sms_login(req: SmsLoginRequest, request: Request, db: DbSession) -> TokenWit
detail="登录尝试过于频繁,请稍后再试",
)
if not verify_code(req.phone, req.code):
try:
ok = verify_code(req.phone, req.code)
except SmsError as e: # provider 校验降级(如阿里云接口异常)→ 原样透出其状态码(503),别误报「验证码错误」
raise HTTPException(status_code=e.status_code, detail=str(e)) from e
if not ok:
# 校验码错误才记风控失败事件(provider 降级 503 已在上面提前 raise,不算「验证失败」)
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_LOGIN,
subject_type="device",
subject_id=subject_id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=req.phone,
client_ip=_client_ip(request),
outcome="failed",
reason="invalid sms code",
)
raise HTTPException(status_code=400, detail="invalid sms code")
user = user_repo.upsert_user_for_login(db, phone=req.phone, register_channel="sms")
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_LOGIN,
subject_type="device",
subject_id=subject_id,
user_id=user.id,
device_id=req.device_id or None,
device_model=req.device_model or None,
phone=req.phone,
client_ip=_client_ip(request),
outcome="success",
)
if user.status != "active":
raise HTTPException(status_code=403, detail="account disabled")
@@ -322,7 +461,11 @@ def wechat_bind_phone_sms(
detail="登录尝试过于频繁,请稍后再试",
)
if not verify_code(req.phone, req.code):
try:
ok = verify_code(req.phone, req.code)
except SmsError as e: # provider 校验降级(如阿里云接口异常)→ 原样透出其状态码(503),别误报「验证码错误」
raise HTTPException(status_code=e.status_code, detail=str(e)) from e
if not ok:
raise HTTPException(status_code=400, detail="invalid sms code")
return _finish_wechat_bind(
+41 -8
View File
@@ -28,19 +28,30 @@ import httpx
from fastapi import APIRouter, HTTPException, Request, status
from fastapi.concurrency import run_in_threadpool
from app.api.deps import OptionalUser
from app.api.deps import DbSession, OptionalUser
from app.core.config import settings
from app.core.logging import trace_id_ctx
from app.core.pricebot_client import get_pricebot_client
from app.core.pricebot_router import pick_pricebot
from app.db.session import SessionLocal
from app.repositories import comparison as crud_compare
from app.repositories import risk as risk_repo
logger = logging.getLogger("shagua.compare")
router = APIRouter(prefix="/api/v1", tags=["compare"])
def _ensure_compare_allowed(user, db) -> None:
if user and risk_repo.is_restricted(
db,
subject_type="user",
subject_id=str(user.id),
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
):
raise HTTPException(status_code=403, detail="账号存在异常,该功能暂不可用")
# ============================================================
# harvest 落库(阻塞 SQLAlchemy → run_in_threadpool,独立 SessionLocal,不阻塞事件循环;
# 任何写库异常都吞掉、绝不连累比价透传返回 —— 同 coupon.py 现有 best-effort 写法)。
@@ -52,10 +63,14 @@ def _harvest_running_blocking(
device_id: str | None, device_info: dict | None, trace_url: str | None,
) -> None:
with SessionLocal() as db:
crud_compare.harvest_running(
rec = crud_compare.harvest_running(
db, trace_id=trace_id, user_id=user_id, business_type=business_type,
device_id=device_id, device_info=device_info, trace_url=trace_url,
)
if rec.user_id is not None:
risk_repo.sync_compare_incident(
db, user_id=rec.user_id, at=rec.created_at
)
logger.info(
"harvest running row (user=%s)", user_id,
extra={"phase": "harvest_running", "status": "running", "user_id": user_id},
@@ -72,6 +87,10 @@ def _harvest_done_blocking(
business_type=business_type, device_id=device_id,
device_info=device_info, trace_url=trace_url,
)
if rec.user_id is not None:
risk_repo.sync_compare_incident(
db, user_id=rec.user_id, at=rec.created_at
)
logger.info(
"harvest done → %s saved=%s newly=%s", rec.status,
rec.saved_amount_cents, newly_success,
@@ -199,25 +218,33 @@ async def _forward(
@router.post("/intent/recognize", summary="外卖比价 Phase 1 意图识别 (透传 + 建 running 行)")
async def intent_recognize(request: Request, user: OptionalUser) -> dict[str, Any]:
async def intent_recognize(
request: Request, user: OptionalUser, db: DbSession
) -> dict[str, Any]:
_ensure_compare_allowed(user, db)
resp, _, _ = await _forward(request, "/api/intent/recognize", user)
return resp
@router.post("/intent/step", summary="外卖比价 Phase 1 多帧意图识别 (透传, 仅淘宝源)")
async def intent_step(request: Request, user: OptionalUser) -> dict[str, Any]:
async def intent_step(request: Request, user: OptionalUser, db: DbSession) -> dict[str, Any]:
_ensure_compare_allowed(user, db)
resp, _, _ = await _forward(request, "/api/intent/step", user)
return resp
@router.post("/intent/precoupon/step", summary="外卖比价 Phase 0 识别前先用券 (透传, 仅美团源)")
async def intent_precoupon_step(request: Request, user: OptionalUser) -> dict[str, Any]:
async def intent_precoupon_step(
request: Request, user: OptionalUser, db: DbSession
) -> dict[str, Any]:
_ensure_compare_allowed(user, db)
resp, _, _ = await _forward(request, "/api/intent/precoupon/step", user)
return resp
@router.post("/price/step", summary="外卖比价 Phase 2 步进 (透传 + done 落库)")
async def price_step(request: Request, user: OptionalUser) -> dict[str, Any]:
async def price_step(request: Request, user: OptionalUser, db: DbSession) -> dict[str, Any]:
_ensure_compare_allowed(user, db)
resp, trace_id, meta = await _forward(request, "/api/price/step", user)
# 最终 done 帧(command=done 且 continue=false)→ harvest 更新成终态。
# (单平台中途 done 被 pricebot 改写成 wait+continue=true,不会命中这里,同 coupon 语义。)
@@ -237,7 +264,10 @@ async def price_step(request: Request, user: OptionalUser) -> dict[str, Any]:
@router.post("/trace/epilogue", summary="比价结果页尾声帧 (透传到 pricebot, 用户视角截图入 trace)")
async def trace_epilogue(request: Request, user: OptionalUser) -> dict[str, Any]:
async def trace_epilogue(
request: Request, user: OptionalUser, db: DbSession
) -> dict[str, Any]:
_ensure_compare_allowed(user, db)
# App 收到 done、渲染完结果页后,把自己页面的截图(base64)传给 pricebot 存进 trace
# 目录并触发重传 —— trace 里补上"用户实际看到的汇总页"(步骤帧只有目标 App 画面)。
# body ~几百 KB(截图 base64), 纯透传壳(harvest_first_frame=False:不建行/不落库,
@@ -250,7 +280,10 @@ async def trace_epilogue(request: Request, user: OptionalUser) -> dict[str, Any]
@router.post("/trace/finalize", summary="比价 trace 收尾上云 (透传 + 夭折落库)")
async def trace_finalize(request: Request, user: OptionalUser) -> dict[str, Any]:
async def trace_finalize(
request: Request, user: OptionalUser, db: DbSession
) -> dict[str, Any]:
_ensure_compare_allowed(user, db)
# 用户终止 / Phase1 未识别没到 done 帧: pricebot 打包半截上云返回 {trace_url};
# app-server 顺手把该 trace 的 running 行更新成夭折终态(**不降级 success**)。
# 客户端 finalize body 带 status(cancelled/failed)+ reason;老客户端只带 trace_id →
+18 -1
View File
@@ -19,6 +19,7 @@ from app.api.deps import CurrentUser, DbSession
from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord
from app.repositories import comparison as crud_compare
from app.repositories import risk as risk_repo
from app.schemas.compare_record import (
CompareStartReserveIn,
CompareStartReserveOut,
@@ -47,8 +48,15 @@ def reserve_compare_start(
user: CurrentUser,
db: DbSession,
) -> CompareStartReserveOut:
if risk_repo.is_restricted(
db,
subject_type="user",
subject_id=str(user.id),
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
):
raise HTTPException(status_code=403, detail="账号存在异常,该功能暂不可用")
try:
_, used = crud_compare.reserve_daily_start(
rec, used = crud_compare.reserve_daily_start(
db,
user_id=user.id,
trace_id=payload.trace_id,
@@ -65,6 +73,8 @@ def reserve_compare_start(
status_code=status.HTTP_409_CONFLICT,
detail="比价任务标识冲突,请重新发起",
) from None
# 风控阈值由后台动态配置,不能再只在固定的 100 次业务上限处同步。
risk_repo.sync_compare_incident(db, user_id=user.id, at=rec.created_at)
return CompareStartReserveOut(
limit=crud_compare.DAILY_COMPARE_START_LIMIT,
used=used,
@@ -83,6 +93,13 @@ def report_record(
db: DbSession,
background_tasks: BackgroundTasks,
) -> ComparisonRecordCreatedOut:
if risk_repo.is_restricted(
db,
subject_type="user",
subject_id=str(user.id),
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
):
raise HTTPException(status_code=403, detail="账号存在异常,该功能暂不可用")
rec = crud_compare.upsert_record(db, user_id=user.id, payload=payload)
# LLM 调用明细异步回填:同机拉 pricebot llm_calls 最长 5s 且是 best-effort,放后台
# 任务做,不阻塞上报响应(顺带给 pricebot 落盘留足余量)。upsert 已 commit,后台用
+8
View File
@@ -11,6 +11,7 @@ import logging
from fastapi import APIRouter, HTTPException, status
from app.api.deps import CurrentUser, DbSession
from app.repositories import risk as risk_repo
from app.repositories import task as crud_task
from app.schemas.welfare import TaskClaimResultOut, TaskListOut, TaskOut
@@ -36,6 +37,13 @@ def list_tasks(user: CurrentUser, db: DbSession) -> TaskListOut:
summary="领取任务奖励",
)
def claim_task(task_key: str, user: CurrentUser, db: DbSession) -> TaskClaimResultOut:
if risk_repo.is_restricted(
db,
subject_type="user",
subject_id=str(user.id),
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
):
raise HTTPException(status_code=403, detail="账号存在异常,该功能暂不可用")
try:
coin, balance = crud_task.claim_task(db, user.id, task_key)
except crud_task.UnknownTaskError as e:
+9 -6
View File
@@ -20,8 +20,9 @@ from app.core.rewards import (
WITHDRAW_MAX_CENTS,
WITHDRAW_MIN_CENTS,
)
from app.repositories import wallet as crud_wallet
from app.models.user import User
from app.repositories import risk as risk_repo
from app.repositories import wallet as crud_wallet
from app.schemas.welfare import (
BindWechatRequest,
BindWechatResultOut,
@@ -204,6 +205,13 @@ def withdraw_info(
dependencies=[Depends(rate_limit(5, 60, "withdraw"))], # IP 级粗限流;用户级未完成单限制在仓库层
)
def withdraw(req: WithdrawRequest, user: CurrentUser, db: DbSession) -> WithdrawResultOut:
if risk_repo.is_restricted(
db,
subject_type="user",
subject_id=str(user.id),
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
):
raise HTTPException(status_code=403, detail="账号存在异常,该功能暂不可用")
# 提现发起本身不调微信(打款在审核通过后),但仍要求微信支付已配置——否则审核通过也打不了款,提前拦
if not settings.wxpay_configured:
raise HTTPException(status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail="wechat pay not configured")
@@ -222,11 +230,6 @@ def withdraw(req: WithdrawRequest, user: CurrentUser, db: DbSession) -> Withdraw
) from e
except crud_wallet.WechatNotBoundError as e:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="请先绑定微信") from e
except crud_wallet.WithdrawTooFrequentError as e:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail="已有提现申请正在审核或打款中,请处理完成后再申请",
) from e
except crud_wallet.WithdrawTierUnavailableError as e:
# 福利页档位闸(7-9):次数满/已选其他额度。正常客户端已按 tiers 预拦,此处兜底防绕过。
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="今日额度已达上限") from e
+25
View File
@@ -141,6 +141,31 @@ class Settings(BaseSettings):
SMS_DAILY_LIMIT_PER_PHONE: int = 10 # 单手机号每日发送上限(防刷 + 控费)
SMS_MAX_VERIFY_ATTEMPTS: int = 5 # 单个验证码最多校验失败次数,超过即作废(防爆破)
# ===== 短信提供商(可切换:极光 / 阿里云号码认证)=====
# jiguang(默认):本服务生成验证码,极光只负责下发,本地内存校验(自管码,现状不变)。
# aliyun:阿里云 dypns 号码认证,阿里云生成+下发+校验(Mode A);缺凭证时 /sms/* 返 503(优雅降级)。
SMS_PROVIDER: Literal["jiguang", "aliyun"] = "jiguang"
ALIYUN_SMS_ACCESS_KEY_ID: str = ""
ALIYUN_SMS_ACCESS_KEY_SECRET: str = ""
ALIYUN_SMS_SIGN_NAME: str = "" # 系统赠送签名(自定义签名下发易失败)
ALIYUN_SMS_TEMPLATE_CODE: str = "" # 赠送模板 CODE(须与赠送签名搭配)
ALIYUN_SMS_SCHEME_NAME: str = "" # 方案名(可空=默认方案);send/check 共用避免不匹配
ALIYUN_SMS_ENDPOINT: str = "dypnsapi.aliyuncs.com"
ALIYUN_SMS_CODE_LENGTH: int = 6 # 验证码位数(CodeLength 4~8)
ALIYUN_SMS_VALID_TIME_SEC: int = 300 # 验证码有效期秒(ValidTime);短信内 min 文案 = //60
ALIYUN_SMS_INTERVAL_SEC: int = 60 # 单号发送频控秒(Interval);核验免费
ALIYUN_SMS_TIMEOUT_SEC: int = 15 # 阿里云 API 读/连超时秒
@property
def aliyun_sms_configured(self) -> bool:
"""阿里云短信凭证齐全(缺则 SMS_PROVIDER=aliyun 时 /sms/* 返 503,而非启动崩)。"""
return bool(
self.ALIYUN_SMS_ACCESS_KEY_ID
and self.ALIYUN_SMS_ACCESS_KEY_SECRET
and self.ALIYUN_SMS_SIGN_NAME
and self.ALIYUN_SMS_TEMPLATE_CODE
)
# ===== 测试账号(release 包全流程联调用)=====
# 配一个固定测试手机号,专供无 SIM 卡 / 不走一键登录时打通全流程:该号登录【免短信验证码】
# (real 模式下也跳过校验)、每次登录【强制重走新手引导】,并设【每日使用次数上限】防被人
+36
View File
@@ -11,6 +11,10 @@ from typing import Any
from app.core import rewards as r
RISK_SMS_HOURLY_THRESHOLD_KEY = "risk_sms_hourly_threshold"
RISK_ONECLICK_DAILY_THRESHOLD_KEY = "risk_oneclick_daily_threshold"
RISK_COMPARE_DAILY_THRESHOLD_KEY = "risk_compare_daily_threshold"
# type 约定(给前端渲染编辑控件用):int / int_list / dict_str_int / bool / enum
# hidden=True:仍是合法可配项(业务照常 get_value / admin 可经专用端点读写),但**不在通用
# 「系统配置」页渲染**(admin/routers/config.py:list_config 按此过滤)。用于把已下线/已改由
@@ -106,4 +110,36 @@ CONFIG_DEFS: dict[str, dict[str, Any]] = {
"default 兜底未登记的模型。改价只影响之后回填的新记录,历史记录用当时价格快照。"
),
},
# 风控监控页专用配置。复用 app_config 持久化,但不在通用「系统配置」页重复展示;
# 由 /admin/api/risk-monitor/rules 专用接口整组校验、审计和即时重算当前窗口。
RISK_SMS_HOURLY_THRESHOLD_KEY: {
"default": 5,
"label": "短信设备每小时告警阈值",
"group": "风控",
"type": "int",
"min": 1,
"max": 5,
"hidden": True,
"help": "同一设备在北京时间同一自然小时内成功下发短信达到该次数时告警;不得高于现有每小时 5 次的发送上限。",
},
RISK_ONECLICK_DAILY_THRESHOLD_KEY: {
"default": 20,
"label": "一键登录设备每日告警阈值",
"group": "风控",
"type": "int",
"min": 1,
"max": 100_000,
"hidden": True,
"help": "同一设备在北京时间同一自然日内一键登录成功和失败合计达到该次数时告警。",
},
RISK_COMPARE_DAILY_THRESHOLD_KEY: {
"default": 100,
"label": "比价账户每日告警阈值",
"group": "风控",
"type": "int",
"min": 1,
"max": 100,
"hidden": True,
"help": "同一账户在北京时间同一自然日内发起比价达到该次数时告警。",
},
}
+33
View File
@@ -0,0 +1,33 @@
"""短信验证码服务 —— provider 分派入口。
对外只暴露 `send_code` / `verify_code` / `SmsError`,api 层无需关心用哪个 provider
provider `settings.SMS_PROVIDER` 选择(**每次调用读取**,支持运行时切换 + 灰度回退):
- `jiguang`(默认):自管码(本服务生成内存存/校验,极光只发) [jiguang.py](jiguang.py)
- `aliyun`:阿里云号码认证(阿里云生成+下发+校验,Mode A) [aliyun.py](aliyun.py)
mock(`SMS_MOCK=true`)与各 provider 的行为差异都封在 provider 内部;本层只做路由
拆包前本模块是单文件 `sms.py`;拆包后极光逻辑迁入 `jiguang` 子模块,行为零改动
"""
from __future__ import annotations
from app.core.config import settings
from . import aliyun, jiguang
from .base import SmsError
__all__ = ["SmsError", "send_code", "verify_code"]
def _provider():
"""按配置选 provider 模块(每次调用读 settings,支持运行时切换 / 测试注入)。"""
return aliyun if settings.SMS_PROVIDER == "aliyun" else jiguang
def send_code(phone: str) -> int:
"""发送验证码,返回距下次可发的冷却秒数;失败抛 SmsError。委托给当前 provider。"""
return _provider().send_code(phone)
def verify_code(phone: str, code: str) -> bool:
"""校验验证码,返回是否通过;provider 异常降级抛 SmsError。委托给当前 provider。"""
return _provider().verify_code(phone, code)
+193
View File
@@ -0,0 +1,193 @@
"""阿里云号码认证(dypns)短信 provider —— Mode A(阿里云托管验证码)。
与极光(自管码)最大不同:**本服务不生成/不存储验证码**,验证码由阿里云生成+存储+下发+校验
- 发码: SendSmsVerifyCode,TemplateParam `{"code":"##code##","min":...}` 占位,阿里云生成
- 校验: CheckSmsVerifyCode,阿里云返回 PASS / UNKNOWN核验免费
天然消除极光路径内存存码 worker 不共享的技术债(发码/校验可落不同 worker,阿里云统一裁决)
**唯一本地态**:per-phone 连续失败计数(`_verify_attempts`),用于复刻极光单码失败
`SMS_MAX_VERIFY_ATTEMPTS` 次即作废的防爆破语义 刻意与极光一致,避免两 provider 行为不同
导致排查困惑其多 worker 降级特性与极光现状同级;另有 API 层登录频控(设备+IP)做硬兜底
单号发送频控(冷却)交给阿里云 `Interval` 参数(命中FREQUENCY_FAIL429),本地不再维护冷却
SDK 交互隔离在 `_call_send` / `_call_check` 两个薄封装(惰性 import + 惰性建 client,仿 wxpay
惰性加载),单测 monkeypatch 这两个即可,不触真 SDK / 网络
"""
from __future__ import annotations
import json
import logging
import time
from threading import Lock
from app.core.config import settings
from .base import SmsError, mock_verify
logger = logging.getLogger("shagua.sms.aliyun")
# 阿里云路径唯一本地态:per-phone 连续失败次数(与极光同语义,防爆破)。
_verify_attempts: dict[str, int] = {} # phone -> 连续失败次数
_verify_seen: dict[str, float] = {} # phone -> 最近触碰 epoch(仅供 GC 老化)
_lock = Lock()
_GC_THRESHOLD = 10000 # 超此阈值,send 时顺手清老于验证码有效期的计数(仿极光 _gc)
# 发码错误码 → (HTTP 码, 用户提示)。未列出的一律 503(供应商不可用)。
_SEND_ERRORS: dict[str, tuple[int, str]] = {
"MOBILE_NUMBER_ILLEGAL": (400, "手机号无效"),
"BUSINESS_LIMIT_CONTROL": (429, "今日发送次数过多,请明天再试"),
"FREQUENCY_FAIL": (429, "发送过于频繁,请稍后再试"),
}
# 需运维介入的配置/开通类错误:打 critical 日志(融合认证未开通 / 参数非法)。
_SEND_CRITICAL_CODES = frozenset({"FUNCTION_NOT_OPENED", "INVALID_PARAMETERS"})
_client = None # 惰性构建的 SDK client(模块级缓存)
# ============================ 对外:发码 / 校验 ============================
def send_code(phone: str) -> int:
"""发送验证码(阿里云生成+下发)。
Returns: 距下次可发的秒数(= ALIYUN_SMS_INTERVAL_SEC,冷却由阿里云 Interval 侧执行)
Raises: SmsError(手机号无效 400 / 过频·天级流控 429 / 未配置·未开通·其他 503)
"""
if settings.SMS_MOCK:
logger.info("[SMS-aliyun-MOCK] to %s**** (不真发)", phone[:3])
return settings.ALIYUN_SMS_INTERVAL_SEC
if not settings.aliyun_sms_configured:
raise SmsError("短信服务未配置(缺阿里云凭证)", status_code=503)
result = _call_send(phone) # 传输/SDK 异常在内部抛 SmsError(503)
if result["success"] and result["code"] == "OK":
now = time.time()
with _lock:
_gc(now) # 顺手清老计数(超阈值才扫)
_verify_attempts.pop(phone, None) # 新码 = 新失败预算
_verify_seen.pop(phone, None)
logger.info("[SMS-aliyun] sent to %s****", phone[:3])
return settings.ALIYUN_SMS_INTERVAL_SEC
code = result["code"]
logger.error("[SMS-aliyun] send failed code=%s msg=%s", code, result["message"])
if code in _SEND_CRITICAL_CODES:
logger.critical("[SMS-aliyun] %s —— 需运维处理(融合认证未开通 / 参数非法)", code)
status, msg = _SEND_ERRORS.get(code, (503, "短信服务暂不可用,请稍后重试"))
raise SmsError(msg, status_code=status)
def verify_code(phone: str, code: str) -> bool:
"""校验验证码(阿里云裁决)。
- **mock**:放行任意 N 位数字(provider 无关,同极光)
- **real**:先查本地失败计数(达上限即本地作废,不调阿里云,与极光一致) CheckSmsVerifyCode:
PASS 清计数返 True(一次性);UNKNOWN 计数 +1 False;接口异常抛 SmsError(503)
"""
if settings.SMS_MOCK:
ok = mock_verify(code)
logger.info("[SMS-aliyun-MOCK] verify %s for %s****", "ok" if ok else "fail", phone[:3])
return ok
# 失败计数是 best-effort:网络调用不持锁(不能锁跨 IO),故并发下同号可能多放行个位数次。
# 无碍——API 层登录频控(设备+IP 5/时)是硬上限,阿里云码有效期 + DuplicatePolicy 亦兜底。
with _lock:
if _verify_attempts.get(phone, 0) >= settings.SMS_MAX_VERIFY_ATTEMPTS:
return False # 已作废:保持计数(直到 send_code 重置),与极光「达上限即作废」一致
result = _call_check(phone, code) # 传输/SDK 异常在内部抛 SmsError(503)
if not (result["success"] and result["code"] == "OK"):
# 接口层失败(非码错):降级 503,别误报「验证码错误」(400),便于区分排查。
logger.error("[SMS-aliyun] check failed code=%s msg=%s", result["code"], result["message"])
raise SmsError("短信服务暂不可用,请稍后重试", status_code=503)
if result["verify_result"] == "PASS":
with _lock:
_verify_attempts.pop(phone, None) # 验过即清(一次性)
_verify_seen.pop(phone, None)
return True
# UNKNOWN:码错 / 过期 → 失败计数 +1(累计到上限即作废)
with _lock:
_verify_attempts[phone] = _verify_attempts.get(phone, 0) + 1
_verify_seen[phone] = time.time()
return False
def _gc(now: float) -> None:
"""超阈值时清理老于验证码有效期的失败计数(码早已在阿里云侧失效,计数无意义)。仅持锁调用。"""
if len(_verify_attempts) <= _GC_THRESHOLD:
return
cutoff = now - settings.ALIYUN_SMS_VALID_TIME_SEC
for p in [p for p, ts in _verify_seen.items() if ts < cutoff]:
_verify_attempts.pop(p, None)
_verify_seen.pop(p, None)
# ============================ SDK 接缝(单测 monkeypatch 这两个)============================
def _get_client():
"""惰性构建 dypns SDK client(仿 wxpay 惰性加载:jiguang-only 部署不加载 alibabacloud)。"""
global _client
if _client is None:
from alibabacloud_dypnsapi20170525.client import Client
from alibabacloud_tea_openapi import models as open_api_models
cfg = open_api_models.Config(
access_key_id=settings.ALIYUN_SMS_ACCESS_KEY_ID,
access_key_secret=settings.ALIYUN_SMS_ACCESS_KEY_SECRET,
read_timeout=settings.ALIYUN_SMS_TIMEOUT_SEC * 1000, # SDK 单位 ms
connect_timeout=settings.ALIYUN_SMS_TIMEOUT_SEC * 1000,
)
cfg.endpoint = settings.ALIYUN_SMS_ENDPOINT
_client = Client(cfg)
return _client
def _call_send(phone: str) -> dict:
"""调 SendSmsVerifyCode。返回归一化 {success, code, message};import/建 client/调用 任一失败抛 SmsError(503)。"""
valid_min = max(1, settings.ALIYUN_SMS_VALID_TIME_SEC // 60)
template_param = json.dumps({"code": "##code##", "min": str(valid_min)}, ensure_ascii=False)
try:
# import + 建 req + 调用 全在 try 内:任一 provider 侧失败都归一成 503(保「provider 出问题→503」不变式)
from alibabacloud_dypnsapi20170525 import models as dypns_models
req = dypns_models.SendSmsVerifyCodeRequest(
phone_number=phone,
sign_name=settings.ALIYUN_SMS_SIGN_NAME,
template_code=settings.ALIYUN_SMS_TEMPLATE_CODE,
template_param=template_param,
code_length=settings.ALIYUN_SMS_CODE_LENGTH,
valid_time=settings.ALIYUN_SMS_VALID_TIME_SEC,
interval=settings.ALIYUN_SMS_INTERVAL_SEC,
scheme_name=settings.ALIYUN_SMS_SCHEME_NAME or None,
)
body = _get_client().send_sms_verify_code(req).body
except Exception as e:
logger.exception("[SMS-aliyun] send_sms_verify_code 调用异常 phone=%s****", phone[:3])
raise SmsError("短信服务暂不可用,请稍后重试", status_code=503) from e
return {"success": bool(body.success), "code": body.code, "message": body.message}
def _call_check(phone: str, code: str) -> dict:
"""调 CheckSmsVerifyCode。返回归一化 {success, code, message, verify_result};import/建 client/调用 任一失败抛 SmsError(503)。"""
try:
# import + 建 req + 调用 全在 try 内:任一 provider 侧失败都归一成 503(保「provider 出问题→503」不变式)
from alibabacloud_dypnsapi20170525 import models as dypns_models
req = dypns_models.CheckSmsVerifyCodeRequest(
phone_number=phone,
verify_code=code,
scheme_name=settings.ALIYUN_SMS_SCHEME_NAME or None,
)
body = _get_client().check_sms_verify_code(req).body
except Exception as e:
logger.exception("[SMS-aliyun] check_sms_verify_code 调用异常 phone=%s****", phone[:3])
raise SmsError("短信服务暂不可用,请稍后重试", status_code=503) from e
verify_result = getattr(body.model, "verify_result", None) if body.model else None
return {
"success": bool(body.success),
"code": body.code,
"message": body.message,
"verify_result": verify_result,
}
+23
View File
@@ -0,0 +1,23 @@
"""短信 provider 共享基座:业务异常 + provider 无关的 mock 校验。
provider(jiguang / aliyun) `from .base import SmsError`,api 层也从包入口拿到同一个
`SmsError` 保证无论用哪个 provider,异常类型与 HTTP 码映射语义一致
"""
from __future__ import annotations
from app.core.config import settings
class SmsError(Exception):
"""业务异常。`status_code` 决定 api 层翻成哪个 HTTP 码:
过频/每日超限 = 429(客户端等会再来),供应商不可用 = 503,手机号无效 = 400
"""
def __init__(self, message: str, status_code: int = 429) -> None:
super().__init__(message)
self.status_code = status_code
def mock_verify(code: str) -> bool:
"""mock 模式校验:放行任意 SMS_CODE_LENGTH 位数字(provider 无关,测试/开发便利,不真校验)。"""
return len(code) == settings.SMS_CODE_LENGTH and code.isdigit()
@@ -1,15 +1,15 @@
"""短信验证码服务
"""极光短信 provider(自管码 Mode B)
两种运行模式由 `SMS_MOCK` 切换:
- **mock**(开发/测试,默认):不真发短信,验证码打到日志;校验**放行任意 N 位数字**
(测试/开发便利)真实校验逻辑(比对存码 / 一次性 / 防爆破) real 分支 + 单测覆盖
- **real**(生产 `SMS_MOCK=false`):本服务生成 N 位验证码 调极光短信 REST
`/v1/messages` 发送(自定义验证码模式,极光只负责发,code 由本服务生成/保管/
- **real**(生产 `SMS_MOCK=false` `SMS_PROVIDER=jiguang`):本服务生成 N 位验证码 调极光
短信 REST `/v1/messages` 发送(自定义验证码模式,极光只负责发,code 由本服务生成/保管/
校验) 鉴权复用极光一键登录的 `JG_APP_KEY`/`JG_MASTER_SECRET`(同一极光应用)
验证码存储:**进程内存**( worker uvicorn 够用)重启丢失(用户重发即可)
worker / 多机时内存不共享 冷却校验都会失效,届时迁移到 DB/Redis
docs/待办与技术债.md
worker / 多机时内存不共享 冷却校验都会失效,届时迁移到 DB/Redis(或改用 aliyun provider,
其验证码由阿里云托管无本地存码) docs/待办与技术债.md
防刷两层(短信花钱 + `/sms/send` 在登录前无法 JWT 鉴权):
1. 单号 `SMS_SEND_INTERVAL_SEC` 冷却(本文件)
@@ -34,17 +34,9 @@ import httpx
from app.core.config import settings
logger = logging.getLogger("shagua.sms")
from .base import SmsError, mock_verify
class SmsError(Exception):
"""业务异常。`status_code` 决定 api 层翻成哪个 HTTP 码:
过频/每日超限 = 429(客户端等会再来),供应商不可用 = 503,手机号无效 = 400
"""
def __init__(self, message: str, status_code: int = 429) -> None:
super().__init__(message)
self.status_code = status_code
logger = logging.getLogger("shagua.sms.jiguang")
@dataclass
@@ -126,7 +118,7 @@ def verify_code(phone: str, code: str) -> bool:
- **real 模式**:比对本服务存的码,匹配即作废(一次性);失败累计到上限也作废(防爆破)
"""
if settings.SMS_MOCK:
ok = len(code) == settings.SMS_CODE_LENGTH and code.isdigit()
ok = mock_verify(code)
logger.info("[SMS-MOCK] verify %s for %s****", "ok" if ok else "fail", phone[:3])
return ok
+8 -7
View File
@@ -13,13 +13,7 @@ from app.models.analytics_selfstat import ( # noqa: F401
)
from app.models.app_config import AppConfig # noqa: F401
from app.models.comparison import ComparisonRecord # noqa: F401
from app.models.cps_activity import CpsActivity # noqa: F401
from app.models.cps_group import CpsGroup # noqa: F401
from app.models.cps_link import CpsClick, CpsLink # noqa: F401
from app.models.cps_order import CpsOrder # noqa: F401
from app.models.cps_wx_user import CpsWxUser # noqa: F401
from app.models.comparison_milestone import ComparisonMilestoneClaim # noqa: F401
from app.models.device import DeviceLiveness # noqa: F401
from app.models.coupon_state import ( # noqa: F401
CouponClaimEvent,
CouponClaimRecord,
@@ -27,6 +21,12 @@ from app.models.coupon_state import ( # noqa: F401
CouponPromptEngagement,
CouponSession,
)
from app.models.cps_activity import CpsActivity # noqa: F401
from app.models.cps_group import CpsGroup # noqa: F401
from app.models.cps_link import CpsClick, CpsLink # noqa: F401
from app.models.cps_order import CpsOrder # noqa: F401
from app.models.cps_wx_user import CpsWxUser # noqa: F401
from app.models.device import DeviceLiveness # noqa: F401
from app.models.feedback import Feedback # noqa: F401
from app.models.guide_video import GuideVideoPlay # noqa: F401
from app.models.inactivity import ( # noqa: F401
@@ -39,11 +39,12 @@ from app.models.launch_confirm_sample import LaunchConfirmSample # noqa: F401
from app.models.meituan_coupon import MeituanCoupon # noqa: F401
from app.models.notification import Notification # noqa: F401
from app.models.onboarding import OnboardingCompletion # noqa: F401
from app.models.phone_rebind_log import PhoneRebindLog # noqa: F401
from app.models.ops_marquee_seed import OpsMarqueeSeed # noqa: F401
from app.models.ops_stat_config import OpsStatConfig # noqa: F401
from app.models.phone_rebind_log import PhoneRebindLog # noqa: F401
from app.models.price_observation import PriceObservation # noqa: F401
from app.models.price_report import PriceReport # noqa: F401
from app.models.risk import BehaviorEvent, RiskIncident, SubjectRestriction # noqa: F401
from app.models.savings import SavingsRecord # noqa: F401
from app.models.signin import SigninRecord # noqa: F401
from app.models.store_mapping import StoreMapping # noqa: F401
+163
View File
@@ -0,0 +1,163 @@
"""通用行为事件、风险事件与主体限制。
这三张表不是风控监控页面专用表
- ``behavior_event`` 保存服务端权威的关键行为流水后续可复用于安全审计
漏斗核查和客诉排查
- ``risk_incident`` 保存规则命中后的可处置事件负责待处理/忽略/封禁/解除
生命周期
- ``subject_restriction`` 保存当前生效的主体限制统一承载设备账号等主体的
业务拦截状态
"""
from __future__ import annotations
from datetime import datetime
from sqlalchemy import (
JSON,
Boolean,
DateTime,
Index,
Integer,
String,
UniqueConstraint,
func,
true,
)
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.orm import Mapped, mapped_column
from app.db.base import Base
_JSON = JSON().with_variant(JSONB(), "postgresql")
class BehaviorEvent(Base):
"""服务端权威行为流水;只增不改。"""
__tablename__ = "behavior_event"
__table_args__ = (
Index("ix_behavior_event_type_time", "event_type", "occurred_at"),
Index(
"ix_behavior_event_subject_time",
"subject_type",
"subject_id",
"occurred_at",
),
Index("ix_behavior_event_user_time", "user_id", "occurred_at"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
event_type: Mapped[str] = mapped_column(String(64), nullable=False)
subject_type: Mapped[str] = mapped_column(String(32), nullable=False)
subject_id: Mapped[str] = mapped_column(String(128), nullable=False)
user_id: Mapped[int | None] = mapped_column(Integer, nullable=True)
device_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
device_model: Mapped[str | None] = mapped_column(String(128), nullable=True)
phone: Mapped[str | None] = mapped_column(String(20), nullable=True)
client_ip: Mapped[str | None] = mapped_column(String(64), nullable=True)
# success / failed / rejected / attempted
outcome: Mapped[str] = mapped_column(String(24), nullable=False, default="success")
reason: Mapped[str | None] = mapped_column(String(256), nullable=True)
details: Mapped[dict | None] = mapped_column(_JSON, nullable=True)
occurred_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), nullable=False, index=True
)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
class RiskIncident(Base):
"""规则命中后形成的可处置风险事件。"""
__tablename__ = "risk_incident"
__table_args__ = (
UniqueConstraint(
"rule_code",
"subject_type",
"subject_id",
"window_key",
name="uq_risk_incident_rule_subject_window",
),
Index("ix_risk_incident_rule_status", "rule_code", "status", "triggered_at"),
Index(
"ix_risk_incident_subject",
"subject_type",
"subject_id",
"triggered_at",
),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
rule_code: Mapped[str] = mapped_column(String(64), nullable=False)
event_type: Mapped[str] = mapped_column(String(64), nullable=False)
subject_type: Mapped[str] = mapped_column(String(32), nullable=False)
subject_id: Mapped[str] = mapped_column(String(128), nullable=False)
window_key: Mapped[str] = mapped_column(String(64), nullable=False)
window_start: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
window_end: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
first_event_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
triggered_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
last_event_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
event_count: Mapped[int] = mapped_column(Integer, nullable=False)
# open / ignored / blocked / resolved
status: Mapped[str] = mapped_column(String(24), nullable=False, default="open")
action_reason: Mapped[str | None] = mapped_column(String(256), nullable=True)
handled_by: Mapped[int | None] = mapped_column(Integer, nullable=True)
handled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
details: Mapped[dict | None] = mapped_column(_JSON, nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
server_default=func.now(),
onupdate=func.now(),
nullable=False,
)
class SubjectRestriction(Base):
"""主体当前限制状态;同一主体同一作用域只有一条,可封禁后再解除/重启。"""
__tablename__ = "subject_restriction"
__table_args__ = (
UniqueConstraint(
"subject_type",
"subject_id",
"scope",
name="uq_subject_restriction_subject_scope",
),
Index(
"ix_subject_restriction_lookup",
"subject_type",
"subject_id",
"scope",
"active",
),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
subject_type: Mapped[str] = mapped_column(String(32), nullable=False)
subject_id: Mapped[str] = mapped_column(String(128), nullable=False)
# auth_device / economic_account / all
scope: Mapped[str] = mapped_column(String(32), nullable=False)
active: Mapped[bool] = mapped_column(
Boolean, nullable=False, default=True, server_default=true()
)
reason: Mapped[str | None] = mapped_column(String(256), nullable=True)
incident_id: Mapped[int | None] = mapped_column(Integer, nullable=True)
created_by: Mapped[int | None] = mapped_column(Integer, nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
revoked_by: Mapped[int | None] = mapped_column(Integer, nullable=True)
revoked_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
details: Mapped[dict | None] = mapped_column(_JSON, nullable=True)
-9
View File
@@ -96,15 +96,6 @@ class WithdrawOrder(Base):
"""
__tablename__ = "withdraw_order"
__table_args__ = (
Index(
"ux_withdraw_order_user_active",
"user_id",
unique=True,
sqlite_where=text("status IN ('reviewing', 'pending')"),
postgresql_where=text("status IN ('reviewing', 'pending')"),
),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
user_id: Mapped[int] = mapped_column(
+569
View File
@@ -0,0 +1,569 @@
"""通用行为事件、风险规则与主体限制仓库。"""
from __future__ import annotations
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from zoneinfo import ZoneInfo
from sqlalchemy import func, select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.core.config_schema import (
RISK_COMPARE_DAILY_THRESHOLD_KEY,
RISK_ONECLICK_DAILY_THRESHOLD_KEY,
RISK_SMS_HOURLY_THRESHOLD_KEY,
)
from app.models.app_config import AppConfig
from app.models.comparison import ComparisonRecord
from app.models.risk import BehaviorEvent, RiskIncident, SubjectRestriction
from app.repositories import app_config
CN_TZ = ZoneInfo("Asia/Shanghai")
RULE_SMS_HOURLY = "sms_device_hourly_5"
RULE_ONECLICK_DAILY = "oneclick_device_daily_20"
RULE_COMPARE_DAILY = "compare_account_daily_100"
EVENT_SMS_SEND = "sms_send"
EVENT_SMS_LOGIN = "sms_login"
EVENT_ONECLICK_LOGIN = "oneclick_login"
SCOPE_AUTH_DEVICE = "auth_device"
SCOPE_ECONOMIC_ACCOUNT = "economic_account"
AUTO_RESOLVED_REASON = "规则阈值调整后不再命中"
MANUAL_RESET_REASON = "管理员重置报警计数"
RISK_RESET_BASELINES_KEY = "risk_monitor_reset_baselines"
@dataclass(frozen=True)
class RuleSpec:
code: str
event_type: str
subject_type: str
threshold_key: str
window: str
count_outcomes: tuple[str, ...]
RULES: dict[str, RuleSpec] = {
RULE_SMS_HOURLY: RuleSpec(
code=RULE_SMS_HOURLY,
event_type=EVENT_SMS_SEND,
subject_type="device",
threshold_key=RISK_SMS_HOURLY_THRESHOLD_KEY,
window="hour",
count_outcomes=("success",),
),
RULE_ONECLICK_DAILY: RuleSpec(
code=RULE_ONECLICK_DAILY,
event_type=EVENT_ONECLICK_LOGIN,
subject_type="device",
threshold_key=RISK_ONECLICK_DAILY_THRESHOLD_KEY,
window="day",
count_outcomes=("success", "failed"),
),
}
RULE_THRESHOLD_KEYS: dict[str, str] = {
RULE_SMS_HOURLY: RISK_SMS_HOURLY_THRESHOLD_KEY,
RULE_ONECLICK_DAILY: RISK_ONECLICK_DAILY_THRESHOLD_KEY,
RULE_COMPARE_DAILY: RISK_COMPARE_DAILY_THRESHOLD_KEY,
}
def get_rule_threshold(db: Session, rule_code: str) -> int:
"""读取规则当前阈值;配置表为空时回退上线前的 5/20/100 默认值。"""
return int(app_config.get_value(db, RULE_THRESHOLD_KEYS[rule_code]))
def utcnow() -> datetime:
return datetime.now(UTC)
def _ensure_aware(value: datetime) -> datetime:
return value.replace(tzinfo=UTC) if value.tzinfo is None else value
def get_rule_reset_at(db: Session, rule_code: str) -> datetime | None:
"""返回规则最近一次全局重置时间;异常旧值按未重置处理。"""
row = db.get(AppConfig, RISK_RESET_BASELINES_KEY)
if row is None or not isinstance(row.value, dict):
return None
raw = row.value.get(rule_code)
if not isinstance(raw, str):
return None
try:
parsed = datetime.fromisoformat(raw.replace("Z", "+00:00"))
except ValueError:
return None
return _ensure_aware(parsed).astimezone(UTC)
def reset_rule_baselines(
db: Session,
*,
rule_codes: tuple[str, ...],
reset_at: datetime,
admin_id: int,
) -> None:
"""持久化全局重置基线;不删除任何行为或比价流水。"""
row = db.get(AppConfig, RISK_RESET_BASELINES_KEY)
value = dict(row.value) if row and isinstance(row.value, dict) else {}
timestamp = _ensure_aware(reset_at).astimezone(UTC).isoformat()
value.update({rule_code: timestamp for rule_code in rule_codes})
if row is None:
db.add(
AppConfig(
key=RISK_RESET_BASELINES_KEY,
value=value,
updated_by_admin_id=admin_id,
)
)
else:
row.value = value
row.updated_by_admin_id = admin_id
db.flush()
def _window_bounds(at: datetime, window: str) -> tuple[str, datetime, datetime]:
local = _ensure_aware(at).astimezone(CN_TZ)
if window == "hour":
start_local = local.replace(minute=0, second=0, microsecond=0)
key = start_local.strftime("%Y-%m-%dT%H")
end_local = start_local + timedelta(hours=1)
elif window == "day":
start_local = local.replace(hour=0, minute=0, second=0, microsecond=0)
key = start_local.strftime("%Y-%m-%d")
end_local = start_local + timedelta(days=1)
else: # pragma: no cover - 规则声明错误应尽早暴露
raise ValueError(f"unsupported risk window: {window}")
return key, start_local.astimezone(UTC), end_local.astimezone(UTC)
def _event_stats(
db: Session,
spec: RuleSpec,
*,
subject_id: str,
start: datetime,
end: datetime,
threshold: int,
) -> tuple[int, datetime | None, datetime | None, datetime | None]:
filters = (
BehaviorEvent.event_type == spec.event_type,
BehaviorEvent.subject_type == spec.subject_type,
BehaviorEvent.subject_id == subject_id,
BehaviorEvent.outcome.in_(spec.count_outcomes),
BehaviorEvent.occurred_at >= start,
BehaviorEvent.occurred_at < end,
)
count, first_at, last_at = db.execute(
select(
func.count(BehaviorEvent.id),
func.min(BehaviorEvent.occurred_at),
func.max(BehaviorEvent.occurred_at),
).where(*filters)
).one()
triggered_at = None
if int(count or 0) >= threshold:
triggered_at = db.execute(
select(BehaviorEvent.occurred_at)
.where(*filters)
.order_by(BehaviorEvent.occurred_at.asc(), BehaviorEvent.id.asc())
.offset(threshold - 1)
.limit(1)
).scalar_one()
return int(count or 0), first_at, last_at, triggered_at
def _upsert_incident(
db: Session,
*,
rule_code: str,
event_type: str,
subject_type: str,
subject_id: str,
window_key: str,
window_start: datetime,
window_end: datetime,
first_event_at: datetime,
triggered_at: datetime,
last_event_at: datetime,
event_count: int,
details: dict | None = None,
) -> RiskIncident:
incident = db.execute(
select(RiskIncident).where(
RiskIncident.rule_code == rule_code,
RiskIncident.subject_type == subject_type,
RiskIncident.subject_id == subject_id,
RiskIncident.window_key == window_key,
)
).scalar_one_or_none()
if incident is None:
incident = RiskIncident(
rule_code=rule_code,
event_type=event_type,
subject_type=subject_type,
subject_id=subject_id,
window_key=window_key,
window_start=window_start,
window_end=window_end,
first_event_at=first_event_at,
triggered_at=triggered_at,
last_event_at=last_event_at,
event_count=event_count,
status="open",
details=details,
)
# 并发首次命中可能同时插入;保存点只回滚重复 incident,不丢行为流水。
try:
with db.begin_nested():
db.add(incident)
db.flush()
except IntegrityError:
incident = db.execute(
select(RiskIncident).where(
RiskIncident.rule_code == rule_code,
RiskIncident.subject_type == subject_type,
RiskIncident.subject_id == subject_id,
RiskIncident.window_key == window_key,
)
).scalar_one()
else:
reopened = incident.status == "resolved" and incident.action_reason in (
AUTO_RESOLVED_REASON,
MANUAL_RESET_REASON,
)
if reopened:
incident.status = "open"
incident.action_reason = None
incident.handled_by = None
incident.handled_at = None
# 同一自然窗口内重置后会复用唯一 incident;明细窗口必须同步切到新基线,
# 否则展开时会把重置前的旧流水也混进来。
incident.window_start = window_start
incident.window_end = window_end
incident.first_event_at = first_event_at
if incident.status == "open":
incident.triggered_at = triggered_at
incident.last_event_at = last_event_at
incident.event_count = event_count
# 已忽略/封禁的事件只更新事实数据,不重新打开。
if details:
incident.details = {**(incident.details or {}), **details}
return incident
def evaluate_behavior_rule(
db: Session, *, rule_code: str, subject_id: str, at: datetime
) -> RiskIncident | None:
spec = RULES[rule_code]
threshold = get_rule_threshold(db, rule_code)
window_key, window_start, end = _window_bounds(at, spec.window)
reset_at = get_rule_reset_at(db, rule_code)
start = max(window_start, reset_at) if reset_at else window_start
count, first_at, last_at, triggered_at = _event_stats(
db,
spec,
subject_id=subject_id,
start=start,
end=end,
threshold=threshold,
)
if count < threshold or first_at is None or last_at is None or triggered_at is None:
return None
return _upsert_incident(
db,
rule_code=spec.code,
event_type=spec.event_type,
subject_type=spec.subject_type,
subject_id=subject_id,
window_key=window_key,
window_start=start,
window_end=end,
first_event_at=first_at,
triggered_at=triggered_at,
last_event_at=last_at,
event_count=count,
)
def reconcile_behavior_rule(
db: Session,
*,
rule_code: str,
at: datetime | None = None,
commit: bool = True,
) -> int:
"""按当前阈值重算短信/一键登录当前窗口,并收起已不再命中的待处理告警。"""
spec = RULES[rule_code]
current = at or utcnow()
threshold = get_rule_threshold(db, rule_code)
window_key, window_start, end = _window_bounds(current, spec.window)
reset_at = get_rule_reset_at(db, rule_code)
start = max(window_start, reset_at) if reset_at else window_start
filters = (
BehaviorEvent.event_type == spec.event_type,
BehaviorEvent.subject_type == spec.subject_type,
BehaviorEvent.outcome.in_(spec.count_outcomes),
BehaviorEvent.occurred_at >= start,
BehaviorEvent.occurred_at < end,
)
qualifying_rows = db.execute(
select(
BehaviorEvent.subject_id,
func.max(BehaviorEvent.occurred_at),
)
.where(*filters)
.group_by(BehaviorEvent.subject_id)
.having(func.count(BehaviorEvent.id) >= threshold)
).all()
qualifying = {str(subject_id) for subject_id, _ in qualifying_rows}
for subject_id, last_at in qualifying_rows:
evaluate_behavior_rule(
db,
rule_code=rule_code,
subject_id=str(subject_id),
at=last_at or current,
)
open_incidents = db.scalars(
select(RiskIncident).where(
RiskIncident.rule_code == rule_code,
RiskIncident.window_key == window_key,
RiskIncident.status == "open",
)
).all()
for incident in open_incidents:
if incident.subject_id not in qualifying:
incident.status = "resolved"
incident.action_reason = AUTO_RESOLVED_REASON
incident.handled_by = None
incident.handled_at = utcnow()
if commit:
db.commit()
return len(qualifying)
def record_behavior_event(
db: Session,
*,
event_type: str,
subject_type: str,
subject_id: str,
user_id: int | None = None,
device_id: str | None = None,
device_model: str | None = None,
phone: str | None = None,
client_ip: str | None = None,
outcome: str = "success",
reason: str | None = None,
details: dict | None = None,
occurred_at: datetime | None = None,
evaluate_rule: str | None = None,
commit: bool = True,
) -> BehaviorEvent:
at = occurred_at or utcnow()
event = BehaviorEvent(
event_type=event_type,
subject_type=subject_type,
subject_id=subject_id,
user_id=user_id,
device_id=device_id,
device_model=device_model,
phone=phone,
client_ip=client_ip,
outcome=outcome,
reason=reason,
details=details,
occurred_at=at,
)
db.add(event)
db.flush()
if evaluate_rule:
evaluate_behavior_rule(db, rule_code=evaluate_rule, subject_id=subject_id, at=at)
if commit:
db.commit()
db.refresh(event)
return event
def sync_compare_incident(
db: Session,
*,
user_id: int,
at: datetime,
threshold: int | None = None,
commit: bool = True,
) -> RiskIncident | None:
effective_threshold = threshold or get_rule_threshold(db, RULE_COMPARE_DAILY)
# comparison_record 的既有写入口统一落“北京时间 naive”时间;这里必须沿用同一
# 口径,否则 SQLite/PG session timezone 不同时会把凌晨记录算到前一天。
local = at.astimezone(CN_TZ).replace(tzinfo=None) if at.tzinfo else at
window_start = local.replace(hour=0, minute=0, second=0, microsecond=0)
end = window_start + timedelta(days=1)
window_key = window_start.strftime("%Y-%m-%d")
reset_at = get_rule_reset_at(db, RULE_COMPARE_DAILY)
reset_local = (
reset_at.astimezone(CN_TZ).replace(tzinfo=None) if reset_at else None
)
start = max(window_start, reset_local) if reset_local else window_start
filters = (
ComparisonRecord.user_id == user_id,
ComparisonRecord.created_at >= start,
ComparisonRecord.created_at < end,
)
count, first_at, last_at = db.execute(
select(
func.count(ComparisonRecord.id),
func.min(ComparisonRecord.created_at),
func.max(ComparisonRecord.created_at),
).where(*filters)
).one()
if int(count or 0) < effective_threshold or first_at is None or last_at is None:
return None
triggered_at = db.execute(
select(ComparisonRecord.created_at)
.where(*filters)
.order_by(ComparisonRecord.created_at.asc(), ComparisonRecord.id.asc())
.offset(effective_threshold - 1)
.limit(1)
).scalar_one()
incident = _upsert_incident(
db,
rule_code=RULE_COMPARE_DAILY,
event_type="compare_start",
subject_type="user",
subject_id=str(user_id),
window_key=window_key,
window_start=start,
window_end=end,
first_event_at=first_at,
triggered_at=triggered_at,
last_event_at=last_at,
event_count=int(count),
)
if commit:
db.commit()
db.refresh(incident)
return incident
def reconcile_compare_rule(
db: Session, *, at: datetime | None = None, commit: bool = True
) -> int:
"""按当前阈值重算北京时间当日比价告警,并收起不再命中的待处理告警。"""
current = (at or utcnow()).astimezone(CN_TZ).replace(tzinfo=None)
window_start = current.replace(hour=0, minute=0, second=0, microsecond=0)
end = window_start + timedelta(days=1)
window_key = window_start.strftime("%Y-%m-%d")
reset_at = get_rule_reset_at(db, RULE_COMPARE_DAILY)
reset_local = (
reset_at.astimezone(CN_TZ).replace(tzinfo=None) if reset_at else None
)
start = max(window_start, reset_local) if reset_local else window_start
threshold = get_rule_threshold(db, RULE_COMPARE_DAILY)
rows = db.execute(
select(
ComparisonRecord.user_id,
func.max(ComparisonRecord.created_at),
)
.where(
ComparisonRecord.user_id.is_not(None),
ComparisonRecord.created_at >= start,
ComparisonRecord.created_at < end,
)
.group_by(ComparisonRecord.user_id)
.having(func.count(ComparisonRecord.id) >= threshold)
).all()
qualifying = {str(user_id) for user_id, _ in rows}
for user_id, last_at in rows:
sync_compare_incident(
db,
user_id=int(user_id),
at=last_at or current,
threshold=threshold,
commit=False,
)
open_incidents = db.scalars(
select(RiskIncident).where(
RiskIncident.rule_code == RULE_COMPARE_DAILY,
RiskIncident.window_key == window_key,
RiskIncident.status == "open",
)
).all()
for incident in open_incidents:
if incident.subject_id not in qualifying:
incident.status = "resolved"
incident.action_reason = AUTO_RESOLVED_REASON
incident.handled_by = None
incident.handled_at = utcnow()
if commit:
db.commit()
return len(qualifying)
def get_active_restriction(
db: Session, *, subject_type: str, subject_id: str, scope: str
) -> SubjectRestriction | None:
return db.execute(
select(SubjectRestriction).where(
SubjectRestriction.subject_type == subject_type,
SubjectRestriction.subject_id == subject_id,
SubjectRestriction.scope.in_((scope, "all")),
SubjectRestriction.active.is_(True),
)
).scalar_one_or_none()
def is_restricted(db: Session, *, subject_type: str, subject_id: str, scope: str) -> bool:
return get_active_restriction(
db, subject_type=subject_type, subject_id=subject_id, scope=scope
) is not None
def activate_restriction(
db: Session,
*,
subject_type: str,
subject_id: str,
scope: str,
reason: str,
incident_id: int | None,
admin_id: int,
) -> SubjectRestriction:
restriction = db.execute(
select(SubjectRestriction).where(
SubjectRestriction.subject_type == subject_type,
SubjectRestriction.subject_id == subject_id,
SubjectRestriction.scope == scope,
)
).scalar_one_or_none()
if restriction is None:
restriction = SubjectRestriction(
subject_type=subject_type,
subject_id=subject_id,
scope=scope,
)
db.add(restriction)
restriction.active = True
restriction.reason = reason
restriction.incident_id = incident_id
restriction.created_by = admin_id
restriction.created_at = utcnow()
restriction.revoked_by = None
restriction.revoked_at = None
db.flush()
return restriction
def revoke_restriction(
db: Session, *, restriction: SubjectRestriction, admin_id: int
) -> None:
restriction.active = False
restriction.revoked_by = admin_id
restriction.revoked_at = utcnow()
db.flush()
+2 -23
View File
@@ -35,7 +35,6 @@ from app.services import notification_events
_WX_STATE_SUCCESS = "SUCCESS"
_WX_STATE_FAILED = {"FAIL", "CANCELLED", "CLOSED"}
_WX_STATE_WAIT_CONFIRM = "WAIT_USER_CONFIRM" # 用户还没在微信确认页确认
_WITHDRAW_ACTIVE_STATUSES = {"reviewing", "pending"}
# 占用新人档「一次性」资格的提现状态:进行中(reviewing/pending)或成功打款(success)。
# 被拒/转账失败/解绑退回(rejected/failed,均已退款、钱没到手)不在此列 → 新人档恢复可提
# (2026-07-16 修正:此前判定不看状态,解绑微信退回后 0.1 被误判已用、资格永久锁死)。
@@ -69,10 +68,6 @@ class InsufficientCashError(Exception):
"""现金余额不足。"""
class WithdrawTooFrequentError(Exception):
"""提现申请过于频繁,或已有未完成提现单。"""
class WithdrawTierUnavailableError(Exception):
"""该档位今日不可提:次数已满,或今天已选了其他额度(7-9 福利页档位规则)。"""
@@ -755,17 +750,8 @@ def create_withdraw(
else:
out_bill_no = uuid.uuid4().hex
active_order_id = db.execute(
select(WithdrawOrder.id).where(
WithdrawOrder.user_id == user_id,
WithdrawOrder.status.in_(_WITHDRAW_ACTIVE_STATUSES),
).limit(1)
).scalar_one_or_none()
if active_order_id is not None:
raise WithdrawTooFrequentError
# 福利页档位闸(7-9):coin_cash 只能提预设档位,且该档今日可提(服务端权威口径,防绕过
# 客户端刷)。放在幂等返回/在途互斥之后:同号重试仍原样返回旧单,不被档位闸误杀。
# 客户端刷)。放在幂等返回之后:同号重试仍原样返回旧单,不被档位闸误杀。
# allow_sub_min(0.01 调试直发)保持原样放行,不受档位约束;invite_cash 本轮无档位概念不校验。
if source == "coin_cash" and not allow_sub_min:
tier_state = next(
@@ -815,6 +801,7 @@ def create_withdraw(
db.commit()
except IntegrityError:
db.rollback()
# 唯一冲突只可能来自 out_bill_no 幂等键并发重试:原样返回既有单;否则未知冲突,上抛。
existing = db.execute(
select(WithdrawOrder).where(
WithdrawOrder.out_bill_no == out_bill_no, WithdrawOrder.user_id == user_id
@@ -822,14 +809,6 @@ def create_withdraw(
).scalar_one_or_none()
if existing is not None:
return existing
active_order_id = db.execute(
select(WithdrawOrder.id).where(
WithdrawOrder.user_id == user_id,
WithdrawOrder.status.in_(_WITHDRAW_ACTIVE_STATUSES),
).limit(1)
).scalar_one_or_none()
if active_order_id is not None:
raise WithdrawTooFrequentError from None
raise
db.refresh(order)
return order # 待管理员审核;**不在此处打款**
+9 -1
View File
@@ -11,7 +11,6 @@ from datetime import datetime
from pydantic import BaseModel, ConfigDict, Field
# ===== 用户对外信息 =====
class UserOut(BaseModel):
@@ -65,6 +64,9 @@ class JverifyLoginRequest(BaseModel):
"", max_length=64,
description="硬件级设备标识(Android ANDROID_ID),用于新手引导按 设备+账号 去重;空=按未完成处理",
)
device_model: str = Field(
"", max_length=128, description="客户端设备型号快照,用于登录安全审计"
)
# ===== 短信验证码 =====
@@ -75,6 +77,9 @@ class SmsSendRequest(BaseModel):
"", max_length=64,
description="硬件级设备标识(Android ANDROID_ID),用于发码防刷按 设备+IP 限流;空=按 IP 聚一桶",
)
device_model: str = Field(
"", max_length=128, description="客户端设备型号快照,用于短信安全审计"
)
class SmsSendResponse(BaseModel):
@@ -90,6 +95,9 @@ class SmsLoginRequest(BaseModel):
"", max_length=64,
description="硬件级设备标识(Android ANDROID_ID),用于新手引导按 设备+账号 去重;空=按未完成处理",
)
device_model: str = Field(
"", max_length=128, description="客户端设备型号快照,用于短信验证安全审计"
)
# ===== Refresh =====
@@ -0,0 +1,224 @@
CheckSmsVerifyCode - 核验验证码
更新时间:2026年3月19日 20:02:53
核验短信验证码并返回核验是否成功的结果。
调试
您可以在OpenAPI Explorer中直接运行该接口,免去您计算签名的困扰。运行成功后,OpenAPI Explorer可以自动生成SDK代码示例。
调试
授权信息
下表是API对应的授权信息,可以在RAM权限策略语句的Action元素中使用,用来给RAM用户或RAM角色授予调用此API的权限。具体说明如下:
操作:是指具体的权限点。
访问级别:是指每个操作的访问级别,取值为写入(Write)、读取(Read)或列出(List)。
资源类型:是指操作中支持授权的资源类型。具体说明如下:
对于必选的资源类型,用前面加 * 表示。
对于不支持资源级授权的操作,用全部资源表示。
条件关键字:是指云产品自身定义的条件关键字。
关联操作:是指成功执行操作所需要的其他权限。操作者必须同时具备关联操作的权限,操作才能成功。
放大查看
操作
访问级别
资源类型
条件关键字
关联操作
dypns:CheckSmsVerifyCode
none
*全部资源
*
无 无
请求参数
放大查看
名称
类型
必填
描述
示例值
SchemeName
string
方案名称,如果不填则为“默认方案”。最多不超过 20 个字符。
重要 如果发送接口的方案名称不为空,请确保该参数不为空且与发送接口的方案名称参数一致
测试方案
CountryCode
string
号码国家编码,默认为 86。
86
PhoneNumber
string
手机号。
186****0000
OutId
string
外部流水号。
12123231
VerifyCode
string
验证码。
说明
SendSmsVerifyCode 接口的字段 TemplateParam,配置方式有 2 种:
{"code":"##code##","min":"5"}
{"code":"123456","min":"5"}
{"code":"##code##","min":"5"}验证码是 api 动态生成的,阿里云接口可以完成校验。
{"code":"123456","min":"5"}验证码是用户配置的不是 api 动态生成,阿里云接口无法校验。
请您按照实际情况传入对应的验证码。
1231
CaseAuthPolicy
integer
验证码大小写字母核验策略。取值:
1:不区分大小写。
2:区分大小写。
1
返回参数
放大查看
名称
类型
描述
示例值
object
AccessDeniedDetail
string
访问被拒绝详细信息。
Message
string
状态码的描述。
成功
Model
object
请求结果数据。
OutId
string
外部流水号。
1212312
VerifyResult
string
短信验证码核验结果。取值:
PASS:短信验证码核验成功。
UNKNOWN:短信验证码核验失败。
PASS
Code
string
接口请求状态码。
返回 OK 代表请求成功。
其他错误码,请参见返回码。
重要 接口请求成功不代表短信验证码核验成功,短信验证码核验结果仅以Model.VerifyResult参数返回值为准。
OK
Success
boolean
接口调用是否成功。取值:
true:接口调用成功。
false:接口调用失败。
重要 接口调用成功不代表短信验证码核验成功,短信验证码核验结果仅以Model.VerifyResult参数返回值为准。
true
RequestId
string
CF8854E5-DB21-3E5D-A9B1-DDC752FD7384
示例
正常返回示例
JSON格式
放大查看复制代码
{
"AccessDeniedDetail": "无",
"Message": "成功",
"Model": {
"OutId": "1212312",
"VerifyResult": "PASS"
},
"Code": "OK",
"Success": true,
"RequestId": "CF8854E5-DB21-3E5D-A9B1-DDC752FD7384"
}
@@ -0,0 +1,396 @@
SendSmsVerifyCode - 发送短信验证码
更新时间:2026年7月3日 09:54:53
发送短信验证码。
接口说明
由于运营商近期加强对短信签名的管控。您自定义的签名面临下发失败问题,推荐您使用号码认证控制台赠送的短信签名和模板进行短信认证。系统赠送签名必须搭配系统赠送模板使用。
请确保在使用该接口前,已充分了解号码认证服务产品的收费方式和价格,短信认证服务仅收取短信发送费用(按运营商回执状态计费,短信提交成功但运营商回执失败时不计费),核验服务免费。
调试
您可以在OpenAPI Explorer中直接运行该接口,免去您计算签名的困扰。运行成功后,OpenAPI Explorer可以自动生成SDK代码示例。
调试
授权信息
下表是API对应的授权信息,可以在RAM权限策略语句的Action元素中使用,用来给RAM用户或RAM角色授予调用此API的权限。具体说明如下:
操作:是指具体的权限点。
访问级别:是指每个操作的访问级别,取值为写入(Write)、读取(Read)或列出(List)。
资源类型:是指操作中支持授权的资源类型。具体说明如下:
对于必选的资源类型,用前面加 * 表示。
对于不支持资源级授权的操作,用全部资源表示。
条件关键字:是指云产品自身定义的条件关键字。
关联操作:是指成功执行操作所需要的其他权限。操作者必须同时具备关联操作的权限,操作才能成功。
放大查看
操作
访问级别
资源类型
条件关键字
关联操作
dypns:SendSmsVerifyCode
create
*全部资源
*
无 无
请求参数
放大查看
名称
类型
必填
描述
示例值
SchemeName
string
方案名称,如果不填则为“默认方案”。最多不超过 20 个字符。
测试方案
CountryCode
string
号码国家编码。默认为 86,目前也仅支持国内号码发送。
86
PhoneNumber
string
短信接收方手机号。
130****0000
SignName
string
签名名称。暂不支持使用自定义签名,请使用系统赠送的签名,您可在赠送签名配置页面选择需要下发的签名。
恒创联众
TemplateCode
string
短信模板 CODE。参数SignName选择赠送签名时,必须搭配赠送模板下发短信。您可在赠送模板配置页面选择适用您业务场景的模板。
100001
TemplateParam
string
短信模板参数。验证码位置有两种传值方式:
可使用"##code##"替代,由参数 CodeType 指定验证码生成规则;
也可直接传入具体的验证码值,直接下发至接收方。
示例:如模板内容为:“您的验证码是${code},有效期${min}分钟,请勿告诉他人。”。
重要 上文中的 code 请替换成您实际申请的验证码模板中的参数名称
该字段可传入{"code":"##code##","min":"5"}由系统根据规则生成验证码;
或直接传入指定的验证码值{"code":"123456","min":"5"}。
说明
{"code":"##code##","min":"5"}验证码是 api 动态生成的,阿里云接口可以完成校验。
{"code":"123456","min":"5"}验证码是用户配置的不是 api 动态生成,阿里云接口无法校验。
说明
如果 JSON 中需要带换行符,请参照标准的 JSON 协议处理。
模板变量规范,请参见短信模板规范。
{"code":"##code##","min":"5"}
SmsUpExtendCode
string
上行短信扩展码。上行短信指发送给通信服务提供商的短信,用于定制某种服务、完成查询,或是办理某种业务等,需要收费,按运营商普通短信资费进行扣费。
说明
扩展码是生成签名时系统自动默认生成的,不支持自行传入。无特殊需要此字段的用户请忽略此字段。如需使用,请联系您的商务经理。
1213123
OutId
string
外部流水号。
外部流水号(透传)
CodeLength
integer
验证码长度支持 4~8 位长度,默认是 4 位。
4
ValidTime
integer
验证码有效时长,单位秒,默认为 300 秒。
300
DuplicatePolicy
integer
核验规则,当有效时间内对同场景内的同号码重复发送验证码时,旧验证码如何处理。
1:覆盖处理(默认),即旧验证码会失效掉。
2:保留,即多个验证码都是在有效期内都可以校验通过。
枚举值:
1 :
覆盖
2 :
保留
1
Interval
integer
时间间隔,单位:秒。即多久间隔可以发送一次验证码,用于频控,默认 60 秒。
60
CodeType
integer
生成的验证码类型。当参数 TemplateParam 传入占位符时,此参数必填,将由系统根据指定的规则生成验证码。取值:
1:纯数字(默认)。
2:纯大写字母。
3:纯小写字母。
4:大小字母混合。
5:数字+大写字母混合。
6:数字+小写字母混合。
7:数字+大小写字母混合。
枚举值:
1 :
纯数字
2 :
纯大写字母
3 :
纯小写字母
4 :
大小字母混合
5 :
数字+大写字母混合
6 :
数字+小写字母混合
7 :
数字+大小写字母混合
1
ReturnVerifyCode
boolean
是否返回验证码。取值:
true:返回。
false:不返回。
true
AutoRetry
integer
是否自动替换签名重试(默认开启),可取值:
1 开启自动重试功能,开启后,在验证码有效期内,当运营商返回明确的失败状态时,允许阿里云尽可能的尝试使用其他方式发送验证码,以提升发送成功率。其他方式包括且不限于:通过其他运营商重试、更换签名重试等
0 不开启自动重试
是否自动重试
返回参数
放大查看
名称
类型
描述
示例值
object
AccessDeniedDetail
string
访问被拒绝详细信息。
Message
string
状态码的描述。
成功
RequestId
string
请求 ID。
CC3BB6D2-2FDF-4321-9DCE-B38165CE4C47
Model
object
请求结果数据。
VerifyCode
string
验证码。
4232
RequestId
string
请求 ID。
a3671ccf-0102-4c8e-8797-a3678e091d09
OutId
string
外部流水号。
1231231313
BizId
string
业务 ID。
112231421412414124123^4
Code
string
请求状态码。返回 OK 代表请求成功。其他错误码,请参见返回码列表。
OK
Success
boolean
请求是否成功。
true:请求成功。
false:请求失败。
true
示例
正常返回示例
JSON格式
放大查看复制代码
{
"AccessDeniedDetail": "无",
"Message": "成功 ",
"RequestId": "CC3BB6D2-2FDF-4321-9DCE-B38165CE4C47",
"Model": {
"VerifyCode": "4232",
"RequestId": "a3671ccf-0102-4c8e-8797-a3678e091d09",
"OutId": "1231231313",
"BizId": "112231421412414124123^4"
},
"Code": "OK",
"Success": true
}
错误码
放大查看
HTTP status code
错误码
错误信息
描述
400 MOBILE_NUMBER_ILLEGAL The mobile number is illegal. 手机号码格式错误
400 BUSINESS_LIMIT_CONTROL The number has exceeded the limit for the day. 触发号码天级流控
400 FREQUENCY_FAIL Check frequency fail. 频控校验未通过
400 INVALID_PARAMETERS parameter is not valid. 非法参数
400 FUNCTION_NOT_OPENED You have not opened this function. 没有开通融合认证功能
+15 -2
View File
@@ -1,9 +1,22 @@
# 短信验证码(sms
> 文件:`app/integrations/sms.py` | 关联接口:[auth-sms-send](../api/auth-sms-send.md) · [auth-sms-login](../api/auth-sms-login.md) | [← 集成索引](./README.md)
> 文件:`app/integrations/sms/`(分派器 `__init__` + `jiguang` / `aliyun` provider + `base`) | 关联接口:[auth-sms-send](../api/auth-sms-send.md) · [auth-sms-login](../api/auth-sms-login.md) | [← 集成索引](./README.md)
## 作用
手机号 + 验证码登录的验证码发送 / 校验。**已接极光短信 REST**,由 `SMS_MOCK` 切 mock / real。
手机号 + 验证码登录的验证码发送 / 校验。支持**可切换 provider**(`SMS_PROVIDER`):`jiguang`(默认,极光自管码)/ `aliyun`(阿里云号码认证托管码)。`SMS_MOCK` 切 mock / real。
## 短信提供商(`SMS_PROVIDER`,可切换 + 灰度回退)
| | `jiguang`(默认) | `aliyun` |
|---|---|---|
| 验证码 | **本服务生成**、极光只下发、**本地内存校验** | **阿里云生成 + 下发 + 校验**(dypns 号码认证,Mode A,核验免费) |
| 发码 | 极光 `/v1/messages` | `SendSmsVerifyCode`(`##code##` 占位) |
| 校验 | 比对本地存码 | `CheckSmsVerifyCode``PASS` / `UNKNOWN` |
| 多 worker | ⚠️ 内存存码不共享(见已知局限) | ✅ 阿里云托管,天然共享 |
| 防爆破 | 单码失败 `SMS_MAX_VERIFY_ATTEMPTS` 次即作废 | **同语义**(本地 per-phone 失败计数,刻意与极光一致) |
| 单号冷却 | 本地 `SMS_SEND_INTERVAL_SEC` | 交给阿里云 `Interval` |
| 校验降级 | 本地无网络调用 | 阿里云接口异常 → `/sms/login`**503**(非 400,便于排查) |
`aliyun` 需在**号码认证控制台开通「融合认证」**,用系统赠送签名 + 赠送模板;配置见 `.env.example``ALIYUN_SMS_*`,SDK 为 `alibabacloud_dypnsapi20170525`。**以下章节描述 `jiguang` provider(自管码)细节**。
| | mock(`SMS_MOCK=true`,默认 / 开发测试) | real(`SMS_MOCK=false`,生产) |
|---|---|---|
@@ -0,0 +1,323 @@
# 允许在途提现时继续提交新申请 Implementation Plan
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** 取消「同一用户同一时刻仅一笔在途提现」限制,已有 `reviewing`/`pending` 单时允许继续提交新的提现申请。
**Architecture:** 该限制由两道闸共同强制——应用层 `create_withdraw` 的在途单检查(`WithdrawTooFrequentError`)与数据库分区唯一索引 `ux_withdraw_order_user_active`。彻底移除两者 + 清理随之失效的死代码;既有约束(建单先扣款、`coin_cash` 每日档位次数、`out_bill_no` 幂等)天然保留,无需改动。测试库由 `Base.metadata.create_all()` 依模型建表,故删模型内索引定义即让测试反映新 schema;另配一条 Alembic 迁移让真实库(dev/prod)落地同一变更。
**Tech Stack:** FastAPI · SQLAlchemy 2.0 · Alembic · pytest · ruff
规格来源:[docs/superpowers/specs/2026-07-24-withdraw-allow-concurrent-design.md](../specs/2026-07-24-withdraw-allow-concurrent-design.md)
---
## 关键事实(实现依据)
- 当前 Alembic head:`d8dd2106e438`(新迁移的 `down_revision`)。
- 索引出处:[`alembic/versions/withdraw_safety_indexes.py`](../../../alembic/versions/withdraw_safety_indexes.py) 同时建了两个索引——本次**只删** `ux_withdraw_order_user_active`,**保留**姊妹索引 `ux_cash_transaction_withdraw_refund_ref`(退款幂等)。
- 档位:`WithdrawTier(50, "0.5", None, 3, False)`——50 分(0.5 元)是**常规档**,每日 3 次,非新人档;测试用它来造多笔在途。
- 测试库走 `create_all`(见 `tests/conftest.py`),**不跑 Alembic**;故迁移的正确性由本计划单独的 upgrade/downgrade 回环验证,不由 pytest 覆盖。
- `app/models/wallet.py``Index`/`text` 仍被其它表(行 54/186/224)使用,删本表 `__table_args__` 后**无需**清理 import。
---
## Task 1: 允许多笔在途提现并存(TDD:模型 + 应用层 + 端点 + 迁移,单次提交)
本变更是一次原子的 schema+行为改动:模型内索引、应用层检查、Alembic 迁移相互依赖,任一缺失都会让"多笔在途"在测试库或真实库其一不成立。故作为**一个任务、一次提交**完成,内部按 TDD 分步。
**Files:**
- Test: `tests/test_withdraw.py`(新增 2 个用例)
- Modify: `app/repositories/wallet.py`(删在途单检查 + IntegrityError 兜底瘦身 + 删死常量/异常)
- Modify: `app/api/v1/wallet.py:225-229`(删失效的 409 处理)
- Modify: `app/models/wallet.py:99-107`(删分区唯一索引)
- Create: `alembic/versions/drop_withdraw_active_unique_index.py`(真实库删索引)
---
- [ ] **Step 1: 写两个失败测试**
`tests/test_withdraw.py` 末尾追加(复用文件内既有 helper `_login`/`_auth`/`_seed_cash`/`_patch_userinfo`):
```python
def test_withdraw_multiple_in_flight_allowed(client, monkeypatch) -> None:
"""取消「同时仅一单」:已有在途(reviewing)时,不同 out_bill_no 可继续提交,两单并存。"""
_patch_userinfo(monkeypatch, "openid_multi_inflight")
token = _login(client, "13800002020")
_seed_cash(client, token, "13800002020", 100) # 够两笔 0.5
client.post("/api/v1/wallet/bind-wechat", json={"code": "c"}, headers=_auth(token))
r1 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billmulti00000001"},
headers=_auth(token),
)
assert r1.status_code == 200, r1.text
assert r1.json()["status"] == "reviewing"
r2 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billmulti00000002"},
headers=_auth(token),
)
assert r2.status_code == 200, r2.text
assert r2.json()["status"] == "reviewing"
# 两张在途单并存
r = client.get("/api/v1/wallet/withdraw-orders", headers=_auth(token))
reviewing = [o for o in r.json()["items"] if o["status"] == "reviewing"]
assert len(reviewing) == 2, r.text
# 余额扣两次:100-50-50=0
r = client.get("/api/v1/wallet/account", headers=_auth(token))
assert r.json()["cash_balance_cents"] == 0
def test_withdraw_second_blocked_only_by_insufficient_cash(client, monkeypatch) -> None:
"""并行放开后第二笔仅受余额约束:余额不足返 409「现金余额不足」,而非旧的「已有提现」拦截。"""
_patch_userinfo(monkeypatch, "openid_multi_insuff")
token = _login(client, "13800002021")
_seed_cash(client, token, "13800002021", 50) # 仅够一笔 0.5
client.post("/api/v1/wallet/bind-wechat", json={"code": "c"}, headers=_auth(token))
r1 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billinsuff0000001"},
headers=_auth(token),
)
assert r1.status_code == 200, r1.text
r2 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billinsuff0000002"},
headers=_auth(token),
)
assert r2.status_code == 409, r2.text
assert "现金余额不足" in r2.json()["detail"]
```
- [ ] **Step 2: 运行新测试,确认失败**
Run: `pytest tests/test_withdraw.py::test_withdraw_multiple_in_flight_allowed tests/test_withdraw.py::test_withdraw_second_blocked_only_by_insufficient_cash -q`
Expected: 两条 FAIL —— `multiple_in_flight` 因第二笔被拦返回 409(期望 200);`second_blocked` 因返回的 409 detail 是「已有提现申请正在审核或打款中」而非「现金余额不足」。
- [ ] **Step 3: 应用层删在途单互斥检查**(`app/repositories/wallet.py` · `create_withdraw`)
删掉在途单预检查(它紧邻档位闸注释之前):
old:
```python
active_order_id = db.execute(
select(WithdrawOrder.id).where(
WithdrawOrder.user_id == user_id,
WithdrawOrder.status.in_(_WITHDRAW_ACTIVE_STATUSES),
).limit(1)
).scalar_one_or_none()
if active_order_id is not None:
raise WithdrawTooFrequentError
# 福利页档位闸(7-9):coin_cash 只能提预设档位,且该档今日可提(服务端权威口径,防绕过
```
new:
```python
# 福利页档位闸(7-9):coin_cash 只能提预设档位,且该档今日可提(服务端权威口径,防绕过
```
- [ ] **Step 4: 应用层给 IntegrityError 兜底瘦身**(同函数末尾 commit 处)
索引移除后不会再因在途单触发唯一冲突,只保留 `out_bill_no` 幂等重试分支:
old:
```python
except IntegrityError:
db.rollback()
existing = db.execute(
select(WithdrawOrder).where(
WithdrawOrder.out_bill_no == out_bill_no, WithdrawOrder.user_id == user_id
)
).scalar_one_or_none()
if existing is not None:
return existing
active_order_id = db.execute(
select(WithdrawOrder.id).where(
WithdrawOrder.user_id == user_id,
WithdrawOrder.status.in_(_WITHDRAW_ACTIVE_STATUSES),
).limit(1)
).scalar_one_or_none()
if active_order_id is not None:
raise WithdrawTooFrequentError from None
raise
```
new:
```python
except IntegrityError:
db.rollback()
# 唯一冲突只可能来自 out_bill_no 幂等键并发重试:原样返回既有单;否则未知冲突,上抛。
existing = db.execute(
select(WithdrawOrder).where(
WithdrawOrder.out_bill_no == out_bill_no, WithdrawOrder.user_id == user_id
)
).scalar_one_or_none()
if existing is not None:
return existing
raise
```
- [ ] **Step 5: 删死常量 `_WITHDRAW_ACTIVE_STATUSES`**(`app/repositories/wallet.py` 模块顶部)
只删常量行,保留其后属于 `_NEWBIE_TIER_HELD_STATUSES` 的注释:
old:
```python
_WITHDRAW_ACTIVE_STATUSES = {"reviewing", "pending"}
# 占用新人档「一次性」资格的提现状态:进行中(reviewing/pending)或成功打款(success)。
```
new:
```python
# 占用新人档「一次性」资格的提现状态:进行中(reviewing/pending)或成功打款(success)。
```
- [ ] **Step 6: 删死异常类 `WithdrawTooFrequentError`**(`app/repositories/wallet.py`)
old:
```python
class WithdrawTooFrequentError(Exception):
"""提现申请过于频繁,或已有未完成提现单。"""
class WithdrawTierUnavailableError(Exception):
```
new:
```python
class WithdrawTierUnavailableError(Exception):
```
- [ ] **Step 7: 删端点内失效的 409 处理**(`app/api/v1/wallet.py` · `withdraw`)
old:
```python
except crud_wallet.WechatNotBoundError as e:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="请先绑定微信") from e
except crud_wallet.WithdrawTooFrequentError as e:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail="已有提现申请正在审核或打款中,请处理完成后再申请",
) from e
except crud_wallet.WithdrawTierUnavailableError as e:
```
new:
```python
except crud_wallet.WechatNotBoundError as e:
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="请先绑定微信") from e
except crud_wallet.WithdrawTierUnavailableError as e:
```
- [ ] **Step 8: 删模型内分区唯一索引**(`app/models/wallet.py` · `WithdrawOrder`)
old:
```python
__tablename__ = "withdraw_order"
__table_args__ = (
Index(
"ux_withdraw_order_user_active",
"user_id",
unique=True,
sqlite_where=text("status IN ('reviewing', 'pending')"),
postgresql_where=text("status IN ('reviewing', 'pending')"),
),
)
```
new:
```python
__tablename__ = "withdraw_order"
```
- [ ] **Step 9: 运行新测试,确认通过**
Run: `pytest tests/test_withdraw.py::test_withdraw_multiple_in_flight_allowed tests/test_withdraw.py::test_withdraw_second_blocked_only_by_insufficient_cash -q`
Expected: 2 passed。
- [ ] **Step 10: 建 Alembic 迁移(真实库删索引)**
创建 `alembic/versions/drop_withdraw_active_unique_index.py`:
```python
"""drop withdraw active-order partial unique index (allow multiple in-flight withdrawals)
Revision ID: drop_withdraw_active_unique_index
Revises: d8dd2106e438
Create Date: 2026-07-24 00:00:00.000000
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
# revision identifiers, used by Alembic.
revision: str = "drop_withdraw_active_unique_index"
down_revision: Union[str, Sequence[str], None] = "d8dd2106e438"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
# 取消「同一用户同一时刻仅一笔在途提现」:允许 reviewing/pending 并存。
# 仅删本索引;姊妹索引 ux_cash_transaction_withdraw_refund_ref(退款幂等)保持不动。
op.drop_index("ux_withdraw_order_user_active", table_name="withdraw_order")
def downgrade() -> None:
# 回滚重建分区唯一索引。注意:若届时某用户已有 ≥2 张在途单,重建会因唯一冲突失败——
# 属预期的回滚代价(取消限制后本就允许多单),需先人工收敛在途单再回滚。
op.create_index(
"ux_withdraw_order_user_active",
"withdraw_order",
["user_id"],
unique=True,
sqlite_where=sa.text("status IN ('reviewing', 'pending')"),
postgresql_where=sa.text("status IN ('reviewing', 'pending')"),
)
```
- [ ] **Step 11: 验证迁移 upgrade + 回环 downgrade/upgrade**
Run: `alembic upgrade head`
Expected: 输出应用 `drop_withdraw_active_unique_index`,无报错。
Run: `alembic downgrade -1 && alembic upgrade head`
Expected: downgrade 重建索引、upgrade 再次删除,均成功(dev 库每用户在途单 ≤1,不会触发唯一冲突)。
- [ ] **Step 12: 全量测试 + Lint**
Run: `pytest -q`
Expected: 全绿(既有 `test_withdraw_idempotent_same_bill_no``tests/test_withdraw_tiers.py``tests/test_invite_cash_withdraw.py` 均不受影响)。
Run: `ruff check .`
Expected: 无新增告警(死常量/异常已连同引用一并删除)。
- [ ] **Step 13: 提交**
```bash
git add tests/test_withdraw.py app/repositories/wallet.py app/api/v1/wallet.py app/models/wallet.py alembic/versions/drop_withdraw_active_unique_index.py
git commit -m "feat(withdraw): 允许在途提现时继续提交新申请
取消「同一用户同时仅一笔在途提现」限制:删应用层在途单检查 +
删 DB 分区唯一索引 ux_withdraw_order_user_active + 清理死代码
(_WITHDRAW_ACTIVE_STATUSES / WithdrawTooFrequentError)。
先扣款、coin_cash 每日档位次数、out_bill_no 幂等等既有约束不变。
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>"
```
---
## 自审(spec 覆盖核对)
- **R1 已有 reviewing/pending 可再提** → Step 3/4(应用层)+ Step 8/10(模型 + 迁移);`test_withdraw_multiple_in_flight_allowed` 证明。✓
- **R2 不加硬上限** → 无新增限制;`coin_cash` 天然封顶由既有档位次数(`tests/test_withdraw_tiers.py`)保障,不改。✓
- **R3 幂等/限额/退款/对账不回退**`test_withdraw_idempotent_same_bill_no` 与档位/邀请现金测试留绿(Step 12);解绑退款、对账按单号维度,未触碰。✓
- **spec §4 四处改动** → Step 3/4/5/6(repo)、Step 7(api)、Step 8(model)、Step 10(迁移)一一对应。✓
- **spec §6 客户端注意点** → 属跨仓 App 事项,文档已记,本计划无对应代码任务(有意为之)。✓
- **占位符扫描**:迁移 `revision` / `down_revision`(`d8dd2106e438`)均为具体值,无 TBD。✓
- **命名一致性**:`WithdrawTooFrequentError`/`_WITHDRAW_ACTIVE_STATUSES` 的删除点与引用点全覆盖(全仓仅这 4 处);`ux_withdraw_order_user_active` 在模型/迁移中拼写一致。✓
@@ -0,0 +1,119 @@
# 允许在途提现时继续提交新的提现申请 设计
- **日期**:2026-07-24
- **状态**:Draft — 待评审
- **所属**:app-server(`app/`),含一处 Alembic 迁移;另有一条 Android/客户端注意点(非本次后端工作)
- **一句话**:取消「同一用户同一时刻只能有一笔在途提现」的限制 —— 已有 `reviewing`(待审核)或 `pending`(打款在途)提现单时,允许再次发起新的提现申请。
---
## 1. 背景与目标
现状:用户发起提现后,在管理员审核通过并打款完成之前(`reviewing` / `pending`),**无法再发起第二笔提现**,会收到 409「已有提现申请正在审核或打款中,请处理完成后再申请」。人工审核有延迟时,用户被卡住、体验差。
**目标**:放开该限制,已有在途提现单时仍可继续提交新的提现申请。
### 非目标(本期不做)
- 不改「先扣款」模型(提现建单即原子扣现金,天然防超提)。
- 不改 `coin_cash` 的每日档位次数限制(仍是独立的限额闸)。
- 不加任何新的并发/次数硬上限(见 §3 决策)。
- 不改 admin 审核/打款/对账逻辑(其全部按单号维度操作,天然支持一人多单)。
---
## 2. 需求
| # | 需求 | 落地 |
|---|---|---|
| R1 | 已有 `reviewing``pending` 单时可再次发起提现 | 删除应用层在途单互斥检查 + 删除 DB 分区唯一索引(§4) |
| R2 | 不引入新的并发上限 | 仅靠既有约束:先扣款余额 + `coin_cash` 每日档位次数(§5) |
| R3 | 既有幂等/限额/退款/对账行为不回退 | 保留 `out_bill_no` 幂等、档位闸、解绑退款、对账(§5) |
---
## 3. 决策记录(来自评审问答)
| 决策点 | 结论 | 理由 |
|---|---|---|
| **放行范围** | `reviewing``pending` **两种在途状态都放行**,彻底取消「同时仅一单」 | 需求即"存在在途提现时可继续提交";人工审核延迟不应卡住用户 |
| **并发上限** | **不加硬上限** | 现金**先扣款**→每笔各需自己的余额,不会超提;`coin_cash` 每日档位次数已是天然上限;`invite_cash` 仅受余额约束,可接受 |
| **实现方式** | **彻底移除**(删检查 + 删索引 + 清理死代码),非配置开关 | 决策明确且无回滚诉求,YAGNI;避免留半死代码与多余配置项 |
| **已失效异常/文案** | 删除 `WithdrawTooFrequentError` 及端点的 409 处理 | 全仓仅这 4 处引用(§4),移除后无残留 |
### 已知取舍(可接受)
- **多笔 `pending` × 未开免确认**:每笔 `pending` 会各自返回一个微信确认页 `package_info`。用户若**未开启免确认**,可能同时存在多笔待确认。属客户端交互问题(§6),后端行为正确;开启免确认后直接到账、无此问题。
- **回滚风险**:迁移 `downgrade` 会重建分区唯一索引;若届时某用户已有 ≥2 张在途单,重建会失败 —— 属预期的回滚代价,在迁移里注释说明。
---
## 4. 现状机制与改动点
「同一时刻仅一笔在途提现」由**两道闸**共同强制,均以 `status IN ('reviewing','pending')` 为口径:
1. **应用层** — [`app/repositories/wallet.py:758-765`](../../../app/repositories/wallet.py) `create_withdraw` 内:查在途单 → `raise WithdrawTooFrequentError` → 端点转 409。
2. **数据库层** — [`app/models/wallet.py:99-107`](../../../app/models/wallet.py) 的**分区唯一索引** `ux_withdraw_order_user_active`(`user_id WHERE status IN ('reviewing','pending')`)。这是硬约束,`create_withdraw``IntegrityError` 兜底([`:814-833`](../../../app/repositories/wallet.py))即捕获它。
> 提现单状态机:`reviewing`(待审核,建单即扣现金)→`pending`(审核通过、打款在途)→`success`/`failed`;或 `reviewing``rejected`(已退款)。
### 改动清单(4 处代码 + 1 个迁移)
**① `app/repositories/wallet.py` · `create_withdraw`**
- 删除在途单互斥检查([`:758-765`](../../../app/repositories/wallet.py)):`select WithdrawOrder.id WHERE status IN (_WITHDRAW_ACTIVE_STATUSES)``raise WithdrawTooFrequentError`
- `IntegrityError` 兜底([`:814-833`](../../../app/repositories/wallet.py)):**保留** `out_bill_no` 幂等重试分支([`:816-824`](../../../app/repositories/wallet.py));**删除**其中的在途单二次检查分支([`:825-832`](../../../app/repositories/wallet.py))(索引移除后不会再因在途单触发 `IntegrityError`);末尾其余情况原样 `raise`(仅剩 `out_bill_no` 唯一冲突等,正常不该出现)。
- 删除模块常量 `_WITHDRAW_ACTIVE_STATUSES`([`:38`](../../../app/repositories/wallet.py))与异常类 `WithdrawTooFrequentError`([`:72-73`](../../../app/repositories/wallet.py))。`_NEWBIE_TIER_HELD_STATUSES`(档位资格判定,含 `success`)与本改动无关,**保留**。
**② `app/models/wallet.py` · `WithdrawOrder`**
- 删除 `__table_args__` 中的分区唯一索引 `ux_withdraw_order_user_active`([`:99-107`](../../../app/models/wallet.py))。该表 `__table_args__` 仅此一项,整块移除。
**③ `app/api/v1/wallet.py` · `withdraw` 端点**
- 删除 `except crud_wallet.WithdrawTooFrequentError`([`:225-229`](../../../app/api/v1/wallet.py))这段已失效的 409 处理。
**④ 新增 Alembic 迁移** `alembic/versions/<...>_drop_withdraw_active_unique_index.py`
- `down_revision` = 当前 head(实现时确定)。
- `upgrade`:`op.drop_index("ux_withdraw_order_user_active", table_name="withdraw_order")`
- `downgrade`:`op.create_index("ux_withdraw_order_user_active", "withdraw_order", ["user_id"], unique=True, sqlite_where=text("status IN ('reviewing','pending')"), postgresql_where=text("status IN ('reviewing','pending')"))`,并注释"若已有用户存在多张在途单则重建失败,属预期回滚代价"。
- 索引的 drop/create 为具名操作,SQLite/PG 均无需 `render_as_batch` 重建表。
---
## 5. 天然保持不变(无需改动)的约束
| 约束 | 为何仍成立 |
|---|---|
| **不会超提** | 建单即原子扣现金(`_try_deduct_cash`,余额不足影响 0 行 → `InsufficientCashError`);每笔并行单各需自己的余额 |
| **`coin_cash` 每日档位次数** | `withdraw_tier_states` 的档位闸独立于在途单检查:常规档按**当天发起即计入(任意状态,含被拒/失败,按 `created_at` 北京日)**、每档每日限次(0.5×3 / 10×1 / 20×1)且当天只选一档;新人档(0.1/0.3)按 `reviewing/pending/success` 一次性占用。故即便并行,`coin_cash` 单日在途仍被档位天然封顶(至多 3 笔 0.5) |
| **`out_bill_no` 幂等** | 幂等分支在被删检查之前,同号重试仍原样返回旧单、不重复扣款 |
| **解绑退款** | `refund_reviewing_withdraws_on_unbind``for` 遍历该用户**所有** `reviewing` 单,天然支持多单 |
| **admin 审核 / 打款 / 对账** | `approve_withdraw`/`reject_withdraw`/`refresh_withdraw_status`/对账全按 `out_bill_no` 单号维度,不假设一人一单 |
---
## 6. 客户端/App 团队注意点(跨仓,非本次后端工作)
> - 允许多笔并行后,每次提交都要生成**新的 `out_bill_no`**(复用旧号会命中幂等、返回旧单)。
> - 用户**未开启免确认**时,多笔 `pending` 会各自返回一个微信确认页 `package_info`,App 需能处理/串行多笔待确认。开启免确认后直接到账、无此问题。
---
## 7. 测试计划(`tests/test_withdraw.py`,沿用现有 helper)
- **新增 · 多笔在途放行**:同一用户 seed ≥1 元现金,用**两个不同 `out_bill_no`** 各提 0.5 元(在 0.5 元档 3 次/天限额内)→ 两次均 200 且 `status=reviewing`;`/withdraw-orders` 返回 2 条;余额正确扣两次(1 元 → 0)。
- **新增 · 第二笔余额不足**:seed 0.5 元,连提两笔 0.5 元 → 第一笔 200、第二笔 409(`InsufficientCashError`),验证每笔各需自己的余额。
- **回归 · 幂等**:`test_withdraw_idempotent_same_bill_no`(同 `out_bill_no` → 同一单、只扣一次)仍绿。
- **回归 · 档位限额**:`tests/test_withdraw_tiers.py`(0.5 元日 3 次、第 4 次 409;跨档互斥)不受影响仍绿。
- 沿用 `tests/conftest.py`(临时 SQLite、`RATE_LIMIT_ENABLED=false`);wxpay 网络调用全部 monkeypatch。
---
## 附:涉及文件清单
**改动**
- `app/repositories/wallet.py` — 删在途单检查 + `IntegrityError` 兜底瘦身 + 删 `_WITHDRAW_ACTIVE_STATUSES` / `WithdrawTooFrequentError`
- `app/models/wallet.py` — 删分区唯一索引 `ux_withdraw_order_user_active`
- `app/api/v1/wallet.py` — 删 `WithdrawTooFrequentError` 的 409 处理
- `tests/test_withdraw.py` — 新增多笔在途放行 / 第二笔余额不足用例
**新增**
- `alembic/versions/<...>_drop_withdraw_active_unique_index.py` — drop 分区唯一索引(downgrade 重建)
@@ -0,0 +1,182 @@
# 阿里云短信验证服务 — 设计方案
- 日期:2026-07-25
- 状态:已定稿,待实现
- 范围:新增阿里云 dypns(号码认证服务)短信验证码 provider,与现有极光短信可切换
## 1. 背景与目标
现有短信验证码服务 `app/integrations/sms.py`:本服务**本地生成**验证码、存**进程内存**、极光 REST 仅负责下发;`verify_code()` 比对本地内存(一次性 + 单码失败 `SMS_MAX_VERIFY_ATTEMPTS` 次即作废)。docstring 已标注"内存存码、多 worker 不共享"为技术债。
阿里云文档(`docs/integrations/aliyun/`)为 **号码认证服务 dypns**`SendSmsVerifyCode` + `CheckSmsVerifyCode`:该产品由**阿里云生成并校验**验证码(`{"code":"##code##"}` 模式),核验免费。
目标:接入阿里云该套接口作为一个新的短信 provider,可与极光切换。
## 2. 关键决策(已确认)
1. **验证码模式 = Mode A(阿里云托管码)**:发码用 `SendSmsVerifyCode` + `##code##` 占位符,阿里云生成/存储/下发;校验用 `CheckSmsVerifyCode`,阿里云返回 `PASS/UNKNOWN`。本服务不再本地生成/存储验证码。
2. **可切换 Provider**:新增 `SMS_PROVIDER=jiguang|aliyun` 开关,`send_code/verify_code` 按 provider 分派;保留极光作回退(短信=花钱+登录关键路径,灰度上线/融合认证未开通时可秒切回极光)。
3. **官方 SDK**:调阿里云 dypns 用 `alibabacloud_dypnsapi20170525`,签名/加签由 SDK 处理。
4. **防爆破与极光一致**(排查一致性):阿里云路径**保留**与极光相同的"单码失败 N 次即作废"本地计数,而非改用 API 层频控,避免两 provider 行为不一致导致排查困惑。
## 3. 模块结构(`sms.py` 单文件升级为 provider 包)
```
app/integrations/sms/
__init__.py # 公开 API + 分派器:send_code / verify_code / SmsError
# - SMS_MOCK=true 短路(不碰任何 provider)
# - 按 settings.SMS_PROVIDER 选 jiguang / aliyun
base.py # SmsError(沿用现定义)+ Provider 协议(send_code/verify_code 签名约定)
jiguang.py # 现有自管码逻辑原样迁入(内存存码/冷却/一次性/防爆破 全保留,行为零改动)
aliyun.py # 新增:SendSmsVerifyCode 发码 + CheckSmsVerifyCode 校验 + 本地失败计数
```
- `__init__.py` 继续 re-export `SmsError / send_code / verify_code`,故 `app/api/v1/auth.py:37`
`from app.integrations.sms import SmsError, send_code, verify_code` **导入不变**
- 纯增量重构:极光逻辑整体迁入 `jiguang.py`,对外行为零变化。
## 4. 数据流 — 阿里云 providerMode A
### 4.1 发码 `aliyun.send_code(phone) -> int`
1. 校验 `settings.aliyun_sms_configured`(缺 AK/SignName/TemplateCode → `SmsError(503)`)。
2. 调 `SendSmsVerifyCode`
- `PhoneNumber=phone`
- `SignName=ALIYUN_SMS_SIGN_NAME``TemplateCode=ALIYUN_SMS_TEMPLATE_CODE`
- `TemplateParam = json({"code":"##code##","min": str(ALIYUN_SMS_VALID_TIME_SEC//60)})`
- `CodeLength=ALIYUN_SMS_CODE_LENGTH``ValidTime=ALIYUN_SMS_VALID_TIME_SEC``Interval=ALIYUN_SMS_INTERVAL_SEC`
- `SchemeName=ALIYUN_SMS_SCHEME_NAME`(可空)
3. 成功(`body.Success and body.Code=="OK"`)→ **清本地失败计数**(新码=新预算)→ 返回 `ALIYUN_SMS_INTERVAL_SEC` 作客户端冷却秒数。
4. 失败 → 按 §6 错误码映射抛 `SmsError`
### 4.2 校验 `aliyun.verify_code(phone, code) -> bool`
1. **本地失败计数**`attempts >= SMS_MAX_VERIFY_ATTEMPTS` → 直接 `False`(码已作废,不调阿里云)。
2. 调 `CheckSmsVerifyCode(PhoneNumber, VerifyCode=code, SchemeName)`
3. `body.Model.VerifyResult`
- `"PASS"` → 清计数,返回 `True`(一次性)。
- `"UNKNOWN"``attempts += 1`,返回 `False`(码错/过期)。
4. 网络错误 / 接口非 `OK` → 抛 `SmsError(503)`**不静默返回 False**,区分"阿里云挂了"与"码错了";网络错误不计入 attempts)。
本服务**不存验证码**,仅存一个 per-phone 失败计数(见 §7)。
## 5. 分派器 & mock`__init__.py`
```
send_code(phone):
if settings.SMS_MOCK: # 短路:不碰 provider(测试/开发)
log placeholder code; return cooldown
return _provider().send_code(phone)
verify_code(phone, code):
if settings.SMS_MOCK: # 放行任意 N 位数字(沿用现 mock 语义)
return len(code)==SMS_CODE_LENGTH and code.isdigit()
return _provider().verify_code(phone, code)
_provider(): jiguang if settings.SMS_PROVIDER=="jiguang" else aliyun
```
- mock 语义提到分派层、provider 无关 → 现有 28 个测试文件(conftest 设 `SMS_MOCK=true`)全部零改动通过。
## 6. 错误映射
### 发码(阿里云错误码 → SmsError.status_code
| 阿里云码 | HTTP | 说明 |
|---|---|---|
| `MOBILE_NUMBER_ILLEGAL` | 400 | 手机号格式错误 |
| `BUSINESS_LIMIT_CONTROL` | 429 | 号码天级流控 |
| `FREQUENCY_FAIL` | 429 | 频控(`Interval` 命中) |
| `FUNCTION_NOT_OPENED` | 503 | 融合认证未开通(**critical 日志**,需运维开通) |
| `INVALID_PARAMETERS` | 503 | 参数错误(配置/模板问题,**critical 日志** |
| 其他非 OK / `Success=false` / 网络错误 | 503 | 供应商不可用 |
### 校验
- `PASS` → True`UNKNOWN` → False;接口异常/网络错误 → `SmsError(503)`
## 7. 防爆破 / 频控分工
| 机制 | 极光(Mode B | 阿里云(Mode A |
|---|---|---|
| 验证码存储 | 本地内存 | **阿里云托管**(消除多 worker 存码债) |
| 单号发送冷却 | 本地 `_last_sent` 60s | **交给阿里云 `Interval`**(无本地状态),命中→429 |
| 单设备+IP 频控 | API 层 5/时、20/天 | 同左,**不变** |
| **防爆破(单码失败 N 次即作废)** | 本地 `_CodeRecord.attempts` | **本地 per-phone 计数**,与极光同语义(§4.2 |
- 阿里云路径的**唯一本地状态** = per-phone 失败计数 `dict[phone,int]` + `Lock` + GC(仿极光 `_gc`)。
- 多 worker 降级:失败计数按 worker 各计,effective 上限 = N×workers;与极光现状**同级**(属刻意保留的一致性),且 API 层登录频控(`sms-login-device` 设备+IP 5/时)提供硬兜底。
- 计数复位:`send_code` 成功清计数、`verify` PASS 清计数(新码/验过即新预算)。
- API 层设备频控与测试账号短路(`app/core/test_account.py`**完全不动**。
## 8. 配置项(`app/core/config.py` 新增)
```python
SMS_PROVIDER: str = "jiguang" # jiguang | aliyun;默认极光(保持现状,上线后切 aliyun)
# --- 阿里云 dypns 号码认证 ---
ALIYUN_SMS_ACCESS_KEY_ID: str = ""
ALIYUN_SMS_ACCESS_KEY_SECRET: str = ""
ALIYUN_SMS_SIGN_NAME: str = "" # 系统赠送签名(自定义签名下发易失败)
ALIYUN_SMS_TEMPLATE_CODE: str = "" # 赠送模板 CODE(须与赠送签名搭配)
ALIYUN_SMS_SCHEME_NAME: str = "" # 方案名(可空=默认方案);send/check 必须一致 → 单一来源
ALIYUN_SMS_ENDPOINT: str = "dypnsapi.aliyuncs.com"
ALIYUN_SMS_CODE_LENGTH: int = 6 # CodeLength 4~8
ALIYUN_SMS_VALID_TIME_SEC: int = 300 # ValidTime;短信内 min 文案 = //60
ALIYUN_SMS_INTERVAL_SEC: int = 60 # Interval 单号发送频控
```
- 新增属性 `aliyun_sms_configured`(仿 `mt_cps_configured`):AK_ID/AK_SECRET/SignName/TemplateCode 齐全才为真;`SMS_PROVIDER=aliyun` 但未配 → `send_code``SmsError(503)`
- 复用现有 `SMS_MOCK``SMS_CODE_LENGTH`mock 校验位数)、`SMS_MAX_VERIFY_ATTEMPTS`(防爆破上限,两 provider 共用)。
### 配置敏感点
- `TemplateParam` 变量名(`code`/`min`)须与控制台所选**赠送模板**一致。融合认证验证码模板通常即 `code`+`min`,按此硬编码并加注释;若模板变量名不同,改 `aliyun.py` 该处即可。
- `SchemeName` 在 send 与 check 必须一致,故用**单一** `ALIYUN_SMS_SCHEME_NAME` 供两处,避免不匹配(CheckSmsVerifyCode 文档明确警告)。
## 9. auth.py 改动(最小)
`verify_code` 现在可能抛 `SmsError`(阿里云降级 503)。两处调用点各包一层 `try/except SmsError → HTTPException(e.status_code)`,与 `send_code` 现有写法一致:
- `app/api/v1/auth.py` `sms_login`(约 L185
- `app/api/v1/auth.py` `wechat_bind_phone_sms`(约 L325
`send_code` 调用点已 try/except `SmsError`,无需改。
## 10. 依赖 & SDK
- `pyproject.toml``alibabacloud_dypnsapi20170525`(连带 `alibabacloud-tea-openapi` 等)。
- SDK 同步阻塞调用 → 与现有 sync 端点 + sync httpx 风格一致(FastAPI 跑 threadpool,无碍)。
- `aliyun.py` 内**惰性 import SDK + 惰性建 client**(仿 `wxpay` 惰性加载证书):`SMS_PROVIDER=jiguang` 时不加载 alibabacloud,启动保持精简。
- SDK 调用形态(实现时按实际包名/字段核对):
```python
from alibabacloud_dypnsapi20170525.client import Client
from alibabacloud_dypnsapi20170525 import models as dypns_models
from alibabacloud_tea_openapi import models as open_api_models
cfg = open_api_models.Config(access_key_id=..., access_key_secret=...)
cfg.endpoint = settings.ALIYUN_SMS_ENDPOINT
client = Client(cfg)
resp = client.send_sms_verify_code(dypns_models.SendSmsVerifyCodeRequest(...))
# resp.body.code / resp.body.success / resp.body.model.verify_code
resp = client.check_sms_verify_code(dypns_models.CheckSmsVerifyCodeRequest(...))
# resp.body.model.verify_result == "PASS"
```
## 11. 测试
- 现有测试:`SMS_MOCK=true` → 分派器短路,全绿不变。
- 新增 `tests/test_sms_aliyun.py`monkeypatch SDK client,不发真网络):
1. 发码成功 → 返回 cooldown、清计数。
2. 各错误码 → 对应 `SmsError.status_code`400/429/503)。
3. 校验 `PASS`→True(清计数)/ `UNKNOWN`→False(计数 +1)/ 接口异常→`SmsError(503)`
4. 失败计数达 `SMS_MAX_VERIFY_ATTEMPTS` → 直接 False,不再调阿里云。
5. `send_code` 成功复位计数。
- 新增分派测试:`SMS_PROVIDER` 切换选中正确 provider`SMS_MOCK` 优先于 provider。
## 12. YAGNI(明确不做)
- ❌ 不做 Redis/DB 存码(Mode A 无需;极光路径内存债维持现状,非本次范围)。
- ❌ 不改极光任何行为、不动 API 层频控/测试账号逻辑。
- ❌ 不做多签名/多模板轮换(单签名单模板足够)。
- ❌ 不把失败计数持久化/跨进程(刻意保留与极光同级的本地态)。
## 13. 验收标准
- `SMS_PROVIDER=aliyun` 且配置齐全时:`/sms/send``SendSmsVerifyCode``/sms/login``CheckSmsVerifyCode`,真机可收码并登录。
- `SMS_PROVIDER=jiguang`(默认):行为与当前完全一致。
- `SMS_MOCK=true`:任意 N 位数字通过,不发真短信。
- 阿里云接口异常时:`/sms/login` 返回 503(非 400),日志可区分。
- `ruff check .` 通过;新增/现有 `pytest` 全绿。
+3
View File
@@ -29,6 +29,9 @@ dependencies = [
# HTTP 客户端 (调极光 REST)
"httpx>=0.27.0",
# 阿里云号码认证(dypns)短信验证码 provider(SMS_PROVIDER=aliyun 时用;签名由 SDK 处理)
"alibabacloud_dypnsapi20170525>=2.0.0",
# multipart form (FastAPI 表单上传依赖)
"python-multipart>=0.0.9",
+274
View File
@@ -0,0 +1,274 @@
"""为本地“风控监控”页面灌入一组可重复的完整演示数据。
仅允许在 APP_ENV=dev 运行脚本会重建 ``risk-demo-*`` 前缀的数据并创建本地
后台账号 ``risk_demo / RiskDemo123!``方便端到端和视觉验收
"""
from __future__ import annotations
from datetime import timedelta
from sqlalchemy import delete, select
from app.admin.repositories import admin_user as admin_repo
from app.core.config import settings
from app.core.security import hash_password
from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord
from app.models.risk import BehaviorEvent, RiskIncident, SubjectRestriction
from app.repositories import risk as risk_repo
from app.repositories import user as user_repo
DEMO_PREFIX = "risk-demo-"
ADMIN_USERNAME = "risk_demo"
ADMIN_PASSWORD = "RiskDemo123!"
def _reset_demo(db) -> None:
demo_user_ids = list(
db.scalars(
select(ComparisonRecord.user_id)
.where(
ComparisonRecord.trace_id.like(f"{DEMO_PREFIX}%"),
ComparisonRecord.user_id.is_not(None),
)
.distinct()
).all()
)
demo_incident_ids = list(
db.scalars(
select(RiskIncident.id).where(
(RiskIncident.subject_id.like(f"{DEMO_PREFIX}%"))
| (
(RiskIncident.rule_code == risk_repo.RULE_COMPARE_DAILY)
& (RiskIncident.subject_id.in_([str(value) for value in demo_user_ids]))
)
)
).all()
)
if demo_incident_ids:
db.execute(
delete(SubjectRestriction).where(
SubjectRestriction.incident_id.in_(demo_incident_ids)
)
)
db.execute(
delete(SubjectRestriction).where(
SubjectRestriction.subject_id.like(f"{DEMO_PREFIX}%")
)
)
db.execute(
delete(RiskIncident).where(
(RiskIncident.subject_id.like(f"{DEMO_PREFIX}%"))
| (
(RiskIncident.rule_code == risk_repo.RULE_COMPARE_DAILY)
& (RiskIncident.subject_id.in_([str(value) for value in demo_user_ids]))
)
)
)
db.execute(
delete(BehaviorEvent).where(
BehaviorEvent.subject_id.like(f"{DEMO_PREFIX}%")
)
)
db.execute(
delete(ComparisonRecord).where(
ComparisonRecord.trace_id.like(f"{DEMO_PREFIX}%")
)
)
db.commit()
def _ensure_admin(db) -> None:
admin = admin_repo.get_by_username(db, ADMIN_USERNAME)
if admin is None:
admin_repo.create_admin(
db,
username=ADMIN_USERNAME,
password=ADMIN_PASSWORD,
role="super_admin",
)
return
admin.password_hash = hash_password(ADMIN_PASSWORD)
admin.role = "super_admin"
admin.status = "active"
db.commit()
def _seed_sms(db, now) -> None:
alerted = [
("risk-demo-sms-oppo-a5", "OPPO A5", "13812343001", 9),
("risk-demo-sms-xiaomi14", "Xiaomi 14", "13812343002", 7),
("risk-demo-sms-iphone15", "iPhone 15", "13812343003", 5),
]
events: list[BehaviorEvent] = []
for device_id, model, phone, count in alerted:
for index in range(count):
events.append(
BehaviorEvent(
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=device_id,
device_id=device_id,
device_model=model,
phone=phone if index == 0 else f"138{index:08d}"[-11:],
client_ip="127.0.0.1",
outcome="success",
details={"demo": True},
occurred_at=now + timedelta(seconds=index),
)
)
# 其余成功下发分散到未达阈值的设备,每台最多 4 条,确保报警设备仍严格为 3 台。
remaining = 1284 - sum(row[3] for row in alerted)
for index in range(remaining):
device_no = index // 4
device_id = f"risk-demo-sms-normal-{device_no:04d}"
events.append(
BehaviorEvent(
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=device_id,
device_id=device_id,
device_model="演示普通设备",
phone=f"137{index % 100_000_000:08d}",
client_ip="127.0.0.1",
outcome="success",
details={"demo": True},
occurred_at=now + timedelta(seconds=index % 1800),
)
)
db.add_all(events)
db.flush()
for device_id, _, phone, _ in alerted:
risk_repo.evaluate_behavior_rule(
db,
rule_code=risk_repo.RULE_SMS_HOURLY,
subject_id=device_id,
at=now,
)
user = user_repo.upsert_user_for_login(
db, phone=phone, register_channel="sms"
)
db.add(
BehaviorEvent(
event_type=risk_repo.EVENT_SMS_LOGIN,
subject_type="device",
subject_id=device_id,
user_id=user.id,
device_id=device_id,
phone=phone,
outcome="success",
details={"demo": True},
occurred_at=now + timedelta(minutes=2),
)
)
db.commit()
def _seed_oneclick(db, now) -> None:
devices = [
("risk-demo-oneclick-oppo-reno", "OPPO Reno", "13912343001", 482),
("risk-demo-oneclick-vivo-y36", "vivo Y36", "13912343002", 481),
]
for device_id, model, phone, count in devices:
user = user_repo.upsert_user_for_login(
db, phone=phone, register_channel="jverify"
)
db.add_all(
[
BehaviorEvent(
event_type=risk_repo.EVENT_ONECLICK_LOGIN,
subject_type="device",
subject_id=device_id,
user_id=user.id,
device_id=device_id,
device_model=model,
phone=phone,
client_ip="127.0.0.1",
outcome="failed" if index % 11 == 0 else "success",
reason="运营商校验失败" if index % 11 == 0 else None,
details={"demo": True},
occurred_at=now + timedelta(seconds=index * 2),
)
for index in range(count)
]
)
db.flush()
risk_repo.evaluate_behavior_rule(
db,
rule_code=risk_repo.RULE_ONECLICK_DAILY,
subject_id=device_id,
at=now,
)
db.commit()
def _seed_compare(db, now) -> None:
local_now = now.astimezone(risk_repo.CN_TZ).replace(tzinfo=None)
accounts = [
("13612343281", "risk-demo-compare-device-a", 1053),
("13511106208", "risk-demo-compare-device-b", 1052),
]
for phone, device_id, count in accounts:
user = user_repo.upsert_user_for_login(
db, phone=phone, register_channel="sms"
)
rows = []
for index in range(count):
rows.append(
ComparisonRecord(
user_id=user.id,
device_id=device_id,
trace_id=f"{DEMO_PREFIX}compare-{user.id}-{index}",
store_name=("春熙路小吃店", "科技园轻食", "万达广场烤肉")[
index % 3
],
product_names=("招牌套餐", "鸡胸沙拉", "双人烤肉")[index % 3],
status="failed" if index % 17 == 0 else "success",
information="演示比价记录",
items=[{"name": "演示菜品", "qty": 1}],
comparison_results=[
{
"platform_id": "meituan",
"platform_name": "美团",
"price": 23.8,
},
{
"platform_id": "taobao",
"platform_name": "淘宝",
"price": 21.5,
},
{
"platform_id": "jd",
"platform_name": "京东",
"price": 22.2,
},
],
skipped_dish_names=[],
saved_amount_cents=230,
raw_payload={"demo": True},
created_at=local_now + timedelta(seconds=index * 2),
)
)
db.add_all(rows)
db.commit()
risk_repo.sync_compare_incident(db, user_id=user.id, at=local_now)
def main() -> None:
if settings.APP_ENV != "dev":
raise SystemExit("拒绝执行:风控演示数据脚本仅允许 APP_ENV=dev")
now = risk_repo.utcnow().replace(minute=10, second=0, microsecond=0)
with SessionLocal() as db:
_reset_demo(db)
_ensure_admin(db)
_seed_sms(db, now)
_seed_oneclick(db, now)
_seed_compare(db, now)
print("风控监控演示数据已重建")
print(f"后台账号:{ADMIN_USERNAME}")
print(f"后台密码:{ADMIN_PASSWORD}")
print("期望卡片:短信 3 / 1284;一键登录 2 / 963;比价 2 / 2105")
if __name__ == "__main__":
main()
+13 -5
View File
@@ -68,10 +68,17 @@ def test_monitoring_audit_catalog_and_api_permissions(
).json()
monitoring = next(group for group in catalog if group["group"] == "监控审计")
assert [page["key"] for page in monitoring["pages"]] == [
"device-liveness", "analytics-health", "event-logs", "audit-logs",
"risk-monitor",
"device-liveness",
"analytics-health",
"event-logs",
"audit-logs",
]
# 运营默认只能查设备存活,不能绕过导航直调技术/审计接口。
# 运营默认可查风控和设备存活,不能绕过导航直调技术/审计接口。
assert admin_client.get(
"/admin/api/risk-monitor/summary", headers=_auth(operator_token)
).status_code == 200
assert admin_client.get(
"/admin/api/device-liveness/stats", headers=_auth(operator_token)
).status_code == 200
@@ -82,8 +89,9 @@ def test_monitoring_audit_catalog_and_api_permissions(
):
assert admin_client.get(path, headers=_auth(operator_token)).status_code == 403
# 技术角色默认拥有监控审计组全部项权限。
# 技术角色默认拥有监控审计组全部项权限。
for path in (
"/admin/api/risk-monitor/summary",
"/admin/api/device-liveness/stats",
"/admin/api/analytics-health/overview?date_from=2026-07-01T00:00:00Z&date_to=2026-07-02T00:00:00Z",
"/admin/api/event-logs",
@@ -184,8 +192,8 @@ def test_builtin_roles_labels_and_pages(admin_client, super_token) -> None:
# 页集对齐 Prototypes/dashboard/permissions.md 的 ROLES
assert set(roles["finance"]["pages"]) == {"dashboard", "ad-revenue-report", "cps", "withdraws"}
assert set(roles["tech"]["pages"]) == {
"dashboard", "device-liveness", "analytics-health", "config", "ad-revenue", "huawei-review",
"event-logs", "audit-logs",
"dashboard", "risk-monitor", "device-liveness", "analytics-health", "config",
"ad-revenue", "huawei-review", "event-logs", "audit-logs",
}
+40 -39
View File
@@ -11,12 +11,13 @@ import time
import pytest
from app.integrations import sms
from app.integrations.sms import jiguang
def _reset(phone: str) -> None:
"""清该号的进程内存状态,隔离 real 模式用例。"""
sms._codes.pop(phone, None)
sms._last_sent.pop(phone, None)
jiguang._codes.pop(phone, None)
jiguang._last_sent.pop(phone, None)
class _OkResp:
@@ -227,16 +228,16 @@ def test_sms_real_send_calls_jiguang(monkeypatch) -> None:
captured.update(url=url, body=json, auth=headers.get("Authorization", ""))
return _OkResp()
monkeypatch.setattr(sms.settings, "SMS_MOCK", False)
monkeypatch.setattr(sms.httpx, "post", _fake_post)
monkeypatch.setattr(jiguang.settings, "SMS_MOCK", False)
monkeypatch.setattr(jiguang.httpx, "post", _fake_post)
sms.send_code(phone)
jiguang.send_code(phone)
assert captured["url"] == sms.settings.SMS_SEND_ENDPOINT
assert captured["url"] == jiguang.settings.SMS_SEND_ENDPOINT
assert captured["body"]["mobile"] == phone
assert captured["body"]["sign_id"] == sms.settings.SMS_SIGN_ID
assert captured["body"]["temp_id"] == sms.settings.SMS_TEMPLATE_ID
assert captured["body"]["temp_para"]["code"] == sms._codes[phone].code
assert captured["body"]["sign_id"] == jiguang.settings.SMS_SIGN_ID
assert captured["body"]["temp_id"] == jiguang.settings.SMS_TEMPLATE_ID
assert captured["body"]["temp_para"]["code"] == jiguang._codes[phone].code
assert captured["auth"].startswith("Basic ")
@@ -244,33 +245,33 @@ def test_sms_real_verify_one_time_and_wrong(monkeypatch) -> None:
"""real 校验:错误码拒(不消费)→ 正确码成功 → 验过即作废。"""
phone = "13455134000"
_reset(phone)
monkeypatch.setattr(sms.settings, "SMS_MOCK", False)
monkeypatch.setattr(sms.httpx, "post", lambda *a, **k: _OkResp())
monkeypatch.setattr(jiguang.settings, "SMS_MOCK", False)
monkeypatch.setattr(jiguang.httpx, "post", lambda *a, **k: _OkResp())
sms.send_code(phone)
code = sms._codes[phone].code
jiguang.send_code(phone)
code = jiguang._codes[phone].code
wrong = "000000" if code != "000000" else "111111"
assert sms.verify_code(phone, wrong) is False
assert sms.verify_code(phone, code) is True
assert sms.verify_code(phone, code) is False # 已作废
assert jiguang.verify_code(phone, wrong) is False
assert jiguang.verify_code(phone, code) is True
assert jiguang.verify_code(phone, code) is False # 已作废
def test_sms_real_verify_attempts_exhausted(monkeypatch) -> None:
"""real 校验:错误次数到上限即作废,正确码也不再通过(防爆破)。"""
phone = "13466134000"
_reset(phone)
monkeypatch.setattr(sms.settings, "SMS_MOCK", False)
monkeypatch.setattr(sms.settings, "SMS_MAX_VERIFY_ATTEMPTS", 3)
monkeypatch.setattr(sms.httpx, "post", lambda *a, **k: _OkResp())
monkeypatch.setattr(jiguang.settings, "SMS_MOCK", False)
monkeypatch.setattr(jiguang.settings, "SMS_MAX_VERIFY_ATTEMPTS", 3)
monkeypatch.setattr(jiguang.httpx, "post", lambda *a, **k: _OkResp())
sms.send_code(phone)
code = sms._codes[phone].code
jiguang.send_code(phone)
code = jiguang._codes[phone].code
wrong = "000000" if code != "000000" else "111111"
for _ in range(3):
assert sms.verify_code(phone, wrong) is False
assert sms.verify_code(phone, code) is False # 超限作废
assert jiguang.verify_code(phone, wrong) is False
assert jiguang.verify_code(phone, code) is False # 超限作废
def test_sms_real_balance_error_keeps_cooldown(monkeypatch) -> None:
@@ -284,36 +285,36 @@ def test_sms_real_balance_error_keeps_cooldown(monkeypatch) -> None:
def json(self):
return {"error": {"code": 50014, "message": "no money"}}
monkeypatch.setattr(sms.settings, "SMS_MOCK", False)
monkeypatch.setattr(sms.httpx, "post", lambda *a, **k: _ErrResp())
monkeypatch.setattr(jiguang.settings, "SMS_MOCK", False)
monkeypatch.setattr(jiguang.httpx, "post", lambda *a, **k: _ErrResp())
with pytest.raises(sms.SmsError) as ei:
sms.send_code(phone)
jiguang.send_code(phone)
assert ei.value.status_code == 503
assert phone not in sms._codes # 没发出去的码已清
assert phone in sms._last_sent # 冷却保留:失败也限速
assert phone not in jiguang._codes # 没发出去的码已清
assert phone in jiguang._last_sent # 冷却保留:失败也限速
# 立即重试 → 被冷却挡下(429),不会再打极光
with pytest.raises(sms.SmsError) as ei2:
sms.send_code(phone)
jiguang.send_code(phone)
assert ei2.value.status_code == 429
def test_sms_gc_purges_stale_only(monkeypatch) -> None:
"""GC 清过期码 / 旧冷却,但不动今天有效的(阈值设 0 强制每次扫)。"""
monkeypatch.setattr(sms, "_GC_THRESHOLD", 0)
sms._codes.clear()
sms._last_sent.clear()
monkeypatch.setattr(jiguang, "_GC_THRESHOLD", 0)
jiguang._codes.clear()
jiguang._last_sent.clear()
now = time.time()
sms._codes["stale"] = sms._CodeRecord(code="111111", expires_at=now - 1)
sms._codes["fresh"] = sms._CodeRecord(code="222222", expires_at=now + 999)
sms._last_sent["old"] = now - 99999
sms._last_sent["recent"] = now
jiguang._codes["stale"] = jiguang._CodeRecord(code="111111", expires_at=now - 1)
jiguang._codes["fresh"] = jiguang._CodeRecord(code="222222", expires_at=now + 999)
jiguang._last_sent["old"] = now - 99999
jiguang._last_sent["recent"] = now
sms._gc(now)
jiguang._gc(now)
assert "stale" not in sms._codes and "fresh" in sms._codes
assert "old" not in sms._last_sent and "recent" in sms._last_sent
assert "stale" not in jiguang._codes and "fresh" in jiguang._codes
assert "old" not in jiguang._last_sent and "recent" in jiguang._last_sent
# ============================ 用户名 / 默认昵称 ============================
+36 -5
View File
@@ -128,8 +128,7 @@ def test_invite_cash_reject_refunds_invite_account(client, monkeypatch) -> None:
def test_two_accounts_withdraw_independent(client, monkeypatch) -> None:
"""两账户各提各的不串:提 invite_cash(拒绝结清),再提 cash,各扣各账户
:一个用户同一时间只能一个活跃提现单(跨账户),故第二笔需先结清第一笔"""
"""两账户各提各的不串:提 invite_cash(拒绝→验证退款回原账户),再提 cash,各扣各账户互不串。"""
_patch_userinfo(monkeypatch, "openid_ic_3")
token = _login(client, "13800004003")
_seed_balances(client, token, "13800004003", cash=400, invite_cash=500)
@@ -140,7 +139,7 @@ def test_two_accounts_withdraw_independent(client, monkeypatch) -> None:
json={"amount_cents": 200, "source": "invite_cash"},
headers=_auth(token),
)
_reject(r1.json()["out_bill_no"]) # 退回 invite_cash + 结清活跃单
_reject(r1.json()["out_bill_no"]) # 拒绝退回 invite_cash(验证退款回原账户)
r2 = client.post(
"/api/v1/wallet/withdraw",
# 50 分 = 0.5 元档(7-9 起 coin_cash 只能提预设档位)
@@ -154,6 +153,39 @@ def test_two_accounts_withdraw_independent(client, monkeypatch) -> None:
assert cash == 350 # 扣了 cash 50
def test_invite_cash_multiple_in_flight_allowed(client, monkeypatch) -> None:
"""取消「同时仅一单」:invite_cash 无档位上限,已有在途时仍可继续提交,多笔并存各扣 invite_cash。"""
_patch_userinfo(monkeypatch, "openid_ic_multi")
token = _login(client, "13800004006")
_seed_balances(client, token, "13800004006", cash=0, invite_cash=500)
client.post("/api/v1/wallet/bind-wechat", json={"code": "c"}, headers=_auth(token))
r1 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 200, "source": "invite_cash", "out_bill_no": "billicflight0001"},
headers=_auth(token),
)
assert r1.status_code == 200, r1.text
assert r1.json()["status"] == "reviewing"
r2 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 200, "source": "invite_cash", "out_bill_no": "billicflight0002"},
headers=_auth(token),
)
assert r2.status_code == 200, r2.text
assert r2.json()["status"] == "reviewing"
# 两笔 invite_cash 在途并存,共扣 400(500→100),金币现金不动
cash, invite_cash = _balances(client, token)
assert invite_cash == 100 and cash == 0
orders = client.get(
"/api/v1/wallet/withdraw-orders", params={"source": "invite_cash"}, headers=_auth(token)
).json()["items"]
reviewing = [o for o in orders if o["status"] == "reviewing"]
assert len(reviewing) == 2
def test_invite_me_returns_reward_stats(client) -> None:
"""/invite/me 返回 reward_balance_cents(可提现奖励金)。"""
token = _login(client, "13800004004")
@@ -170,12 +202,11 @@ def test_withdraw_orders_source_filter(client, monkeypatch) -> None:
_seed_balances(client, token, "13800004005", cash=400, invite_cash=500)
client.post("/api/v1/wallet/bind-wechat", json={"code": "c"}, headers=_auth(token))
r1 = client.post(
client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 200, "source": "invite_cash"},
headers=_auth(token),
)
_reject(r1.json()["out_bill_no"]) # 结清,才能提第二笔
client.post(
"/api/v1/wallet/withdraw",
# 50 分 = 0.5 元档(7-9 起 coin_cash 只能提预设档位)
+541
View File
@@ -0,0 +1,541 @@
"""风控监控:三类规则、明细及处置闭环。"""
from __future__ import annotations
from datetime import datetime, timedelta
from uuid import uuid4
import pytest
from fastapi.testclient import TestClient
from app.admin.main import admin_app
from app.admin.repositories import admin_user as admin_repo
from app.api.v1 import compare as compare_api
from app.core.config_schema import (
RISK_COMPARE_DAILY_THRESHOLD_KEY,
RISK_ONECLICK_DAILY_THRESHOLD_KEY,
RISK_SMS_HOURLY_THRESHOLD_KEY,
)
from app.core.security import issue_token_pair
from app.db.session import SessionLocal
from app.models.app_config import AppConfig
from app.models.comparison import ComparisonRecord
from app.models.device import DeviceLiveness
from app.models.risk import RiskIncident
from app.repositories import risk as risk_repo
from app.repositories import user as user_repo
@pytest.fixture(autouse=True)
def _reset_risk_rule_config():
keys = (
RISK_SMS_HOURLY_THRESHOLD_KEY,
RISK_ONECLICK_DAILY_THRESHOLD_KEY,
RISK_COMPARE_DAILY_THRESHOLD_KEY,
risk_repo.RISK_RESET_BASELINES_KEY,
)
with SessionLocal() as db:
db.query(AppConfig).filter(AppConfig.key.in_(keys)).delete(
synchronize_session=False
)
db.commit()
yield
with SessionLocal() as db:
db.query(AppConfig).filter(AppConfig.key.in_(keys)).delete(
synchronize_session=False
)
db.commit()
def _admin_token() -> str:
username = "risk_monitor_admin"
with SessionLocal() as db:
if admin_repo.get_by_username(db, username) is None:
admin_repo.create_admin(
db, username=username, password="risk-pass-123", role="super_admin"
)
with TestClient(admin_app) as client:
response = client.post(
"/admin/api/auth/login",
json={"username": username, "password": "risk-pass-123"},
)
assert response.status_code == 200
return response.json()["access_token"]
def _headers() -> dict[str, str]:
return {"Authorization": f"Bearer {_admin_token()}"}
def _seed_incidents() -> tuple[int, int, int, int, str]:
suffix = uuid4().hex[:8]
sms_subject = f"risk_sms_{suffix}"
oneclick_subject = f"risk_oneclick_{suffix}"
now = (risk_repo.utcnow() - timedelta(minutes=5)).replace(microsecond=0)
with SessionLocal() as db:
user = user_repo.upsert_user_for_login(
db,
phone=f"139{int(suffix, 16) % 100_000_000:08d}",
register_channel="sms",
)
second_user = user_repo.upsert_user_for_login(
db,
phone=f"137{int(suffix, 16) % 100_000_000:08d}",
register_channel="sms",
)
db.add(
DeviceLiveness(
user_id=user.id,
device_id=sms_subject,
created_at=now - timedelta(days=10),
)
)
for index in range(5):
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=sms_subject,
device_id=sms_subject,
device_model="OPPO A5",
phone=f"13877771{index:03d}",
outcome="success",
occurred_at=now + timedelta(seconds=index),
evaluate_rule=risk_repo.RULE_SMS_HOURLY,
)
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_LOGIN,
subject_type="device",
subject_id=sms_subject,
user_id=user.id,
device_id=sms_subject,
device_model="OPPO A5",
phone="13877771000",
outcome="success",
occurred_at=now + timedelta(seconds=30),
)
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_LOGIN,
subject_type="device",
subject_id=sms_subject,
user_id=second_user.id,
device_id=sms_subject,
device_model="OPPO A5",
phone="13877771001",
outcome="success",
occurred_at=now + timedelta(seconds=31),
)
for index in range(21):
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_ONECLICK_LOGIN,
subject_type="device",
subject_id=oneclick_subject,
user_id=second_user.id if index == 20 else user.id,
device_id=oneclick_subject,
device_model="OPPO Reno",
phone=user.phone,
outcome="success" if index % 3 else "failed",
reason="mock provider failure" if index % 3 == 0 else None,
occurred_at=now + timedelta(seconds=index),
evaluate_rule=risk_repo.RULE_ONECLICK_DAILY,
)
local_now = now.astimezone(risk_repo.CN_TZ).replace(tzinfo=None)
for index in range(100):
db.add(
ComparisonRecord(
user_id=user.id,
device_id="risk_compare_device",
trace_id=f"risk-monitor-{suffix}-{index}",
store_name="测试门店",
product_names="测试菜品",
status="success",
items=[{"name": "测试菜品", "qty": 1}],
comparison_results=[
{"platform_id": "meituan", "platform_name": "美团", "price": 12.5},
{"platform_id": "taobao", "platform_name": "淘宝", "price": 11.8},
],
skipped_dish_names=[],
saved_amount_cents=70,
created_at=local_now + timedelta(seconds=index),
)
)
db.commit()
risk_repo.sync_compare_incident(db, user_id=user.id, at=local_now)
incidents = {
incident.rule_code: incident.id
for incident in db.query(RiskIncident)
.filter(
RiskIncident.subject_id.in_(
(sms_subject, oneclick_subject, str(user.id))
)
)
.all()
}
return (
incidents[risk_repo.RULE_SMS_HOURLY],
incidents[risk_repo.RULE_ONECLICK_DAILY],
incidents[risk_repo.RULE_COMPARE_DAILY],
user.id,
sms_subject,
)
def test_risk_monitor_summary_lists_and_details() -> None:
sms_id, oneclick_id, compare_id, _, _ = _seed_incidents()
headers = _headers()
with TestClient(admin_app) as client:
summary = client.get("/admin/api/risk-monitor/summary", headers=headers)
assert summary.status_code == 200
cards = {card["kind"]: card for card in summary.json()["cards"]}
assert cards["sms"]["alert_subject_count"] >= 1
assert cards["sms"]["today_total"] >= 5
assert cards["oneclick"]["alert_subject_count"] >= 1
assert cards["oneclick"]["today_total"] >= 20
assert cards["compare"]["alert_subject_count"] >= 1
assert cards["compare"]["today_total"] >= 100
sms_list = client.get("/admin/api/risk-monitor/incidents/sms", headers=headers)
assert sms_list.status_code == 200
sms_item = next(
item for item in sms_list.json()["items"] if item["incident_id"] == sms_id
)
assert sms_item["device_model"] == "OPPO A5"
assert sms_item["event_count"] == 5
assert sms_item["first_used_at"].startswith(
(risk_repo.utcnow() - timedelta(days=10)).date().isoformat()
)
sms_detail = client.get(
f"/admin/api/risk-monitor/incidents/sms/{sms_id}/details",
headers=headers,
)
assert sms_detail.status_code == 200
assert sms_detail.json()["total"] == 5
assert sms_detail.json()["distinct_accounts"] == 2
assert any(item["verified"] for item in sms_detail.json()["items"])
oneclick_detail = client.get(
f"/admin/api/risk-monitor/incidents/oneclick/{oneclick_id}/details?limit=1",
headers=headers,
)
assert oneclick_detail.status_code == 200
assert oneclick_detail.json()["total"] == 21
assert oneclick_detail.json()["distinct_accounts"] == 2
compare_detail = client.get(
f"/admin/api/risk-monitor/incidents/compare/{compare_id}/details",
headers=headers,
)
assert compare_detail.status_code == 200
assert compare_detail.json()["total"] == 100
assert compare_detail.json()["items"][0]["prices"]["美团"] == 12.5
def test_risk_monitor_ignore_and_block_are_enforced(client: TestClient) -> None:
sms_id, oneclick_id, compare_id, user_id, sms_subject = _seed_incidents()
headers = _headers()
with TestClient(admin_app) as admin_client:
ignored = admin_client.post(
f"/admin/api/risk-monitor/incidents/{oneclick_id}/ignore",
json={"reason": "测试确认正常"},
headers=headers,
)
assert ignored.status_code == 200
assert ignored.json()["status"] == "ignored"
blocked_device = admin_client.post(
f"/admin/api/risk-monitor/incidents/{sms_id}/block",
json={"reason": "测试设备异常"},
headers=headers,
)
assert blocked_device.status_code == 200
assert blocked_device.json()["restriction_id"]
blocked_user = admin_client.post(
f"/admin/api/risk-monitor/incidents/{compare_id}/block",
json={"reason": "测试账号异常"},
headers=headers,
)
assert blocked_user.status_code == 200
device_restriction_id = blocked_device.json()["restriction_id"]
user_restriction_id = blocked_user.json()["restriction_id"]
blocked_sms = admin_client.get(
"/admin/api/risk-monitor/incidents/sms?status=blocked",
headers=headers,
)
assert blocked_sms.status_code == 200
assert blocked_sms.json()["items"][0]["restriction_id"] == device_restriction_id
assert blocked_sms.json()["items"][0]["restricted"] is True
blocked_compare = admin_client.get(
"/admin/api/risk-monitor/incidents/compare?status=blocked",
headers=headers,
)
assert blocked_compare.status_code == 200
assert blocked_compare.json()["items"][0]["restriction_id"] == user_restriction_id
assert blocked_compare.json()["items"][0]["restricted"] is True
denied_sms = client.post(
"/api/v1/auth/sms/send",
json={
"phone": "13877779999",
"device_id": sms_subject,
"device_model": "OPPO A5",
},
)
assert denied_sms.status_code == 403
assert denied_sms.json()["detail"] == "当前设备环境异常,暂无法发送验证码"
tokens = issue_token_pair(user_id)
denied_compare = client.post(
"/api/v1/compare/start",
json={"trace_id": "risk-blocked-new-trace", "device_id": "risk_compare_device"},
headers={"Authorization": f"Bearer {tokens['access_token']}"},
)
assert denied_compare.status_code == 403
assert denied_compare.json()["detail"] == "账号存在异常,该功能暂不可用"
denied_compare_step = client.post(
"/api/v1/intent/recognize",
json={},
headers={"Authorization": f"Bearer {tokens['access_token']}"},
)
assert denied_compare_step.status_code == 403
assert denied_compare_step.json()["detail"] == "账号存在异常,该功能暂不可用"
with TestClient(admin_app) as admin_client:
assert (
admin_client.post(
f"/admin/api/risk-monitor/restrictions/{device_restriction_id}/revoke",
headers=headers,
).status_code
== 200
)
assert (
admin_client.post(
f"/admin/api/risk-monitor/restrictions/{user_restriction_id}/revoke",
headers=headers,
).status_code
== 200
)
allowed_sms = client.post(
"/api/v1/auth/sms/send",
json={
"phone": "13877779998",
"device_id": sms_subject,
"device_model": "OPPO A5",
},
)
assert allowed_sms.status_code == 200
with SessionLocal() as db:
assert not risk_repo.is_restricted(
db,
subject_type="user",
subject_id=str(user_id),
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
)
assert db.get(RiskIncident, compare_id).status == "resolved"
def test_admin_can_reset_all_alerts_and_sms_restarts_from_zero() -> None:
sms_id, oneclick_id, compare_id, _, sms_subject = _seed_incidents()
headers = _headers()
with TestClient(admin_app) as client:
reset = client.post("/admin/api/risk-monitor/reset", headers=headers)
assert reset.status_code == 200
payload = reset.json()
assert payload["reset_incident_count"] >= 3
assert payload["reset_counts"]["sms"] >= 1
assert payload["reset_counts"]["oneclick"] >= 1
assert payload["reset_counts"]["compare"] >= 1
summary = client.get("/admin/api/risk-monitor/summary", headers=headers)
assert summary.status_code == 200
cards = {card["kind"]: card for card in summary.json()["cards"]}
assert all(card["alert_subject_count"] == 0 for card in cards.values())
# 重置只清待处理报警,不删除今日事实总量。
assert cards["sms"]["today_total"] >= 5
assert cards["oneclick"]["today_total"] >= 20
assert cards["compare"]["today_total"] >= 100
reset_at = datetime.fromisoformat(payload["reset_at"])
with SessionLocal() as db:
for incident_id in (sms_id, oneclick_id, compare_id):
incident = db.get(RiskIncident, incident_id)
assert incident.status == "resolved"
assert incident.action_reason == risk_repo.MANUAL_RESET_REASON
for index in range(4):
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=sms_subject,
device_id=sms_subject,
phone=f"13688880{index:03d}",
outcome="success",
occurred_at=reset_at + timedelta(seconds=index + 1),
evaluate_rule=risk_repo.RULE_SMS_HOURLY,
)
assert db.get(RiskIncident, sms_id).status == "resolved"
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=sms_subject,
device_id=sms_subject,
phone="13688880999",
outcome="success",
occurred_at=reset_at + timedelta(seconds=5),
evaluate_rule=risk_repo.RULE_SMS_HOURLY,
)
incident = db.get(RiskIncident, sms_id)
assert incident.status == "open"
assert incident.event_count == 5
assert incident.window_start == reset_at.replace(tzinfo=None)
def test_real_compare_harvest_creates_incident_at_threshold() -> None:
suffix = uuid4().hex[:8]
local_now = risk_repo.utcnow().astimezone(risk_repo.CN_TZ).replace(tzinfo=None)
with SessionLocal() as db:
user = user_repo.upsert_user_for_login(
db,
phone=f"136{int(suffix, 16) % 100_000_000:08d}",
register_channel="sms",
)
for index in range(99):
db.add(
ComparisonRecord(
user_id=user.id,
device_id="real-harvest-device",
trace_id=f"risk-real-harvest-{suffix}-{index}",
status="running",
created_at=local_now + timedelta(milliseconds=index),
)
)
db.commit()
user_id = user.id
compare_api._harvest_running_blocking(
f"risk-real-harvest-{suffix}-99",
user_id,
"food",
"real-harvest-device",
{"model": "OPPO A5"},
None,
)
with SessionLocal() as db:
incident = db.query(RiskIncident).filter(
RiskIncident.rule_code == risk_repo.RULE_COMPARE_DAILY,
RiskIncident.subject_id == str(user_id),
RiskIncident.window_key == local_now.strftime("%Y-%m-%d"),
).one()
assert incident.event_count == 100
def test_admin_can_edit_rules_and_current_window_is_reconciled() -> None:
suffix = uuid4().hex[:8]
subject = f"risk_dynamic_sms_{suffix}"
now = risk_repo.utcnow().replace(microsecond=0)
with SessionLocal() as db:
for index in range(3):
risk_repo.record_behavior_event(
db,
event_type=risk_repo.EVENT_SMS_SEND,
subject_type="device",
subject_id=subject,
device_id=subject,
device_model="动态规则测试机",
phone=f"13800001{index:03d}",
outcome="success",
occurred_at=now + timedelta(seconds=index),
)
headers = _headers()
with TestClient(admin_app) as client:
defaults = client.get("/admin/api/risk-monitor/rules", headers=headers)
assert defaults.status_code == 200
assert defaults.json() == {
"sms_hourly_threshold": 5,
"oneclick_daily_threshold": 20,
"compare_daily_threshold": 100,
}
lowered = client.patch(
"/admin/api/risk-monitor/rules",
headers=headers,
json={
"sms_hourly_threshold": 3,
"oneclick_daily_threshold": 20,
"compare_daily_threshold": 100,
},
)
assert lowered.status_code == 200
with SessionLocal() as db:
incident = db.query(RiskIncident).filter(
RiskIncident.rule_code == risk_repo.RULE_SMS_HOURLY,
RiskIncident.subject_id == subject,
).one()
assert incident.status == "open"
assert incident.event_count == 3
raised = client.patch(
"/admin/api/risk-monitor/rules",
headers=headers,
json={
"sms_hourly_threshold": 4,
"oneclick_daily_threshold": 20,
"compare_daily_threshold": 100,
},
)
assert raised.status_code == 200
with SessionLocal() as db:
incident = db.query(RiskIncident).filter(
RiskIncident.rule_code == risk_repo.RULE_SMS_HOURLY,
RiskIncident.subject_id == subject,
).one()
assert incident.status == "resolved"
assert incident.action_reason == risk_repo.AUTO_RESOLVED_REASON
reopened = client.patch(
"/admin/api/risk-monitor/rules",
headers=headers,
json={
"sms_hourly_threshold": 3,
"oneclick_daily_threshold": 20,
"compare_daily_threshold": 100,
},
)
assert reopened.status_code == 200
with SessionLocal() as db:
incident = db.query(RiskIncident).filter(
RiskIncident.rule_code == risk_repo.RULE_SMS_HOURLY,
RiskIncident.subject_id == subject,
).one()
assert incident.status == "open"
assert incident.triggered_at == (
now + timedelta(seconds=2)
).replace(tzinfo=None)
invalid = client.patch(
"/admin/api/risk-monitor/rules",
headers=headers,
json={
"sms_hourly_threshold": 21,
"oneclick_daily_threshold": 20,
"compare_daily_threshold": 100,
},
)
assert invalid.status_code == 422
+185
View File
@@ -0,0 +1,185 @@
"""阿里云短信 provider(Mode A)单元测试。
SDK 交互隔离在 aliyun._call_send / aliyun._call_check 两个薄封装,本文件全程 monkeypatch
它们(返回归一化结果 dict 或抛 SmsError) 不触真 SDK不发网络测的是 provider 的可映射逻辑:
错误码HTTP PASS/UNKNOWN 解释本地失败计数(与极光同语义)mock 短路
"""
from __future__ import annotations
import pytest
from app.core.config import settings
from app.integrations.sms import aliyun
from app.integrations.sms.base import SmsError
PHONE = "13800138000"
def _configure(monkeypatch, *, mock: bool = False) -> None:
"""配齐阿里云凭证 + 设 SMS_MOCK;清本地失败计数隔离用例。"""
monkeypatch.setattr(settings, "SMS_MOCK", mock)
monkeypatch.setattr(settings, "ALIYUN_SMS_ACCESS_KEY_ID", "ak")
monkeypatch.setattr(settings, "ALIYUN_SMS_ACCESS_KEY_SECRET", "sk")
monkeypatch.setattr(settings, "ALIYUN_SMS_SIGN_NAME", "恒创联众")
monkeypatch.setattr(settings, "ALIYUN_SMS_TEMPLATE_CODE", "SMS_100001")
aliyun._verify_attempts.clear()
def _send_ok(phone):
return {"success": True, "code": "OK", "message": "成功", "verify_code": "1234"}
def _check(result):
def _f(phone, code):
return {"success": True, "code": "OK", "message": "成功", "verify_result": result}
return _f
# ============================ 发码 ============================
def test_send_success_returns_interval_and_resets_attempts(monkeypatch) -> None:
_configure(monkeypatch)
aliyun._verify_attempts[PHONE] = 3 # 旧失败计数
monkeypatch.setattr(aliyun, "_call_send", _send_ok)
assert aliyun.send_code(PHONE) == settings.ALIYUN_SMS_INTERVAL_SEC
assert PHONE not in aliyun._verify_attempts # 新码 = 新预算
@pytest.mark.parametrize(
"code,expected",
[
("MOBILE_NUMBER_ILLEGAL", 400),
("BUSINESS_LIMIT_CONTROL", 429),
("FREQUENCY_FAIL", 429),
("FUNCTION_NOT_OPENED", 503),
("INVALID_PARAMETERS", 503),
("SOME_UNEXPECTED_CODE", 503),
],
)
def test_send_maps_error_codes(monkeypatch, code, expected) -> None:
_configure(monkeypatch)
monkeypatch.setattr(
aliyun, "_call_send",
lambda phone: {"success": False, "code": code, "message": code, "verify_code": None},
)
with pytest.raises(SmsError) as ei:
aliyun.send_code(PHONE)
assert ei.value.status_code == expected
def test_send_not_configured_raises_503_without_calling_aliyun(monkeypatch) -> None:
monkeypatch.setattr(settings, "SMS_MOCK", False)
monkeypatch.setattr(settings, "ALIYUN_SMS_ACCESS_KEY_ID", "") # 凭证缺
def _boom(phone):
raise AssertionError("未配置时不应调用阿里云")
monkeypatch.setattr(aliyun, "_call_send", _boom)
with pytest.raises(SmsError) as ei:
aliyun.send_code(PHONE)
assert ei.value.status_code == 503
def test_send_transport_error_propagates_503(monkeypatch) -> None:
_configure(monkeypatch)
def _boom(phone):
raise SmsError("network down", status_code=503)
monkeypatch.setattr(aliyun, "_call_send", _boom)
with pytest.raises(SmsError) as ei:
aliyun.send_code(PHONE)
assert ei.value.status_code == 503
def test_send_mock_returns_interval_no_network(monkeypatch) -> None:
_configure(monkeypatch, mock=True)
def _boom(phone):
raise AssertionError("mock 不应调用阿里云")
monkeypatch.setattr(aliyun, "_call_send", _boom)
assert aliyun.send_code(PHONE) == settings.ALIYUN_SMS_INTERVAL_SEC
# ============================ 校验 ============================
def test_verify_pass_true_and_clears_attempts(monkeypatch) -> None:
_configure(monkeypatch)
aliyun._verify_attempts[PHONE] = 2
monkeypatch.setattr(aliyun, "_call_check", _check("PASS"))
assert aliyun.verify_code(PHONE, "1234") is True
assert PHONE not in aliyun._verify_attempts # 验过即清
def test_verify_unknown_false_and_increments(monkeypatch) -> None:
_configure(monkeypatch)
monkeypatch.setattr(aliyun, "_call_check", _check("UNKNOWN"))
assert aliyun.verify_code(PHONE, "0000") is False
assert aliyun._verify_attempts[PHONE] == 1
assert aliyun.verify_code(PHONE, "0000") is False
assert aliyun._verify_attempts[PHONE] == 2
def test_verify_attempts_cap_short_circuits(monkeypatch) -> None:
_configure(monkeypatch)
aliyun._verify_attempts[PHONE] = settings.SMS_MAX_VERIFY_ATTEMPTS
def _boom(phone, code):
raise AssertionError("达失败上限后不应再调阿里云")
monkeypatch.setattr(aliyun, "_call_check", _boom)
assert aliyun.verify_code(PHONE, "1234") is False # 本地作废
def test_verify_api_error_raises_503(monkeypatch) -> None:
_configure(monkeypatch)
monkeypatch.setattr(
aliyun, "_call_check",
lambda phone, code: {"success": False, "code": "SYSTEM_ERROR",
"message": "err", "verify_result": None},
)
with pytest.raises(SmsError) as ei:
aliyun.verify_code(PHONE, "1234")
assert ei.value.status_code == 503
def test_verify_transport_error_raises_503(monkeypatch) -> None:
_configure(monkeypatch)
def _boom(phone, code):
raise SmsError("network down", status_code=503)
monkeypatch.setattr(aliyun, "_call_check", _boom)
with pytest.raises(SmsError) as ei:
aliyun.verify_code(PHONE, "1234")
assert ei.value.status_code == 503
def test_verify_mock_passes_any_ndigit(monkeypatch) -> None:
_configure(monkeypatch, mock=True)
def _boom(phone, code):
raise AssertionError("mock 不应调用阿里云")
monkeypatch.setattr(aliyun, "_call_check", _boom)
assert aliyun.verify_code(PHONE, "123456") is True # 6 位数字放行
assert aliyun.verify_code(PHONE, "12345") is False # 位数不对
# ============================ 端点:阿里云降级 → 503(auth.py 包 try/except)============================
def test_sms_login_aliyun_outage_returns_503(client, monkeypatch) -> None:
"""SMS_PROVIDER=aliyun 且校验时阿里云异常 → /sms/login 返 503(而非 400/500),便于区分排查。"""
_configure(monkeypatch) # 配齐凭证 + SMS_MOCK=False + 清计数
monkeypatch.setattr(settings, "SMS_PROVIDER", "aliyun")
def _boom(phone, code):
raise SmsError("aliyun down", status_code=503)
monkeypatch.setattr(aliyun, "_call_check", _boom)
r = client.post("/api/v1/auth/sms/login", json={"phone": "13812345678", "code": "1234"})
assert r.status_code == 503, r.text
+47
View File
@@ -0,0 +1,47 @@
"""SMS 分派器:按 settings.SMS_PROVIDER 路由到正确 provider。
契约:send_code / verify_code **每次调用** settings.SMS_PROVIDER provider(支持运行时切换 /
灰度回退);默认 jiguang此处 monkeypatch provider 的实现为标记函数,断言路由命中 + 可秒切
"""
from __future__ import annotations
from app.core.config import settings
from app.integrations import sms
from app.integrations.sms import aliyun, jiguang
def test_send_code_routes_by_provider_and_switches_per_call(monkeypatch) -> None:
calls: list[str] = []
monkeypatch.setattr(jiguang, "send_code", lambda phone: (calls.append("jiguang"), 60)[1])
monkeypatch.setattr(aliyun, "send_code", lambda phone: (calls.append("aliyun"), 60)[1])
monkeypatch.setattr(settings, "SMS_PROVIDER", "jiguang")
assert sms.send_code("13800138000") == 60
monkeypatch.setattr(settings, "SMS_PROVIDER", "aliyun")
assert sms.send_code("13800138000") == 60
assert calls == ["jiguang", "aliyun"] # 每次按当前 provider 路由,运行时可切
def test_verify_code_routes_by_provider(monkeypatch) -> None:
calls: list[str] = []
monkeypatch.setattr(jiguang, "verify_code", lambda phone, code: (calls.append("jiguang"), True)[1])
monkeypatch.setattr(aliyun, "verify_code", lambda phone, code: (calls.append("aliyun"), True)[1])
monkeypatch.setattr(settings, "SMS_PROVIDER", "jiguang")
assert sms.verify_code("13800138000", "123456") is True
monkeypatch.setattr(settings, "SMS_PROVIDER", "aliyun")
assert sms.verify_code("13800138000", "123456") is True
assert calls == ["jiguang", "aliyun"]
def test_unknown_provider_falls_back_to_jiguang(monkeypatch) -> None:
"""SMS_PROVIDER 非 aliyun 一律走 jiguang(默认兜底,防误配把登录打挂)。"""
calls: list[str] = []
monkeypatch.setattr(jiguang, "send_code", lambda phone: (calls.append("jiguang"), 60)[1])
monkeypatch.setattr(aliyun, "send_code", lambda phone: (calls.append("aliyun"), 60)[1])
monkeypatch.setattr(settings, "SMS_PROVIDER", "jiguang")
sms.send_code("13800138000")
assert calls == ["jiguang"]
+36
View File
@@ -190,6 +190,42 @@ def test_exempt_from_device_ip_send_rate_limit(client, enabled, monkeypatch) ->
)
def test_send_is_audit_only_and_does_not_trigger_risk_alert(client, enabled) -> None:
"""测试号不真发短信:可留审计流水,但不能算正式短信下发或触发设备风控。"""
from app.db.session import SessionLocal
from app.models.risk import BehaviorEvent, RiskIncident
from app.repositories import risk as risk_repo
device_id = "qa-test-account-device"
for _ in range(6):
response = client.post(
"/api/v1/auth/sms/send",
json={"phone": TEST_PHONE, "device_id": device_id},
)
assert response.status_code == 200
with SessionLocal() as db:
events = (
db.query(BehaviorEvent)
.filter(
BehaviorEvent.event_type == risk_repo.EVENT_SMS_SEND,
BehaviorEvent.subject_id == device_id,
)
.all()
)
assert len(events) == 6
assert {event.outcome for event in events} == {"test"}
assert (
db.query(RiskIncident)
.filter(
RiskIncident.rule_code == risk_repo.RULE_SMS_HOURLY,
RiskIncident.subject_id == device_id,
)
.count()
== 0
)
# ============================ 计数单元逻辑 ============================
def test_quota_resets_across_day(enabled, monkeypatch) -> None:
+55
View File
@@ -268,6 +268,61 @@ def test_withdraw_idempotent_same_bill_no(client, monkeypatch) -> None:
assert sum(1 for t in r.json()["items"] if t["biz_type"] == "withdraw") == 1
def test_withdraw_multiple_in_flight_allowed(client, monkeypatch) -> None:
"""取消「同时仅一单」:已有在途(reviewing)时,不同 out_bill_no 可继续提交,两单并存。"""
_patch_userinfo(monkeypatch, "openid_multi_inflight")
token = _login(client, "13800002020")
_seed_cash(client, token, "13800002020", 100) # 够两笔 0.5
client.post("/api/v1/wallet/bind-wechat", json={"code": "c"}, headers=_auth(token))
r1 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billmulti00000001"},
headers=_auth(token),
)
assert r1.status_code == 200, r1.text
assert r1.json()["status"] == "reviewing"
r2 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billmulti00000002"},
headers=_auth(token),
)
assert r2.status_code == 200, r2.text
assert r2.json()["status"] == "reviewing"
# 两张在途单并存
r = client.get("/api/v1/wallet/withdraw-orders", headers=_auth(token))
reviewing = [o for o in r.json()["items"] if o["status"] == "reviewing"]
assert len(reviewing) == 2, r.text
# 余额扣两次:100-50-50=0
r = client.get("/api/v1/wallet/account", headers=_auth(token))
assert r.json()["cash_balance_cents"] == 0
def test_withdraw_second_blocked_only_by_insufficient_cash(client, monkeypatch) -> None:
"""并行放开后第二笔仅受余额约束:余额不足返 409「现金余额不足」,而非旧的「已有提现」拦截。"""
_patch_userinfo(monkeypatch, "openid_multi_insuff")
token = _login(client, "13800002021")
_seed_cash(client, token, "13800002021", 50) # 仅够一笔 0.5
client.post("/api/v1/wallet/bind-wechat", json={"code": "c"}, headers=_auth(token))
r1 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billinsuff0000001"},
headers=_auth(token),
)
assert r1.status_code == 200, r1.text
r2 = client.post(
"/api/v1/wallet/withdraw",
json={"amount_cents": 50, "out_bill_no": "billinsuff0000002"},
headers=_auth(token),
)
assert r2.status_code == 409, r2.text
assert "现金余额不足" in r2.json()["detail"]
def test_withdraw_ambiguous_timeout_then_success_no_refund(client, monkeypatch) -> None:
"""#3 转账调用超时(异常),但查单确认已 SUCCESS → 不退款,单置 success。"""
monkeypatch.setattr("app.integrations.wxpay.code_to_userinfo", lambda code: {"openid": "openid_amb", "nickname": None, "avatar_url": None, "raw": {}})