Compare commits

..

3 Commits

Author SHA1 Message Date
左辰勇 3250a465c7 '赚钱中心的接口修改' 2026-07-20 21:15:42 +08:00
左辰勇 de1fd58749 feat(ad): 膨胀弹窗改用服务端权威金额 + 本轮累计口径,下线 signin_boost
要守住的不变量:弹窗数字 == 本轮实际到账之和 == 用户看到的余额涨幅。三者对不上,
用户就会认为少发了钱(走查现象:弹窗 240、余额只涨 40)。

- reward-result:按 ad_session_id 查本次实发金币,替代余额差 / coin_per_ad 估算。
  S2S 异步未到账返 200+pending 而非 404(404 只表示路由不存在,混在一起客户端没法
  区分「后端没部署」和「再等等」);同 session 多条时显式优先 granted——客户端先报
  closed_early、S2S 后到时,granted 反而是后写的。

- boost_round_id:客户端经 mediaExtra 透传「这条广告属于哪一轮膨胀」,穿山甲 S2S 原样
  带回后随发奖记录落库。**纯标签,不参与发奖判定**。reward-result 新增 round_coin,按
  (user_id, boost_round_id) 对 granted 记录求和。之所以由服务端求和而非客户端自己累加
  ——客户端进程被杀/重建后本地累计会丢,发奖记录不会。
  · 求和恒带 user_id:轮 id 是客户端生成的,不带就等于让任何人拿别人的轮 id 查别人发了多少。
  · 本条非 granted(capped 等)时仍返本轮累计、该条按 0 计,让限额 toast 有数可显。
  · test-grant 加可选 boost_round_id:它不经 S2S 拿不到 extra,不补则 debug 包验不了累计。
  · 客户端复用同一轮 id 只会把展示数字滚大,求和的是已发生的记录,不产生新入账,无资损。

- 下线 signin_boost(签到膨胀):它按固定 3000 金币发、与广告实际收益脱钩,产品确认从来
  不是设计内的口径——奖励只有「签到」和「看视频」两种。签到弹窗的「看广告膨胀」改与福利页
  看视频同走 reward_video(按 eCPM 公式)。摘除回调分支、POST /signin/boost、
  SigninBoostRecord、signin_boost_coin 配置,并 drop signin_boost_record 表。
  **coin_transaction.biz_type='signin_boost' 的历史流水保留不动**——钱是真发过的,账必须
  留得住;admin 大盘那两项改从金币流水统计(一次膨胀 = 一笔,与原口径等价),继续能查回历史。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-20 16:55:25 +08:00
zhuzihao e9fd51d119 fix(auth): 发码防刷改为按成功计数 + 新增每设备每日发码上限 (#136)
- 修复:发码限流原为原子「判+记」,被单号 60s 冷却挡下的重发也占设备额度
  → 正常用户连点重发可能被误锁 1 小时。改为「先判后记、只对成功发码计数」:
  check 判在真发之前(超限直接 429、不真发),record 只在 send_code 成功后调;
  被单号冷却 / 供应商失败抛 429 时直接返回、不计数。
- 新增:同一设备(device_id)+ IP 每天最多 20 次发码上限,与原每小时 5 次两道闸并存,
  均按成功计数,叠一层日封顶挡低频长时间轰炸。
- 基建:ratelimit.py 新增 RateLimitRule + check_rate_limits / record_rate_limits
  (peek/commit 拆分);原子的 enforce_rate_limit 仍保留给登录爆破(失败也计)不变。
- 测试:补 2 个用例(冷却挡下不占额度 / 每日上限)。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>

---------

Co-authored-by: zzhyyyyy <2685922758@qq.com>
Reviewed-on: #136
Co-authored-by: zhuzihao <zhuzihao@wonderable.ai>
Co-committed-by: zhuzihao <zhuzihao@wonderable.ai>
2026-07-16 09:40:56 +08:00
53 changed files with 1290 additions and 1421 deletions
@@ -0,0 +1,47 @@
"""ad_reward_record.boost_round_id(金币膨胀「本轮累计」标签)
看完一条激励视频后客户端要弹「本轮累计获得 N 金币」,N 必须等于这一轮实际到账之和(否则用户
会认为少发了钱)。单条到账额 reward-result 已经能给,但「一轮」的边界只有客户端知道(点「放弃
赚钱」才算结束),客户端自己累加又会在进程被杀后丢失。
解法:客户端把轮次 id 随 mediaExtra 透传,穿山甲 S2S 原样带回,发奖时打在记录上;
reward-result 按 (user_id, boost_round_id) 对 granted 记录求和返回 round_coin。
本列是**纯标签**:不参与发奖判定,发多少/发不发完全不受影响。客户端就算一直复用同一个 id,
也只是把展示数字滚大,不产生任何新入账(求和的是已发生的发奖记录),无资损风险。
Revision ID: ad_reward_boost_round_id
Revises: comparison_llm_cost
Create Date: 2026-07-20
"""
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
revision: str = "ad_reward_boost_round_id"
down_revision: str | Sequence[str] | None = "comparison_llm_cost"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
# 可空、无默认:SQLite 原生支持 ADD COLUMN,不需要 batch_alter_table(同 comparison_llm_cost)。
# 存量行留 NULL = 「不属于任何一轮」,求和时天然不参与,老客户端行为不变。
op.add_column(
"ad_reward_record",
sa.Column("boost_round_id", sa.String(length=64), nullable=True),
)
# 求和恒带 user_id(轮 id 是客户端生成的,不能跨用户信任),故建复合索引而非单列
op.create_index(
"ix_ad_reward_user_boost_round",
"ad_reward_record",
["user_id", "boost_round_id"],
unique=False,
)
def downgrade() -> None:
op.drop_index("ix_ad_reward_user_boost_round", table_name="ad_reward_record")
op.drop_column("ad_reward_record", "boost_round_id")
@@ -0,0 +1,56 @@
"""下线签到膨胀:drop signin_boost_record
产品 2026-07 确认「固定 3000 金币的签到膨胀」从来不是设计内的口径 —— 奖励只有「签到」和
「看视频」两种。签到弹窗里的「看广告膨胀」改与福利页看视频走同一条 reward_video 路径
(按 eCPM 公式发,记在 ad_reward_record),signin_boost 场景整体摘除。
⚠️ **只 drop 这张表,不动 coin_transaction**:`biz_type='signin_boost'` 的金币流水是真发过的
钱,账必须留得住(admin 大盘的 signin_boost_coin_total / signin_boost_watch_count 改为从
coin_transaction 统计,继续能查回历史)。本表只是「哪天膨胀过」的业务留痕,金额与去向都能
从流水还原,drop 掉不影响对账。
downgrade 只重建空表结构,**不恢复数据** —— 真要回滚得先从备份捞行。
Revision ID: drop_signin_boost_record
Revises: ad_reward_boost_round_id
Create Date: 2026-07-20
"""
from collections.abc import Sequence
import sqlalchemy as sa
from alembic import op
revision: str = "drop_signin_boost_record"
down_revision: str | Sequence[str] | None = "ad_reward_boost_round_id"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
with op.batch_alter_table("signin_boost_record", schema=None) as batch_op:
batch_op.drop_index(batch_op.f("ix_signin_boost_record_user_id"))
op.drop_table("signin_boost_record")
def downgrade() -> None:
# 只还结构不还数据(见模块 docstring)
op.create_table(
"signin_boost_record",
sa.Column("id", sa.Integer(), nullable=False),
sa.Column("user_id", sa.Integer(), nullable=False),
sa.Column("signin_date", sa.Date(), nullable=False),
sa.Column("coin_awarded", sa.Integer(), nullable=False),
sa.Column("ad_ref_id", sa.String(length=64), nullable=True),
sa.Column(
"created_at", sa.DateTime(timezone=True),
server_default=sa.text("(CURRENT_TIMESTAMP)"), nullable=False,
),
sa.ForeignKeyConstraint(["user_id"], ["user.id"]),
sa.PrimaryKeyConstraint("id"),
sa.UniqueConstraint("user_id", "signin_date", name="uq_signin_boost_user_date"),
)
with op.batch_alter_table("signin_boost_record", schema=None) as batch_op:
batch_op.create_index(
batch_op.f("ix_signin_boost_record_user_id"), ["user_id"], unique=False
)
+9 -2
View File
@@ -25,7 +25,7 @@ from app.models.coupon_state import (
from app.models.cps_order import CpsOrder
from app.models.feedback import Feedback
from app.models.savings import SavingsRecord
from app.models.signin import SigninBoostRecord, SigninRecord
from app.models.signin import SigninRecord
from app.models.user import User
from app.models.wallet import CoinTransaction, WithdrawOrder
@@ -590,7 +590,14 @@ def dashboard_overview(
CoinTransaction.amount > 0,
CoinTransaction.biz_type == "signin_boost",
),
"signin_boost_watch_count": _count(SigninBoostRecord),
# 签到膨胀 2026-07 已下线,signin_boost_record 表随之 drop。这两项保留为**历史口径**
# (钱是真发过的,账要能查回)。次数改数金币流水:一次膨胀 = 一笔 signin_boost 流水,
# 与原来数 signin_boost_record 行数等价。
"signin_boost_watch_count": _count(
CoinTransaction,
CoinTransaction.biz_type == "signin_boost",
CoinTransaction.amount > 0,
),
},
"cash": {
"withdraw_success_cents": _sum(
-58
View File
@@ -1,58 +0,0 @@
"""会话存档独立 poller → app-server 的内部信号端点(server→server, 非客户端接口)。
会话存档轮询在【独立进程】跑(WeWorkFinanceSdk 是 Go c-shared 库, 嵌进 app-server 会 segfault
把主进程带崩, 见 scripts/wx_finance_poller.py)。poller 拉到美团卡片/截图后 POST 到这里, 由
app-server 打比价信号 —— set_pending 在 app-server 进程内存, 心跳才 pop 得到, 故必须回到本进程。
靠共享密钥头 X-Internal-Secret(== settings.INTERNAL_API_SECRET)校验;未配置 → 503。
"""
from __future__ import annotations
import hmac
import logging
from typing import Annotated
from fastapi import APIRouter, Header, HTTPException, status
from app.core import wx_poc_signal
from app.core.config import settings
from app.schemas.wx_finance import WxFinancePendingIn, WxFinancePendingOut
logger = logging.getLogger("shagua.internal.wx_finance")
router = APIRouter(prefix="/internal", tags=["internal"])
def _check_secret(x_internal_secret: str | None) -> None:
"""共享密钥校验。未配置 → 503(挡裸奔);不匹配 → 401(常量时间比较)。"""
configured = settings.INTERNAL_API_SECRET
if not configured:
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="internal api not configured",
)
if not x_internal_secret or not hmac.compare_digest(x_internal_secret, configured):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED, detail="invalid internal secret"
)
@router.post(
"/wx-finance/pending",
response_model=WxFinancePendingOut,
summary="会话存档 poller 触发比价信号(独立进程→app-server, 打 set_pending)",
)
def wx_finance_pending(
payload: WxFinancePendingIn,
x_internal_secret: Annotated[str | None, Header()] = None,
) -> WxFinancePendingOut:
_check_secret(x_internal_secret)
device_id = settings.WX_POC_TEST_DEVICE_ID
if not device_id:
logger.warning("wx_finance pending: WX_POC_TEST_DEVICE_ID 未配置, 忽略信号")
return WxFinancePendingOut(ok=False)
wx_poc_signal.set_pending(device_id, payload.source)
logger.info(
"wx_finance PoC: 已给测试设备 %s 打比价信号 source=%s (触发=%s)",
device_id, payload.source, payload.kind,
)
return WxFinancePendingOut(ok=True)
+79 -96
View File
@@ -3,6 +3,8 @@
路由前缀 `/api/v1/ad`:
GET /pangle-callback 穿山甲 S2S 发奖回调(**无 JWT,靠验签**),穿山甲服务器调
GET /reward-status 客户端查今日看广告发奖进度(Bearer)
GET /reward-result/{ad_session_id}
客户端按会话查本次广告实发金币(Bearer,只读,弹窗金额用)
发奖走服务端:激励视频播完穿山甲回调本接口,验签通过后幂等发金币。客户端只负责
看完后刷新余额,不参与发奖,被破解也刷不到钱。
@@ -13,7 +15,7 @@ import json
import logging
import uuid
from fastapi import APIRouter, Depends, HTTPException, Request, status
from fastapi import APIRouter, Depends, HTTPException, Path, Request, status
from app.api.deps import CurrentUser, DbSession
from app.core import rewards
@@ -25,8 +27,8 @@ from app.repositories import ad_feed_reward as crud_feed
from app.repositories import ad_reward as crud_ad
from app.repositories import ad_watch as crud_watch
from app.repositories import app_config
from app.repositories import signin as crud_signin
from app.schemas.ad import (
AdRewardResultOut,
AdRewardStatusOut,
EcpmReportIn,
EcpmReportOut,
@@ -52,11 +54,14 @@ REASON_BAD_PARAMS = 1 # 验签过但缺 trans_id / user_id 非数字
REASON_UNKNOWN_USER = 2 # user_id 不存在(可能伪造)
REWARD_SCENE_REWARD_VIDEO = "reward_video"
REWARD_SCENE_SIGNIN_BOOST = "signin_boost"
# 提现看视频:看完才能提现的「硬门槛」广告,**不发金币**,只记一条幂等记录(收益由 eCPM 上报口径
# ad_type="withdrawal_video" 单独统计)。故意不放进 SUPPORTED_REWARD_SCENES——它不走发币分支。
REWARD_SCENE_WITHDRAWAL_AD = "withdrawal_ad"
SUPPORTED_REWARD_SCENES = {REWARD_SCENE_REWARD_VIDEO, REWARD_SCENE_SIGNIN_BOOST}
# 2026-07 下线 signin_boost(签到膨胀):它按固定 3000 金币发,与广告实际收益脱钩,产品确认
# 从来不是设计内的口径。签到弹窗里的「看广告膨胀」现在与福利页看视频走同一条 reward_video
# 路径(按 eCPM 公式发),奖励只剩「签到」+「看视频」两种。历史发币流水(coin_transaction
# .biz_type='signin_boost')保留不动——钱是真发过的,账必须留。
SUPPORTED_REWARD_SCENES = {REWARD_SCENE_REWARD_VIDEO}
def _parse_extra(raw_extra: str | None) -> dict[str, str]:
@@ -118,6 +123,11 @@ def pangle_callback(request: Request, db: DbSession) -> PangleCallbackOut:
extra.update(_parse_extra(params.get(extra_key)))
reward_scene = extra.get("reward_scene") or REWARD_SCENE_REWARD_VIDEO
ad_session_id = extra.get("ad_session_id")
# 「这条广告属于哪一轮膨胀」。纯标签:不参与发奖判定,只让 reward-result 能把同一轮求和成
# 弹窗要显示的累计值(见 crud_ad.round_coin_total)。老客户端不带 → NULL → 累计值返 null。
boost_round_id = (extra.get("boost_round_id") or None)
if boost_round_id is not None:
boost_round_id = boost_round_id[:64]
ecpm = params.get("ecpm")
# 环境隔离:激励视频 mediaExtra 里带「这次观看属于哪个后端环境」(srv_env=dev/prod,客户端按
@@ -169,50 +179,11 @@ def pangle_callback(request: Request, db: DbSession) -> PangleCallbackOut:
user_id, trans_id, reward_scene,
)
return PangleCallbackOut(is_verify=False, reason=REASON_BAD_PARAMS)
if reward_scene == REWARD_SCENE_SIGNIN_BOOST:
try:
boost, _balance = crud_signin.boost_today_signin(
db, user_id, ad_ref_id=trans_id, commit=False
)
except crud_signin.NotSignedTodayError:
db.rollback()
rec = crud_ad.record_external_reward(
db, user_id, trans_id, coin=0, reward_scene=reward_scene,
ad_session_id=ad_session_id, ecpm=ecpm,
reward_name=params.get("reward_name"), raw=raw[:1024],
status="not_signed",
)
except crud_signin.AlreadyBoostedError:
db.rollback()
rec = crud_ad.record_external_reward(
db, user_id, trans_id, coin=0, reward_scene=reward_scene,
ad_session_id=ad_session_id, ecpm=ecpm,
reward_name=params.get("reward_name"), raw=raw[:1024],
status="already_boosted",
)
except crud_signin.LastCycleDayBoostBlockedError:
db.rollback()
rec = crud_ad.record_external_reward(
db, user_id, trans_id, coin=0, reward_scene=reward_scene,
ad_session_id=ad_session_id, ecpm=ecpm,
reward_name=params.get("reward_name"), raw=raw[:1024],
status="last_day",
)
else:
rec = crud_ad.record_external_reward(
db, user_id, trans_id, coin=boost.coin_awarded,
reward_scene=reward_scene, ad_session_id=ad_session_id, ecpm=ecpm,
reward_name=params.get("reward_name"), raw=raw[:1024],
commit=False,
)
db.commit()
db.refresh(rec)
else:
rec = crud_ad.grant_ad_reward(
db, user_id, trans_id, ecpm=ecpm, ad_session_id=ad_session_id,
reward_scene=REWARD_SCENE_REWARD_VIDEO,
reward_name=params.get("reward_name"), raw=raw[:1024],
)
rec = crud_ad.grant_ad_reward(
db, user_id, trans_id, ecpm=ecpm, ad_session_id=ad_session_id,
reward_scene=REWARD_SCENE_REWARD_VIDEO, boost_round_id=boost_round_id,
reward_name=params.get("reward_name"), raw=raw[:1024],
)
except crud_ad.UnknownUserError:
logger.warning("pangle callback unknown user_id=%d trans_id=%s", user_id, trans_id)
return PangleCallbackOut(is_verify=False, reason=REASON_UNKNOWN_USER)
@@ -242,6 +213,46 @@ def reward_status(user: CurrentUser, db: DbSession) -> AdRewardStatusOut:
)
@router.get(
"/reward-result/{ad_session_id}",
response_model=AdRewardResultOut,
summary="按 ad_session_id 查本次广告的权威发奖结果",
dependencies=[Depends(rate_limit(120, 60, "ad-reward-result"))],
)
def reward_result(
user: CurrentUser,
db: DbSession,
ad_session_id: str = Path(..., min_length=8, max_length=64, description="本次广告会话 id"),
) -> AdRewardResultOut:
"""客户端看完激励视频后轮询本接口拿**本次实发金币 + 本轮累计**用于弹窗,不再用余额差 /
coin_per_ad 估算(修「弹窗数值与真实金币对不上」)。
round_coin 是「恭喜累计获得奖励」弹窗真正显示的数:本轮(= 客户端的 boost_round_id)所有
granted 记录之和。由服务端求和而不是客户端自己累加——客户端进程被杀/重建后本地累计会丢,
发奖记录不会。取不到轮 id(pending / 老客户端 / extra 丢失)时为 null,客户端退回显示单条。
S2S 回调异步:查不到记录 = 回调还没到 → 返 200 + status='pending' 让客户端继续重试,
**不返 404**(404 只表示路由不存在)。纯只读:发奖仍只由验签过的 S2S 回调完成,
这里不写库、不产生任何奖励,被刷也只是查自己的记录。
"""
rec = crud_ad.find_by_session(db, user.id, ad_session_id)
if rec is None:
# 连记录都没有 → 不知道属于哪一轮,round_coin 一并为 null(不是 0,0 会被当成"本轮没赚到")
return AdRewardResultOut(
ad_session_id=ad_session_id, status="pending", coin=None, round_coin=None,
)
# 本条不是 granted 时**仍返本轮累计**(这条按 0 计):第 3 条撞每日上限那下,客户端的限额
# toast 要显示的是前两条已到账的总额,不是空。
round_coin = (
crud_ad.round_coin_total(db, user.id, rec.boost_round_id)
if rec.boost_round_id
else None
)
return AdRewardResultOut(
ad_session_id=ad_session_id, status=rec.status, coin=rec.coin, round_coin=round_coin,
)
@router.post(
"/watch-report",
response_model=WatchReportOut,
@@ -317,55 +328,27 @@ def test_grant(user: CurrentUser, db: DbSession, payload: TestGrantIn | None = N
if reward_scene not in SUPPORTED_REWARD_SCENES:
raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, detail="bad reward_scene")
# 每次新 trans_id,模拟一次独立的穿山甲发奖回调(幂等键各不相同 → 每次都发,直到当日上限/今日膨胀一次)
# 每次新 trans_id,模拟一次独立的穿山甲发奖回调(幂等键各不相同 → 每次都发,直到当日上限)
trans_id = f"test-{user.id}-{uuid.uuid4().hex}"
if reward_scene == REWARD_SCENE_SIGNIN_BOOST:
try:
boost, _balance = crud_signin.boost_today_signin(
db, user.id, ad_ref_id=trans_id, commit=False
)
except crud_signin.NotSignedTodayError:
db.rollback()
rec = crud_ad.record_external_reward(
db, user.id, trans_id, coin=0, reward_scene=reward_scene,
raw="client debug test-grant signin_boost", status="not_signed",
)
except crud_signin.AlreadyBoostedError:
db.rollback()
rec = crud_ad.record_external_reward(
db, user.id, trans_id, coin=0, reward_scene=reward_scene,
raw="client debug test-grant signin_boost", status="already_boosted",
)
except crud_signin.LastCycleDayBoostBlockedError:
db.rollback()
rec = crud_ad.record_external_reward(
db, user.id, trans_id, coin=0, reward_scene=reward_scene,
raw="client debug test-grant signin_boost", status="last_day",
)
else:
rec = crud_ad.record_external_reward(
db, user.id, trans_id, coin=boost.coin_awarded,
reward_scene=reward_scene, reward_name="测试签到膨胀",
raw="client debug test-grant signin_boost", commit=False,
)
db.commit()
db.refresh(rec)
else:
# 优先用客户端按 ad_session_id 上报的真实 eCPM(走与正式发奖相同的公式);
# 取不到或 eCPM≤0(测试应用常返 0/假值)时兜底 200,保证本地联调仍能验出非零金币。
ad_session_id = payload.ad_session_id if payload is not None else None
ecpm_val = "200"
if ad_session_id:
ecpm_rec = crud_ecpm.find_by_session(db, user_id=user.id, ad_session_id=ad_session_id)
if ecpm_rec is not None and rewards.parse_ecpm_fen(ecpm_rec.ecpm_raw) > 0:
ecpm_val = ecpm_rec.ecpm_raw
try:
rec = crud_ad.grant_ad_reward(
db, user.id, trans_id, ecpm=ecpm_val, ad_session_id=ad_session_id,
reward_name="测试发奖", raw=f"client debug test-grant ecpm={ecpm_val}",
)
except crud_ad.UnknownUserError as e:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="user not found") from e
# 正式链路的轮次 id 走 S2S 的 mediaExtra;本接口不经 S2S,只能由 body 补,否则 debug 包
# 的 reward-result 恒返 round_coin=null,「弹窗 40 → 60」那套累计验收在本地跑不起来。
boost_round_id = (payload.boost_round_id if payload is not None else None) or None
# 优先用客户端按 ad_session_id 上报的真实 eCPM(走与正式发奖相同的公式);
# 取不到或 eCPM≤0(测试应用常返 0/假值)时兜底 200,保证本地联调仍能验出非零金币。
ad_session_id = payload.ad_session_id if payload is not None else None
ecpm_val = "200"
if ad_session_id:
ecpm_rec = crud_ecpm.find_by_session(db, user_id=user.id, ad_session_id=ad_session_id)
if ecpm_rec is not None and rewards.parse_ecpm_fen(ecpm_rec.ecpm_raw) > 0:
ecpm_val = ecpm_rec.ecpm_raw
try:
rec = crud_ad.grant_ad_reward(
db, user.id, trans_id, ecpm=ecpm_val, ad_session_id=ad_session_id,
boost_round_id=boost_round_id,
reward_name="测试发奖", raw=f"client debug test-grant ecpm={ecpm_val}",
)
except crud_ad.UnknownUserError as e:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="user not found") from e
(used, limit, coin_per, round_count, cooldown_until,
_watched, _watch_limit) = crud_ad.today_status(db, user.id)
+24 -15
View File
@@ -12,11 +12,16 @@ from __future__ import annotations
import logging
from fastapi import APIRouter, HTTPException, Request, status
from fastapi import APIRouter, HTTPException, Request
from app.api.deps import CurrentUser, DbSession
from app.core import test_account
from app.core.ratelimit import enforce_rate_limit
from app.core.ratelimit import (
RateLimitRule,
check_rate_limits,
enforce_rate_limit,
record_rate_limits,
)
from app.core.security import TokenError, decode_token, issue_token_pair
from app.integrations.jiguang import JiguangError, mask_phone, verify_and_get_phone
from app.integrations.sms import SmsError, send_code, verify_code
@@ -40,9 +45,10 @@ router = APIRouter(prefix="/api/v1/auth", tags=["auth"])
# 手机号登录防刷:同一设备(device_id) + 同一 IP 每小时最多的登录尝试次数(成功/失败都计)。
SMS_LOGIN_MAX_PER_HOUR = 5
# 发码防刷:同一设备(device_id) + 同一 IP 每小时最多的发码次数。
# 发码防刷(同一设备 device_id + 同一 IP,**只按成功发码计数**;被单号 60s 冷却挡下的重发不占额度):
# 堵「换手机号绕开单号 60s 冷却」的洞 —— 冷却是单号维度,一机换号能绕开。
SMS_SEND_MAX_PER_HOUR_PER_DEVICE = 5
SMS_SEND_MAX_PER_HOUR_PER_DEVICE = 5 # 每小时上限
SMS_SEND_MAX_PER_DAY_PER_DEVICE = 20 # 每天上限(再叠一层日封顶,挡低频长时间轰炸)
def _login_response(
@@ -99,23 +105,26 @@ def sms_send(req: SmsSendRequest, request: Request) -> SmsSendResponse:
logger.info("test_account sms_send short-circuit (不真发)")
return SmsSendResponse(sent=True, mock=True, cooldown_sec=0)
# 防刷:同一设备(device_id) + 同一 IP 每小时最多 SMS_SEND_MAX_PER_HOUR_PER_DEVICE 次发码
# 补「换手机号绕开单号 60s 冷却」的洞(冷却是单号维度,一机换号能绕);设备维度按机器封顶,
# 挡短信轰炸/烧钱。放在真发(send_code)之前 → 超限直接拦下、不真发短信。
enforce_rate_limit(
request,
scope="sms-send-device",
subject=req.device_id,
limit=SMS_SEND_MAX_PER_HOUR_PER_DEVICE,
window_sec=3600,
detail="操作过于频繁,请稍后再试",
)
# 发码防刷:同一设备(device_id) + 同一 IP,每小时 / 每天两道闸,**均只按成功发码计数**
# 补「换手机号绕开单号 60s 冷却」的洞(冷却是单号维度,一机换号能绕);设备维度按机器封顶,挡短信轰炸/烧钱。
# 关键:被单号 60s 冷却挡下的重发是「没真发、没烧钱」→ 不该占额度。故 check(先判)放在真发之前
# (超限直接 429、不真发),record(计数)只在 send_code 成功后调 —— 冷却/供应商失败抛 429 时直接返回、不计数。
send_rules = [
RateLimitRule("sms-send-device", SMS_SEND_MAX_PER_HOUR_PER_DEVICE, 3600,
"操作过于频繁,请稍后再试"),
RateLimitRule("sms-send-device-daily", SMS_SEND_MAX_PER_DAY_PER_DEVICE, 86400,
"今日验证码发送次数过多,请明天再试"),
]
check_rate_limits(request, subject=req.device_id, rules=send_rules)
try:
cooldown = send_code(req.phone)
except SmsError as e:
raise HTTPException(status_code=e.status_code, detail=str(e)) from e
# 发码成功 → 两道闸各 +1(被单号冷却挡下的重发走不到这里,故不占额度)
record_rate_limits(request, subject=req.device_id, rules=send_rules)
from app.core.config import settings # 局部 import 避免循环
return SmsSendResponse(sent=True, mock=settings.SMS_MOCK, cooldown_sec=cooldown)
+3 -22
View File
@@ -14,28 +14,18 @@ import logging
from fastapi import APIRouter
from app.api.deps import CurrentUser, DbSession
from app.core import wx_poc_signal
from app.repositories import device as device_repo
from app.schemas.device import (
DeviceOut,
DeviceRegisterRequest,
HeartbeatRequest,
HeartbeatResponse,
LivenessAckRequest,
LivenessOut,
OkResponse,
PendingCompare,
)
logger = logging.getLogger("shagua.device")
# PoC:源平台代号 → Android 包名(前端 launch 用)。poller 按分享的小程序卡区分平台。
# 客户端只认 source_package:用它反查平台名(弹窗"原平台")+ launch 对应 App,故新增平台只需在此登记包名。
_SOURCE_PACKAGES = {
"meituan": "com.sankuai.meituan",
"jd": "com.jingdong.app.mall", # 京东主 App(含京东秒送/外卖),后端 intent skill=jd_waimai
}
router = APIRouter(prefix="/api/v1/device", tags=["device"])
@@ -62,12 +52,12 @@ def register_device(
return DeviceOut.model_validate(device)
@router.post("/heartbeat", response_model=HeartbeatResponse, summary="上报心跳")
@router.post("/heartbeat", response_model=OkResponse, summary="上报心跳")
def report_heartbeat(
req: HeartbeatRequest,
user: CurrentUser,
db: DbSession,
) -> HeartbeatResponse:
) -> OkResponse:
device_repo.touch_heartbeat(
db,
user_id=user.id,
@@ -75,16 +65,7 @@ def report_heartbeat(
accessibility_enabled=req.accessibility_enabled,
registration_id=req.registration_id,
)
# PoC:该设备有"待比价"信号则带回(只对写死的测试设备生效, 取走即清、只弹一次)
source = wx_poc_signal.pop_pending(req.device_id)
pending = None
if source:
pending = PendingCompare(
source_platform=source,
source_package=_SOURCE_PACKAGES.get(source, ""),
)
logger.info("heartbeat 下发比价信号 device=%s source=%s", req.device_id, source)
return HeartbeatResponse(pending_compare=pending)
return OkResponse()
@router.get("/liveness", response_model=LivenessOut, summary="查询本机掉线告警(后置检测)")
+6 -41
View File
@@ -1,9 +1,12 @@
"""签到 endpoint。
路由前缀 `/api/v1/signin`:
GET /status 今日签到状态 + 14 天档位
GET /status 今日签到状态 + 7 天档位
POST / 执行今日签到
POST /boost 签到后看广告膨胀金币
2026-07 下线 `POST /boost`(签到膨胀):它按固定 3000 金币补发与广告实际收益脱钩
签到弹窗里的看广告膨胀改与福利页看视频走同一条 reward_video 路径( eCPM ,
`/ad/pangle-callback` `/ad/reward-result` 取金额),奖励只剩签到+看视频两种
"""
from __future__ import annotations
@@ -12,15 +15,8 @@ import logging
from fastapi import APIRouter, HTTPException, status
from app.api.deps import CurrentUser, DbSession
from app.repositories import ad_reward as crud_ad
from app.repositories import signin as crud_signin
from app.repositories import wallet as crud_wallet
from app.schemas.welfare import (
SigninBoostRequest,
SigninBoostResultOut,
SigninResultOut,
SigninStatusOut,
)
from app.schemas.welfare import SigninResultOut, SigninStatusOut
logger = logging.getLogger("shagua.signin")
@@ -50,34 +46,3 @@ def do_signin(user: CurrentUser, db: DbSession) -> SigninResultOut:
streak=record.streak,
coin_balance=balance,
)
@router.post("/boost", response_model=SigninBoostResultOut, summary="签到后看广告膨胀金币")
def boost_signin(
payload: SigninBoostRequest, user: CurrentUser, db: DbSession
) -> SigninBoostResultOut:
if not payload.ad_ref_id:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="ad reward required")
ad_rec = crud_ad.find_by_trans(db, payload.ad_ref_id)
if (
ad_rec is None
or ad_rec.user_id != user.id
or ad_rec.reward_scene != "signin_boost"
or ad_rec.status != "granted"
):
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="ad reward not verified")
record = crud_signin.boost_by_ad_ref(db, user.id, payload.ad_ref_id)
if record is None:
raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail="signin boost not granted")
acc = crud_wallet.get_or_create_account(db, user.id)
balance = acc.coin_balance
logger.info(
"signin boost ok user_id=%d date=%s coin=%d",
user.id, record.signin_date, record.coin_awarded,
)
return SigninBoostResultOut(
coin_awarded=record.coin_awarded,
coin_balance=balance,
signin_date=record.signin_date.isoformat(),
)
-159
View File
@@ -1,159 +0,0 @@
"""微信客服(企业微信)消息接收回调 /wx/kf/callback(裸路径, 非 /api/v1)。
对应企业微信微信客服 客服账号 接收消息回调里填的 URL是服务号版(wx_mp.py)
平行实现:下游(set_pending 心跳下发 前端弹窗)完全一致, 只有上游"如何拿到用户消息"不同
服务号是微信把消息直接 POST 给你;微信客服是 POST 一个事件通知, 你再用 <Token> 主动 sync_msg
GET : URL 接入验证 msg_signature 验签(含密文 echostr) 解密 echostr 返回明文
与公众号明文模式不同:企业微信 echostr 是密文, 必须解密后返回
POST : 收事件(密文) msg_signature 验签 AES 解密 解析 XML
Event=kf_msg_or_event <Token>/<OpenKfId> sync_msg 增量拉消息
客户(origin=3)发来的图片(image) set_pending(测试设备, "meituan") + 下载落盘
其余消息/事件记日志忽略一律返回 "success"(异常也吞掉, 不让企业微信重试轰炸)
回调加密与公众号安全模式同源, 直接复用 integrations/wx_mp_crypto(验签 + AES-256-CBC 解密),
差别只在 receiveid corpid(公众号是 appid)MVP 只做"收图打信号 + 落盘证明链路通"
"""
from __future__ import annotations
import logging
import xml.etree.ElementTree as ET
from pathlib import Path
from typing import Any
from fastapi import APIRouter, Request
from fastapi.responses import PlainTextResponse
from app.core import wx_poc_signal
from app.core.config import settings
from app.integrations import wx_kf_client, wx_mp_crypto
logger = logging.getLogger("shagua.wx_kf")
router = APIRouter(prefix="/wx/kf", tags=["wx-kf"])
# 收到的截图暂存目录(MVP 验证用;接上识别后可改为不落盘或定期清理)
_INBOX = Path(settings.MEDIA_ROOT) / "wx_kf_inbox"
@router.get("/callback", include_in_schema=False)
async def verify(
msg_signature: str = "", timestamp: str = "", nonce: str = "", echostr: str = ""
) -> PlainTextResponse:
"""企业微信 URL 接入验证:验签(密文 echostr 参与)→ 解密 echostr → 返回明文。"""
if not settings.wx_kf_callback_configured:
logger.warning("wx_kf verify: 回调凭证未配齐(CorpId/Token/AESKey/Secret)")
return PlainTextResponse("", status_code=503)
if not wx_mp_crypto.verify_msg_signature(
settings.WX_KF_TOKEN, timestamp, nonce, echostr, msg_signature
):
logger.warning("wx_kf verify: signature 校验失败 ts=%s nonce=%s", timestamp, nonce)
return PlainTextResponse("invalid signature", status_code=403)
try:
plain = wx_mp_crypto.decrypt_message(
settings.WX_KF_AES_KEY, settings.WX_KF_CORP_ID, echostr
)
except Exception:
logger.exception("wx_kf verify: echostr 解密失败")
return PlainTextResponse("decrypt failed", status_code=403)
return PlainTextResponse(plain)
@router.post("/callback", include_in_schema=False)
async def receive(request: Request) -> PlainTextResponse:
"""收事件通知(安全模式)。任何异常都吞掉返回 success, 不让企业微信重试轰炸, 靠日志排查。"""
if not settings.wx_kf_callback_configured:
logger.warning("wx_kf receive: 回调凭证未配齐")
return PlainTextResponse("success")
try:
body = (await request.body()).decode("utf-8")
qp = request.query_params
# 安全模式外层 XML 只有 ToUserName + Encrypt。来源=企业微信服务器(HTTPS)+ 下面验签,
# 且 stdlib ET 不扩展外部实体, XXE 不适用。
encrypt = ET.fromstring(body).findtext("Encrypt") or ""
if not wx_mp_crypto.verify_msg_signature(
settings.WX_KF_TOKEN,
qp.get("timestamp", ""),
qp.get("nonce", ""),
encrypt,
qp.get("msg_signature", ""),
):
logger.warning("wx_kf receive: msg_signature 校验失败")
return PlainTextResponse("success")
xml = wx_mp_crypto.decrypt_message(
settings.WX_KF_AES_KEY, settings.WX_KF_CORP_ID, encrypt
)
await _handle_event(ET.fromstring(xml))
except Exception:
logger.exception("wx_kf receive: 处理异常")
return PlainTextResponse("success")
async def _handle_event(root: ET.Element) -> None:
"""客服事件:MsgType=event 且 Event=kf_msg_or_event → 用 Token 增量拉消息。"""
event = root.findtext("Event") or ""
if event != "kf_msg_or_event":
logger.info("wx_kf 收到事件(暂忽略): event=%s", event)
return
token = root.findtext("Token") or ""
open_kfid = root.findtext("OpenKfId") or ""
msgs = await wx_kf_client.sync_messages(token, open_kfid)
logger.info("wx_kf sync_msg 拉到 %d 条 open_kfid=%s", len(msgs), open_kfid)
for msg in msgs:
await _handle_message(msg)
async def _handle_message(msg: dict[str, Any]) -> None:
"""处理单条客服消息。只关心【客户发来的】(origin=3)图片 / 小程序卡片:打比价信号(+ 落盘)。
origin:3=客户发的 / 4=系统推送 / 5=接待人员发的接待人员或系统消息一律不触发比价
两种触发物对应场景:image=用户发结算页截图;miniprogram=用户分享美团小程序卡片(见需求截图)
"""
if msg.get("origin") != 3:
return
msgtype = msg.get("msgtype") or ""
external_userid = msg.get("external_userid") or ""
if msgtype == "image":
media_id = (msg.get("image") or {}).get("media_id", "")
logger.info(
"wx_kf 收到图片: external_userid=%s media_id=%s", external_userid, media_id
)
_fire_pending("image")
await _save_media(media_id, external_userid)
elif msgtype == "miniprogram":
mp = msg.get("miniprogram") or {}
logger.info(
"wx_kf 收到小程序卡片: external_userid=%s title=%s appid=%s",
external_userid,
mp.get("title"),
mp.get("appid"),
)
# 卡片无可下载媒体, 只打信号;源平台/订单识别(可据 appid 判平台)留后续, 现写死美团。
_fire_pending("miniprogram")
else:
logger.info(
"wx_kf 收到消息(暂忽略): external_userid=%s type=%s", external_userid, msgtype
)
def _fire_pending(kind: str) -> None:
"""PoC:给写死的测试设备打"从美团比价"信号(下次心跳带回前端弹窗, 与服务号版完全一致)。"""
poc_dev = settings.WX_POC_TEST_DEVICE_ID
if not poc_dev:
return
wx_poc_signal.set_pending(poc_dev, "meituan")
logger.info("wx_kf PoC: 已给测试设备 %s 打比价信号 source=meituan (触发=%s)", poc_dev, kind)
async def _save_media(media_id: str, external_userid: str) -> None:
"""图片消息只给 media_id, 走临时素材下载接口拉回落盘(PoC 验证用)。"""
content = await wx_kf_client.download_media(media_id)
if not content:
return
try:
_INBOX.mkdir(parents=True, exist_ok=True)
dest = _INBOX / f"{external_userid[:16]}_{media_id[:16]}.jpg"
dest.write_bytes(content)
logger.info("wx_kf 图片已落盘: %s (%d bytes)", dest, len(content))
except Exception:
logger.exception("wx_kf 图片落盘失败 media_id=%s", media_id)
-112
View File
@@ -1,112 +0,0 @@
"""微信服务号消息接收回调 /wx/mp/callback(裸路径, 非 /api/v1)。
对应公众平台"设置与开发 → 基本配置 → 服务器配置"里填的 URL
GET : URL 接入验证 校验 signature 原样返回 echostr(明文)
POST : 收用户消息(安全模式密文) msg_signature 验签 AES 解密 解析 XML
图片消息(MsgType=image): 打日志(openid/PicUrl/MediaId) + PicUrl 下载落盘
其余消息/事件: 记一条日志忽略一律返回 "success"(微信据此不再重试)
MVP 只做"收到图片并落盘"证明链路通;后续接:识别平台+订单 pending_compare 标记
心跳带回前端弹窗,以及 openiddevice 绑定任何异常都返回 "success"不让微信重试轰炸
"""
from __future__ import annotations
import logging
import xml.etree.ElementTree as ET
from pathlib import Path
import httpx
from fastapi import APIRouter, Request
from fastapi.responses import PlainTextResponse
from app.core import wx_poc_signal
from app.core.config import settings
from app.integrations import wx_mp_crypto
logger = logging.getLogger("shagua.wx_mp")
router = APIRouter(prefix="/wx/mp", tags=["wx-mp"])
# 收到的截图暂存目录(MVP 验证用;接上识别后可改为不落盘或定期清理)
_INBOX = Path(settings.MEDIA_ROOT) / "wx_inbox"
@router.get("/callback", include_in_schema=False)
async def verify(
signature: str = "", timestamp: str = "", nonce: str = "", echostr: str = ""
) -> PlainTextResponse:
"""微信 URL 接入验证:校验 signature 通过则原样返回 echostr。"""
if not settings.WX_MP_TOKEN:
logger.warning("wx_mp verify: WX_MP_TOKEN 未配置")
return PlainTextResponse("", status_code=503)
if wx_mp_crypto.verify_url_signature(
settings.WX_MP_TOKEN, timestamp, nonce, signature
):
return PlainTextResponse(echostr)
logger.warning("wx_mp verify: signature 校验失败 ts=%s nonce=%s", timestamp, nonce)
return PlainTextResponse("invalid signature", status_code=403)
@router.post("/callback", include_in_schema=False)
async def receive(request: Request) -> PlainTextResponse:
"""收用户消息(安全模式)。任何异常都吞掉返回 success,不让微信重试轰炸,靠日志排查。"""
if not settings.wx_mp_callback_configured:
logger.warning("wx_mp receive: 回调凭证未配齐(Token/AESKey/AppID)")
return PlainTextResponse("success")
try:
body = (await request.body()).decode("utf-8")
qp = request.query_params
# 安全模式外层 XML 只有 ToUserName + Encrypt。来源=微信服务器(HTTPS)+ 下面验签,
# 且 stdlib ET 不扩展外部实体, XXE 不适用。
encrypt = ET.fromstring(body).findtext("Encrypt") or ""
if not wx_mp_crypto.verify_msg_signature(
settings.WX_MP_TOKEN,
qp.get("timestamp", ""),
qp.get("nonce", ""),
encrypt,
qp.get("msg_signature", ""),
):
logger.warning("wx_mp receive: msg_signature 校验失败")
return PlainTextResponse("success")
xml = wx_mp_crypto.decrypt_message(
settings.WX_MP_AES_KEY, settings.WX_MP_APPID, encrypt
)
await _handle_message(ET.fromstring(xml))
except Exception:
logger.exception("wx_mp receive: 处理异常")
return PlainTextResponse("success")
async def _handle_message(root: ET.Element) -> None:
openid = root.findtext("FromUserName") or ""
msg_type = root.findtext("MsgType") or ""
if msg_type == "image":
pic_url = root.findtext("PicUrl") or ""
media_id = root.findtext("MediaId") or ""
logger.info(
"wx_mp 收到图片: openid=%s media_id=%s pic_url=%s", openid, media_id, pic_url
)
# PoC: 写死"该测试设备要从美团比价", 下次心跳带回前端弹选平台窗
poc_dev = settings.WX_POC_TEST_DEVICE_ID
if poc_dev:
wx_poc_signal.set_pending(poc_dev, "meituan")
logger.info("wx_mp PoC: 已给测试设备 %s 打比价信号 source=meituan", poc_dev)
await _download(pic_url, openid, media_id)
else:
logger.info("wx_mp 收到消息(暂忽略): openid=%s type=%s", openid, msg_type)
async def _download(pic_url: str, openid: str, media_id: str) -> None:
"""用 PicUrl 直接下载图片落盘。PicUrl 是临时公网链接,无需 access_token / IP 白名单。"""
if not pic_url:
return
try:
_INBOX.mkdir(parents=True, exist_ok=True)
async with httpx.AsyncClient(timeout=15) as client:
resp = await client.get(pic_url)
resp.raise_for_status()
dest = _INBOX / f"{openid[:12]}_{media_id[:16]}.jpg"
dest.write_bytes(resp.content)
logger.info("wx_mp 图片已落盘: %s (%d bytes)", dest, len(resp.content))
except Exception:
logger.exception("wx_mp 图片下载失败 pic_url=%s", pic_url)
-54
View File
@@ -136,40 +136,6 @@ class Settings(BaseSettings):
# (不跳授权、不拿 openid),整套微信代码保留。审核通过后 .env 置 true + 重启即启用,无需改代码。
WX_MP_OAUTH_ENABLED: bool = False
# ===== 微信服务号(消息接收回调) =====
# 用户发消息给服务号 → 微信 POST 到 /wx/mp/callback(安全模式:Token 验签 + AESKey 解密)。
# 区别于上面的网页授权(那是落地页拿 openid);这里是被动收用户主动发来的消息(截图比价入口)。
WX_MP_TOKEN: str = "" # 公众平台"服务器配置"的 Token(URL 验签 + 消息验签)
WX_MP_AES_KEY: str = "" # EncodingAESKey(43 位, 消息 AES-256-CBC 加解密)
# ===== 微信截图比价 PoC(临时验证链路, 验证后删) =====
# 只对这台写死的测试设备下发"从美团比价"信号;空=不触发任何设备(部署上线零风险)。
# 收图 → 给此 device_id 打 pending → 该设备下次心跳带回 → 前端弹选平台窗。
WX_POC_TEST_DEVICE_ID: str = ""
# ===== 微信客服(企业微信;截图比价上游入口 —— 服务号版的平行实现) =====
# 用户在「微信客服」会话发消息 → 企业微信 POST 事件到 /wx/kf/callback → 后端 sync_msg 拉消息。
# 回调加密与公众号「安全模式」同源(复用 wx_mp_crypto), 但 receiveid 用 corpid 且 echostr 需解密。
# 拉消息 / 下载图片走 qyapi, 需 access_token(corpid + secret 换取)。全空=回调端点拒收, 零风险。
WX_KF_CORP_ID: str = "" # 企业微信 corpid(回调解密 receiveid + 换 access_token)
WX_KF_SECRET: str = "" # 「微信客服」Secret(换 access_token)
WX_KF_TOKEN: str = "" # 客服「接收消息」回调的 Token(URL 验签 + 消息验签)
WX_KF_AES_KEY: str = "" # 客服「接收消息」回调的 EncodingAESKey(AES-256-CBC 解密)
# ===== 微信会话内容存档(企业微信;"发给成员"那条路的截图/美团卡片比价上游入口) =====
# 用户把美团小程序卡片/截图发给企业微信【成员】(非客服) → 会话存档 GetChatData 轮询拉取 →
# RSA+AES 解密 → 识别 weapp/image → 给测试设备打比价信号。下游(心跳→弹窗→比价)复用。
# 需后台开通「会话内容存档」+ 上传 RSA 公钥 + 配可信IP + 拿存档 Secret;.so 放服务器。
# corpid 复用上面的 WX_KF_CORP_ID(同一企业)。ENABLED 默认关 → 未配好前零影响。
WX_FINANCE_ENABLED: bool = False # 总开关(worker 是否启动)
WX_FINANCE_SECRET: str = "" # 会话存档【专用】Secret(≠微信客服/自建应用)
WX_FINANCE_SDK_PATH: str = "./libWeWorkFinanceSdk_C.so" # WeWorkFinanceSdk C 库(Linux .so)
WX_FINANCE_PRIVATE_KEY_PATH: str = "./secrets/wx_finance_private.pem" # RSA 私钥(解 encrypt_random_key)
WX_FINANCE_POLL_INTERVAL_SEC: float = 0.5 # GetChatData 轮询间隔秒(无回调、自己定频;支持亚秒如 0.5)
WX_FINANCE_SEQ_FILE: str = "./data/wx_finance_seq.txt" # seq 游标持久化(防丢/重复)
WX_FINANCE_RECEIVER_USERID: str = "" # 接收成员 userid(万朗杰);只处理【非他发】的消息, 空=不过滤
WX_FINANCE_INTERNAL_URL: str = "http://127.0.0.1:8770" # poller 通知 app-server 打信号的内部地址(同机)
@property
def wx_mp_configured(self) -> bool:
"""服务号网页授权凭证齐全(缺则落地页不发起授权,降级为无 openid)。"""
@@ -180,26 +146,6 @@ class Settings(BaseSettings):
"""落地页是否真正发起微信授权 = 凭证齐全 且 总开关开。"""
return self.wx_mp_configured and self.WX_MP_OAUTH_ENABLED
@property
def wx_mp_callback_configured(self) -> bool:
"""消息接收回调凭证齐全(Token + AESKey + AppID)。缺则回调端点拒绝处理消息。"""
return bool(self.WX_MP_TOKEN and self.WX_MP_AES_KEY and self.WX_MP_APPID)
@property
def wx_kf_callback_configured(self) -> bool:
"""微信客服回调凭证齐全(CorpId + Token + AESKey + Secret)。缺则回调端点拒绝处理。"""
return bool(
self.WX_KF_CORP_ID
and self.WX_KF_TOKEN
and self.WX_KF_AES_KEY
and self.WX_KF_SECRET
)
@property
def wx_finance_configured(self) -> bool:
"""会话存档 worker 是否该启动 = 总开关开 且 corpid + 存档 Secret 齐全。"""
return bool(self.WX_FINANCE_ENABLED and self.WX_KF_CORP_ID and self.WX_FINANCE_SECRET)
# ===== 微信支付(商家转账到零钱 / 提现)=====
# 真实凭证放 .env(已 gitignore),证书 .pem 放 secrets/。WECHAT_APP_ID 同时用于
# 微信登录(code 换 openid)与转账,必须与 App 端开放平台 appid 一致。
-5
View File
@@ -66,11 +66,6 @@ CONFIG_DEFS: dict[str, dict[str, Any]] = {
"default": r.VIDEO_ROUND_COOLDOWN_SECONDS, "label": "广告关闭后冷却(秒)",
"group": "看广告", "type": "int", "help": "点击退出广告后,下次点击观看前的冷却时间,默认 3 秒。",
},
"signin_boost_coin": {
"default": r.SIGNIN_BOOST_COIN, "label": "签到膨胀固定金币",
"group": "签到", "type": "int",
"help": "Day1-Day6 签到后看完激励视频额外发放的固定金币;Day7 不展示也不允许膨胀。",
},
"comparing_ad_enabled": {
"default": True, "label": "比价/领券期信息流广告",
"group": "看广告", "type": "bool", "hidden": True,
+93 -8
View File
@@ -9,29 +9,41 @@ from __future__ import annotations
import threading
import time
from typing import NamedTuple
from fastapi import HTTPException, Request, status
from app.core.config import settings
# key -> (window_start_ts, count)
_buckets: dict[str, tuple[float, int]] = {}
# key -> (window_start_ts, count, window_sec)
# 存每个 key 自己的 window_sec:_buckets 混着不同窗口(60s 广告 / 3600s 登录 / 86400s 日闸)的 key,
# GC 必须按各 key 自己的窗口判过期(见 [_purge_expired]),否则短窗口调用触发的 GC 会误删长窗口 key。
_buckets: dict[str, tuple[float, int, float]] = {}
_lock = threading.Lock()
_GC_THRESHOLD = 10000 # _buckets 超此阈值才顺手清过期 key(仿 sms.py;测试可 monkeypatch 调小强制每次扫)
def _purge_expired(now: float) -> None:
"""清过期 key(**仅在持有 _lock 时调用**)。按每个 key 自己存的 window_sec 判过期,而非调用方的窗口
_buckets 是全局共享混着 60s(广告)/3600s(登录)/86400s(日闸)不同窗口的 key;若用调用方窗口,
高频的 60s 广告端点触发 GC 时会把本该活 3600s/86400s 的登录/日闸计数一并删掉,使其在规模上(超阈值才
触发本清理)被反复清零而失效仅在超阈值时扫,低频开销可忽略"""
if len(_buckets) <= _GC_THRESHOLD:
return
for k in [k for k, (s, _, w) in _buckets.items() if now - s >= w]:
_buckets.pop(k, None)
def _hit(key: str, limit: int, window_sec: float) -> bool:
"""记一次访问。返回 True=放行,False=超限。"""
now = time.monotonic()
with _lock:
start, count = _buckets.get(key, (now, 0))
start, count, _ = _buckets.get(key, (now, 0, window_sec))
if now - start >= window_sec: # 窗口过期,重置
start, count = now, 0
count += 1
_buckets[key] = (start, count)
# 顺手清过期 key,防内存无限涨(低频访问足够)
if len(_buckets) > 10000:
for k in [k for k, (s, _) in _buckets.items() if now - s >= window_sec]:
_buckets.pop(k, None)
_buckets[key] = (start, count, window_sec)
_purge_expired(now) # 顺手清过期 key(按各自窗口),防内存无限涨
return count <= limit
@@ -83,3 +95,76 @@ def enforce_rate_limit(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail=detail,
)
# ===================== 先判 / 后记(只按「成功」计数)=====================
# _hit 是原子「判+记」:一调用就 +1,适合登录爆破(失败尝试也要计)。但对「短信发码」这类
# **只想给成功动作计数**的场景不合适 —— 被单号冷却挡下的重发没真发、没烧钱,不该占额度。
# 故拆成 _peek(只判不记)+ _commit(只记):check_rate_limits 先判 → 动作 → 成功后 record。
class RateLimitRule(NamedTuple):
"""一条限流规则。scope 区分不同闸(不同 key 前缀);同一 (subject, IP) 在 window_sec
内最多 limit ,超限抛 429 detail 文案
(scope, window_sec) 成对绑在一条规则里 check(先判) record(计数)复用同一条,
避免两处把窗口/scope 写歪导致 key 对不上
"""
scope: str
limit: int
window_sec: float
detail: str = "操作过于频繁,请稍后再试"
def _peek(key: str, limit: int, window_sec: float) -> bool:
"""只读:当前窗口内是否还没到上限(count < limit)。**不改计数**。
[_commit] 配对实现先判后记只在动作成功后才 _commit"""
now = time.monotonic()
with _lock:
start, count, _ = _buckets.get(key, (now, 0, window_sec))
if now - start >= window_sec: # 窗口已过期 → 视作已重置(count 归零)
count = 0
return count < limit
def _commit(key: str, window_sec: float) -> None:
"""记一次访问(+1)。窗口过期则以本次为起点重置。仅在动作成功后调用。"""
now = time.monotonic()
with _lock:
start, count, _ = _buckets.get(key, (now, 0, window_sec))
if now - start >= window_sec: # 窗口过期,重置
start, count = now, 0
_buckets[key] = (start, count + 1, window_sec)
_purge_expired(now) # 顺手清过期 key(按各自窗口,同 [_hit])
def check_rate_limits(request: Request, subject: str, rules: list[RateLimitRule]) -> None:
"""【先判】一组限流:任一规则已达上限即抛 429,且**不改计数**。
配合 [record_rate_limits] 实现只按成功计数: check 所有闸(全未超才继续) 执行动作
动作**成功后** record动作被下游挡下(如短信单号冷却)没真正发生时不 record 不占额度
key = `scope:subject:client_ip`( [enforce_rate_limit] 同款)
"""
if not settings.RATE_LIMIT_ENABLED:
return
ip = _client_ip(request)
for rule in rules:
if not _peek(f"{rule.scope}:{subject}:{ip}", rule.limit, rule.window_sec):
raise HTTPException(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail=rule.detail,
)
def record_rate_limits(request: Request, subject: str, rules: list[RateLimitRule]) -> None:
"""【记一次】一组限流(每条规则 +1)。仅在动作成功后调用,与 [check_rate_limits] 配对。
check动作record 非原子:并发突发下计数可能略超 limit(每个在途请求各 +1)
防脚本/防轰炸的安全网定位可接受;要精确配额需迁 Redis(见模块 docstring)
"""
if not settings.RATE_LIMIT_ENABLED:
return
ip = _client_ip(request)
for rule in rules:
_commit(f"{rule.scope}:{subject}:{ip}", rule.window_sec)
+2 -4
View File
@@ -239,8 +239,8 @@ def calculate_ad_reward_coin(ecpm: str | int | float | None, count_after_this: i
return max(0, round(yuan * COIN_PER_YUAN))
# 签到看广告膨胀:S2S 固定补发(原型 2026-06 由 2000 提到 3000,对应 CTA「看广告最高膨胀至3000金币」)。
SIGNIN_BOOST_COIN: int = 3000
# 签到膨胀(SIGNIN_BOOST_COIN,固定 3000)已于 2026-07 下线:它与广告实际收益脱钩,产品确认
# 非设计内口径。签到弹窗的「看广告膨胀」现与福利页看视频同走 calculate_ad_reward_coin。
# ===== 看激励视频发金币(穿山甲 S2S 服务端回调发奖)=====
@@ -334,5 +334,3 @@ def get_ad_cooldown_sec(db) -> int: # noqa: ANN001
return int(_cfg(db, "ad_cooldown_sec"))
def get_signin_boost_coin(db) -> int: # noqa: ANN001
return int(_cfg(db, "signin_boost_coin"))
-27
View File
@@ -1,27 +0,0 @@
"""微信截图比价 PoC:进程内存的"待比价"信号(单 worker 够用, 重启即失效, PoC 可接受)。
收图端点 set_pending(device_id, source) 该设备下次心跳 pop_pending 取走并清除
心跳响应带回 前端弹选平台窗只对 settings.WX_POC_TEST_DEVICE_ID 写入, 不碰其他设备
后续接真识别 / openiddevice 绑定时整体替换本模块
"""
from __future__ import annotations
import threading
_lock = threading.Lock()
_pending: dict[str, str] = {} # device_id -> source_platform(如 "meituan")
def set_pending(device_id: str, source_platform: str) -> None:
if not device_id:
return
with _lock:
_pending[device_id] = source_platform
def pop_pending(device_id: str) -> str | None:
"""取出并清除该设备的待比价信号;无则 None。心跳每帧调, 取到即消费(只弹一次)。"""
if not device_id:
return None
with _lock:
return _pending.pop(device_id, None)
+2 -1
View File
@@ -13,7 +13,8 @@ worker / 多机时内存不共享 → 冷却、校验都会失效,届时迁移
防刷两层(短信花钱 + `/sms/send` 在登录前无法 JWT 鉴权):
1. 单号 `SMS_SEND_INTERVAL_SEC` 冷却(本文件)
2. 单设备(device_id)每小时频控(api auth.sms_send enforce_rate_limit)+ 极光控制台 IP 白名单/防轰炸(运维侧)
2. 单设备(device_id)+ IP 每小时 / 每天频控(api auth.sms_send check/record_rate_limits,
**只按成功发码计数** 被本文件单号冷却挡下的重发不占额度)+ 极光控制台 IP 白名单/防轰炸(运维侧)
IP 频控(rate_limit 依赖)2026-06-26 按产品要求删除改设备维度; device_id 客户端可伪造/轮换,
脚本轮换 id 能绕过本层 挡脚本狂发主要靠极光控制台侧(+ 可选 nginx 限流)
单号每日上限2026-07-03 按精简要求删除(mentor :登录风控只留单号冷却 + 单设备频控);
-100
View File
@@ -1,100 +0,0 @@
"""微信会话内容存档 SDK 封装(WeWorkFinanceSdk C 库的 ctypes 绑定 + 消息解密)。
企业微信会话内容存档是拿"成员↔外部用户私聊消息正文"(含美团小程序卡片)的官方途径:
主动轮询 GetChatData(seq) 拉取每条消息两段密文:
- encrypt_random_key:企业自持 RSA 私钥解出 AES 密钥
- encrypt_chat_msg :SDK DecryptData 用该密钥解出明文消息 JSON
本模块封装:加载 .so Init GetChatData RSA 解密随机密钥 DecryptData纯外部依赖,
不含业务逻辑(识别 weapp/image 打信号在 core/wx_finance_worker.py)
.so 需放服务器(Linux `libWeWorkFinanceSdk_C.so`),路径 settings.WX_FINANCE_SDK_PATH
CDLL 在实例化时才加载( import ),故本模块在无 .so 的机器上也能安全 import
"""
from __future__ import annotations
import base64
import ctypes
import json
import logging
from typing import Any
from cryptography.hazmat.primitives.asymmetric import padding
from cryptography.hazmat.primitives.serialization import load_pem_private_key
logger = logging.getLogger("shagua.wx_finance")
class WxFinanceError(RuntimeError):
"""SDK 调用返回非 0 / 业务 errcode 时抛。"""
class WxFinanceSdk:
"""WeWorkFinanceSdk 薄封装。实例化即加载 .so + Init(需 corpid/secret 正确 + 服务器 IP 在可信IP)。"""
def __init__(self, sdk_path: str, corpid: str, secret: str, private_key_pem: bytes) -> None:
self._priv = load_pem_private_key(private_key_pem, password=None)
self._lib = ctypes.CDLL(sdk_path) # 缺 .so / 缺 libssl 依赖会在此抛 OSError
self._bind()
self._sdk = self._lib.NewSdk()
ret = self._lib.Init(self._sdk, corpid.encode(), secret.encode())
if ret != 0:
raise WxFinanceError(f"Init 失败 ret={ret}(检查 corpid / 会话存档 Secret / 可信IP)")
def _bind(self) -> None:
lib = self._lib
lib.NewSdk.restype = ctypes.c_void_p
lib.Init.argtypes = [ctypes.c_void_p, ctypes.c_char_p, ctypes.c_char_p]
lib.Init.restype = ctypes.c_int
lib.GetChatData.argtypes = [
ctypes.c_void_p, ctypes.c_ulonglong, ctypes.c_uint,
ctypes.c_char_p, ctypes.c_char_p, ctypes.c_int, ctypes.c_void_p,
]
lib.GetChatData.restype = ctypes.c_int
# ⚠️ DecryptData 不吃 sdk 句柄(与 GetChatData 不同, 它是纯解密函数)——只有 3 个参数,
# 多传 sdk 会让参数错位、encrypt_key 收到 sdk 指针 → DecryptData 返 10008(解析 encrypt_key 出错)。
lib.DecryptData.argtypes = [ctypes.c_char_p, ctypes.c_char_p, ctypes.c_void_p]
lib.DecryptData.restype = ctypes.c_int
lib.NewSlice.restype = ctypes.c_void_p
lib.FreeSlice.argtypes = [ctypes.c_void_p]
lib.GetContentFromSlice.argtypes = [ctypes.c_void_p]
lib.GetContentFromSlice.restype = ctypes.c_void_p
lib.GetSliceLen.argtypes = [ctypes.c_void_p]
lib.GetSliceLen.restype = ctypes.c_int
lib.DestroySdk.argtypes = [ctypes.c_void_p]
def _slice_bytes(self, slc: int) -> bytes:
ptr = self._lib.GetContentFromSlice(slc)
length = self._lib.GetSliceLen(slc)
return ctypes.string_at(ptr, length)
def get_chat_data(self, seq: int, limit: int = 1000, timeout: int = 10) -> list[dict[str, Any]]:
"""拉 seq 之后的消息(返回从 seq+1 起)。返回 chatdata 列表(每项含 seq/encrypt_random_key/encrypt_chat_msg)。"""
slc = self._lib.NewSlice()
try:
ret = self._lib.GetChatData(self._sdk, int(seq), int(limit), None, None, int(timeout), slc)
if ret != 0:
raise WxFinanceError(f"GetChatData 失败 ret={ret}")
data = json.loads(self._slice_bytes(slc).decode("utf-8"))
finally:
self._lib.FreeSlice(slc)
if data.get("errcode"):
raise WxFinanceError(f"GetChatData errcode={data.get('errcode')} {data.get('errmsg')}")
return data.get("chatdata", [])
def decrypt(self, encrypt_random_key_b64: str, encrypt_chat_msg: str) -> dict[str, Any]:
"""RSA 私钥解 encrypt_random_key → 得 AES 密钥 → DecryptData 解 encrypt_chat_msg → 明文 dict。"""
# 企业微信用你上传的 RSA 公钥(PKCS1)加密随机密钥, 这里用对应私钥解出, 原样交给 DecryptData。
aes_key = self._priv.decrypt(base64.b64decode(encrypt_random_key_b64), padding.PKCS1v15())
slc = self._lib.NewSlice()
try:
ret = self._lib.DecryptData(aes_key, encrypt_chat_msg.encode(), slc)
if ret != 0:
raise WxFinanceError(f"DecryptData 失败 ret={ret}")
return json.loads(self._slice_bytes(slc).decode("utf-8"))
finally:
self._lib.FreeSlice(slc)
def close(self) -> None:
if getattr(self, "_sdk", None):
self._lib.DestroySdk(self._sdk)
self._sdk = None
-126
View File
@@ -1,126 +0,0 @@
"""微信客服(企业微信)API 客户端:access_token 缓存 + sync_msg 拉消息 + 临时素材下载。
对应 api/v1/wx_kf.py:收到 `kf_msg_or_event` 回调后, 用这里的函数增量拉取用户消息
- access_token(7200s)进程内缓存, 提前 300s 刷新( worker 够用, 重启即失效)
- sync_msg 游标 next_cursor open_kfid 进程内存续(PoC 重启后从最近 3 天重拉, 可接受)
凭证来自 settings(WX_KF_CORP_ID / WX_KF_SECRET);纯外部 HTTP, 不含业务逻辑
"""
from __future__ import annotations
import asyncio
import logging
import time
from typing import Any
import httpx
from app.core.config import settings
logger = logging.getLogger("shagua.wx_kf")
_QYAPI = "https://qyapi.weixin.qq.com/cgi-bin"
# access_token 进程内缓存(单 worker)。token 空 / 未到刷新点直接用。
_token_lock = asyncio.Lock()
_token_cache: dict[str, Any] = {"token": "", "expire_at": 0.0}
# sync_msg 游标:open_kfid -> next_cursor(PoC 进程内存;重启从最近 3 天重拉)
_cursors: dict[str, str] = {}
async def _get_access_token() -> str:
"""取企业微信 access_token, 进程内缓存, 提前 300s 过期刷新。失败返回空串(调用方降级)。"""
now = time.time()
if _token_cache["token"] and _token_cache["expire_at"] - 300 > now:
return _token_cache["token"]
async with _token_lock:
now = time.time() # 拿锁后复检, 避免并发重复刷新
if _token_cache["token"] and _token_cache["expire_at"] - 300 > now:
return _token_cache["token"]
try:
async with httpx.AsyncClient(timeout=10) as client:
resp = await client.get(
f"{_QYAPI}/gettoken",
params={
"corpid": settings.WX_KF_CORP_ID,
"corpsecret": settings.WX_KF_SECRET,
},
)
data = resp.json()
except Exception:
logger.exception("wx_kf gettoken 请求失败")
return ""
if data.get("errcode"):
logger.warning(
"wx_kf gettoken errcode=%s errmsg=%s", data.get("errcode"), data.get("errmsg")
)
return ""
_token_cache["token"] = data.get("access_token", "")
_token_cache["expire_at"] = now + int(data.get("expires_in", 7200))
return _token_cache["token"]
async def sync_messages(token: str, open_kfid: str) -> list[dict[str, Any]]:
"""收到 kf_msg_or_event 回调后调用:用回调带的一次性 token 增量拉消息, has_more 循环拉完。
token: 回调事件里的 <Token>(消息拉取凭证, 非配置 Token;不传会有严格频控)
open_kfid:客服账号 id(回调里的 <OpenKfId>)
返回本次新拉到的 msg_list(已合并多页);游标按 open_kfid 进程内存续
"""
access_token = await _get_access_token()
if not access_token:
return []
out: list[dict[str, Any]] = []
cursor = _cursors.get(open_kfid, "")
try:
async with httpx.AsyncClient(timeout=15) as client:
for _ in range(20): # 最多 20 页护栏, 防异常时死循环
body: dict[str, Any] = {"cursor": cursor, "token": token, "limit": 1000}
if open_kfid:
body["open_kfid"] = open_kfid
resp = await client.post(
f"{_QYAPI}/kf/sync_msg",
params={"access_token": access_token},
json=body,
)
data = resp.json()
if data.get("errcode"):
logger.warning(
"wx_kf sync_msg errcode=%s errmsg=%s",
data.get("errcode"),
data.get("errmsg"),
)
break
out.extend(data.get("msg_list", []))
cursor = data.get("next_cursor", cursor)
_cursors[open_kfid] = cursor
if not data.get("has_more"):
break
except Exception:
logger.exception("wx_kf sync_msg 请求失败 open_kfid=%s", open_kfid)
return out
async def download_media(media_id: str) -> bytes | None:
"""临时素材下载(图片消息只给 media_id)。成功返回二进制, 失败返回 None。"""
if not media_id:
return None
access_token = await _get_access_token()
if not access_token:
return None
try:
async with httpx.AsyncClient(timeout=20) as client:
resp = await client.get(
f"{_QYAPI}/media/get",
params={"access_token": access_token, "media_id": media_id},
)
resp.raise_for_status()
# 成功=二进制文件流;失败=JSON(errcode)。据 Content-Type 区分。
ctype = resp.headers.get("Content-Type", "")
if "application/json" in ctype or "text/plain" in ctype:
logger.warning("wx_kf media/get 非文件响应: %s", resp.text[:200])
return None
return resp.content
except Exception:
logger.exception("wx_kf media/get 下载失败 media_id=%s", media_id)
return None
-69
View File
@@ -1,69 +0,0 @@
"""微信服务号消息接收回调的验签与解密(安全模式)。
服务号"服务器配置"选安全模式后:
- URL 接入验证(GET): sha1(sort(token, timestamp, nonce)) == signature 原样返回 echostr
- 消息(POST): body <Encrypt> 是密文;
msg_signature = sha1(sort(token, timestamp, nonce, encrypt))
密文 AES-256-CBC 解出: random(16B) + msg_len(4B big-endian) + msg + appid, PKCS7(=32) padding
只做验签 + 解密(收消息)被动回复(加密)MVP 不需要 收到后走客服消息异步回执
密钥/口令来自 settings(WX_MP_TOKEN / WX_MP_AES_KEY / WX_MP_APPID),本模块只做纯算法不读配置
"""
from __future__ import annotations
import base64
import hashlib
import hmac
from cryptography.hazmat.primitives.ciphers import Cipher, algorithms, modes
def verify_url_signature(token: str, timestamp: str, nonce: str, signature: str) -> bool:
"""GET 接入验证:token/timestamp/nonce 三者字典序排序拼接后 sha1。"""
return _consteq(_sha1(token, timestamp, nonce), signature)
def verify_msg_signature(
token: str, timestamp: str, nonce: str, encrypt: str, msg_signature: str
) -> bool:
"""POST 消息验签:四者(含密文 encrypt)字典序排序拼接后 sha1。"""
return _consteq(_sha1(token, timestamp, nonce, encrypt), msg_signature)
def decrypt_message(aes_key_b64: str, expected_appid: str, encrypt_b64: str) -> str:
"""解密 <Encrypt> 密文, 返回明文消息 XML。appid 不符抛 ValueError。
aes_key_b64: EncodingAESKey(43 , 不含结尾 '='), '=' base64 解出 32 字节 AES-256 key
"""
aes_key = base64.b64decode(aes_key_b64 + "=") # 43 → 32 bytes
iv = aes_key[:16]
decryptor = Cipher(algorithms.AES(aes_key), modes.CBC(iv)).decryptor()
plain = decryptor.update(base64.b64decode(encrypt_b64)) + decryptor.finalize()
plain = _pkcs7_unpad(plain)
# random(16) + msg_len(4, big-endian) + msg(msg_len) + from_appid
content = plain[16:]
msg_len = int.from_bytes(content[:4], "big")
msg = content[4 : 4 + msg_len]
from_appid = content[4 + msg_len :].decode("utf-8")
if expected_appid and from_appid != expected_appid:
raise ValueError(f"appid mismatch: {from_appid!r} != {expected_appid!r}")
return msg.decode("utf-8")
def _sha1(*parts: str) -> str:
return hashlib.sha1("".join(sorted(parts)).encode("utf-8")).hexdigest()
def _consteq(a: str, b: str) -> bool:
return hmac.compare_digest(a, b)
def _pkcs7_unpad(data: bytes) -> bytes:
"""微信用块大小 32 的 PKCS7,末字节即 padding 长度(1..32)。越界则原样返回(容错)。"""
if not data:
return data
pad = data[-1]
if pad < 1 or pad > 32:
return data
return data[:-pad]
-9
View File
@@ -19,7 +19,6 @@ from app.api.internal.app_version import router as internal_app_version_router
from app.api.internal.launch_confirm import router as internal_launch_confirm_router
from app.api.internal.price import router as internal_price_router
from app.api.internal.store import router as internal_store_router
from app.api.internal.wx_finance import router as internal_wx_finance_router
from app.api.v1.ad import router as ad_router
from app.api.v1.analytics import router as analytics_router
from app.api.v1.auth import router as auth_router
@@ -40,8 +39,6 @@ from app.api.v1.signin import router as signin_router
from app.api.v1.tasks import router as tasks_router
from app.api.v1.user import router as user_router
from app.api.v1.wallet import router as wallet_router
from app.api.v1.wx_kf import router as wx_kf_router
from app.api.v1.wx_mp import router as wx_mp_router
from app.api.v1.wxpay import router as wxpay_router
from app.core.config import settings
from app.core.daily_exchange_worker import (
@@ -140,15 +137,9 @@ app.include_router(internal_price_router)
app.include_router(internal_store_router)
app.include_router(internal_app_version_router)
app.include_router(internal_launch_confirm_router)
# 会话存档独立 poller → 打比价信号(server→server, X-Internal-Secret;轮询进程在 app 外, 见 scripts/wx_finance_poller.py)
app.include_router(internal_wx_finance_router)
app.include_router(platform_router)
# CPS 群发短链跳转 /c/{code}(公网无鉴权:记点击 → 302 跳美团)
app.include_router(cps_redirect_router)
# 微信服务号消息接收回调 /wx/mp/callback(公网无鉴权:微信服务器验签, 截图比价入口)
app.include_router(wx_mp_router)
# 微信客服(企业微信)消息接收回调 /wx/kf/callback(服务号版的平行实现, 收事件→sync_msg 拉图)
app.include_router(wx_kf_router)
# 用户上传文件(头像)静态服务。生产可改由 nginx 直接 serve MEDIA_ROOT。
_media_root = Path(settings.MEDIA_ROOT)
+1 -1
View File
@@ -37,7 +37,7 @@ from app.models.ops_stat_config import OpsStatConfig # noqa: F401
from app.models.price_observation import PriceObservation # noqa: F401
from app.models.price_report import PriceReport # noqa: F401
from app.models.savings import SavingsRecord # noqa: F401
from app.models.signin import SigninBoostRecord, SigninRecord # noqa: F401
from app.models.signin import SigninRecord # noqa: F401
from app.models.store_mapping import StoreMapping # noqa: F401
from app.models.task import UserTask # noqa: F401
from app.models.user import User # noqa: F401
+9 -1
View File
@@ -8,7 +8,7 @@ from __future__ import annotations
from datetime import datetime
from sqlalchemy import DateTime, ForeignKey, Integer, String, func
from sqlalchemy import DateTime, ForeignKey, Index, Integer, String, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db.base import Base
@@ -16,6 +16,10 @@ from app.db.base import Base
class AdRewardRecord(Base):
__tablename__ = "ad_reward_record"
__table_args__ = (
# 「本轮膨胀累计发了多少」= SUM(coin) WHERE user_id=? AND boost_round_id=? AND status='granted'
Index("ix_ad_reward_user_boost_round", "user_id", "boost_round_id"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
# 穿山甲交易号,幂等键(同号回调不重复发奖)
@@ -31,6 +35,10 @@ class AdRewardRecord(Base):
reward_scene: Mapped[str] = mapped_column(String(32), nullable=False, default="reward_video")
# 客户端生成并通过 extra 透传的广告会话 id
ad_session_id: Mapped[str | None] = mapped_column(String(64), index=True, nullable=True)
# 客户端生成并通过 extra 透传的「膨胀轮」id:一轮 = 用户点「去膨胀」到点「放弃赚钱」之间连看的
# 若干条广告。纯标签,不影响发多少/发不发,只用于把同一轮的发奖记录求和成弹窗要显示的累计值。
# 轮次边界完全由客户端定(它才知道用户点了放弃);老客户端/extra 丢失时为 NULL → 累计值返 null。
boost_round_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
# 本次发奖采用的 eCPM 原始值(回调自带或按 ad_session_id 匹配的客户端上报)
ecpm_raw: Mapped[str | None] = mapped_column(String(32), nullable=True)
# 来源(广告收益报表用):我们的应用环境 prod/test + 我们配置的代码位 104xxx。
+5 -27
View File
@@ -1,6 +1,10 @@
"""签到记录表。
每次签到一行,(user_id, signin_date) 唯一,天然防一天签两次
2026-07 下线 `signin_boost_record`(签到膨胀):膨胀按固定 3000 金币补发与广告实际收益
脱钩,产品确认非设计内口径签到弹窗的看广告膨胀改走 reward_video( eCPM ,记在
`ad_reward_record`)历史发币流水 `coin_transaction.biz_type='signin_boost'` 保留不动
- cycle_day: 1..7,7 天循环里今天落在第几档,决定发多少金币;断签后重置回 1
(周期长度 = rewards.SIGNIN_CYCLE_LEN,2026-06 14 天改 7 天一轮)
- streak: 连续签到天数(不封顶),用于"已连续签到 N 天"展示;断签后重置回 1
@@ -9,7 +13,7 @@ from __future__ import annotations
from datetime import date, datetime
from sqlalchemy import Date, DateTime, ForeignKey, Integer, String, UniqueConstraint, func
from sqlalchemy import Date, DateTime, ForeignKey, Integer, UniqueConstraint, func
from sqlalchemy.orm import Mapped, mapped_column
from app.db.base import Base
@@ -42,29 +46,3 @@ class SigninRecord(Base):
)
class SigninBoostRecord(Base):
"""签到后看广告膨胀记录。
一天最多膨胀一次,补发金额等于当天签到原始奖励独立表用于防并发重复补发,
后续接入真实 S2S 广告 session 时可把 ad_ref_id 回填为广告会话/交易号
"""
__tablename__ = "signin_boost_record"
__table_args__ = (
UniqueConstraint("user_id", "signin_date", name="uq_signin_boost_user_date"),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
user_id: Mapped[int] = mapped_column(
Integer, ForeignKey("user.id"), index=True, nullable=False
)
signin_date: Mapped[date] = mapped_column(Date, nullable=False)
coin_awarded: Mapped[int] = mapped_column(Integer, nullable=False)
ad_ref_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
def __repr__(self) -> str: # pragma: no cover
return f"<SigninBoostRecord user_id={self.user_id} date={self.signin_date} coin={self.coin_awarded}>"
+70 -5
View File
@@ -41,6 +41,61 @@ def find_by_trans(db: Session, trans_id: str) -> AdRewardRecord | None:
return _find_by_trans(db, trans_id)
def find_by_session(db: Session, user_id: int, ad_session_id: str) -> AdRewardRecord | None:
"""按广告会话 id 查该用户本次广告的发奖记录,供客户端轮询弹窗金额(reward-result)。
同一 ad_session_id 可能命中多条,**必须显式优先 granted**,不能只取最近一条:
- 客户端先上报 closed_earlyS2S 随后才姗姗来迟 两条,granted 反而是后写的;
- record_reward_noshow 只在写入前查 granted,挡不住这种后到的竞态;
- 本地联调重复调 test-grant session 多条 granted(trans_id 各不相同)
granted 是唯一真发了钱的状态,取它才是权威金额;都没有再取最近一条,让客户端
知道没发的原因(capped/closed_early) user_id 收窄,防止拿别人的 session 探测结果
"""
granted = db.execute(
select(AdRewardRecord)
.where(
AdRewardRecord.user_id == user_id,
AdRewardRecord.ad_session_id == ad_session_id,
AdRewardRecord.status == "granted",
)
.order_by(AdRewardRecord.created_at.desc())
.limit(1)
).scalars().first()
if granted is not None:
return granted
return db.execute(
select(AdRewardRecord)
.where(
AdRewardRecord.user_id == user_id,
AdRewardRecord.ad_session_id == ad_session_id,
)
.order_by(AdRewardRecord.created_at.desc())
.limit(1)
).scalars().first()
def round_coin_total(db: Session, user_id: int, boost_round_id: str) -> int:
"""本轮膨胀累计已发金币 = 该轮所有 granted 记录的 coin 之和(含刚发的这条)。
客户端弹窗要显示的就是它: 1 条弹 40 2 条弹 60(=40+20),放弃赚钱后余额涨 60,
三个数必须相等之所以由服务端求和而不是客户端自己累加客户端进程被杀/低内存重建后
本地累计就丢了,而发奖记录不会丢
**必须带 user_id**:boost_round_id 是客户端生成的,不带 user_id 就等于让任何人拿别人的
id 查别人发了多少未发奖的状态(capped/closed_early/ecpm_missing)coin 本就是 0,
这里按 status 过滤只是让意图显式
"""
return int(
db.execute(
select(func.coalesce(func.sum(AdRewardRecord.coin), 0)).where(
AdRewardRecord.user_id == user_id,
AdRewardRecord.boost_round_id == boost_round_id,
AdRewardRecord.status == "granted",
)
).scalar_one()
)
def _granted_today(db: Session, user_id: int, reward_date: str) -> int:
return db.execute(
select(func.count())
@@ -78,8 +133,12 @@ def grant_ad_reward(
reward_scene: str = "reward_video",
reward_name: str | None = None,
raw: str | None = None,
boost_round_id: str | None = None,
) -> AdRewardRecord:
"""福利页激励视频发奖(幂等 + 每日限额 + eCPM 公式)。"""
"""福利页激励视频发奖(幂等 + 每日限额 + eCPM 公式)。
boost_round_id 只是随记录存下的标签( round_coin_total),**不参与任何发奖判定**
"""
# #2 幂等:同 trans_id 已处理过 → 原样返回,不重复发
existing = _find_by_trans(db, trans_id)
if existing is not None:
@@ -112,7 +171,7 @@ def grant_ad_reward(
trans_id=trans_id, user_id=user_id, coin=0, status="capped",
reward_date=today, reward_name=reward_name, raw=raw,
reward_scene=reward_scene, ad_session_id=ad_session_id, ecpm_raw=ecpm,
app_env=src_app_env, our_code_id=src_code_id,
app_env=src_app_env, our_code_id=src_code_id, boost_round_id=boost_round_id,
)
return _commit_record(db, rec, trans_id)
@@ -123,7 +182,7 @@ def grant_ad_reward(
trans_id=trans_id, user_id=user_id, coin=0, status="ecpm_missing",
reward_date=today, reward_name=reward_name, raw=raw,
reward_scene=reward_scene, ad_session_id=ad_session_id, ecpm_raw=None,
app_env=src_app_env, our_code_id=src_code_id,
app_env=src_app_env, our_code_id=src_code_id, boost_round_id=boost_round_id,
)
return _commit_record(db, rec, trans_id)
@@ -140,7 +199,7 @@ def grant_ad_reward(
trans_id=trans_id, user_id=user_id, coin=coin, status="granted",
reward_date=today, reward_name=reward_name, raw=raw,
reward_scene=reward_scene, ad_session_id=ad_session_id, ecpm_raw=ecpm_raw,
app_env=src_app_env, our_code_id=src_code_id,
app_env=src_app_env, our_code_id=src_code_id, boost_round_id=boost_round_id,
)
return _commit_record(db, rec, trans_id)
@@ -206,8 +265,13 @@ def record_external_reward(
raw: str | None = None,
status: str = "granted",
commit: bool = True,
boost_round_id: str | None = None,
) -> AdRewardRecord:
"""记录非普通看视频场景的 S2S 回调幂等,发币由调用方业务仓储完成。"""
"""记录非普通看视频场景的 S2S 回调幂等,发币由调用方业务仓储完成。
boost_round_id grant_ad_reward:纯标签签到膨胀场景的 coin 也会计入本轮累计
(它的 coin 就是实发额),所以这里也要存,否则一轮里混了膨胀就会漏算
"""
existing = _find_by_trans(db, trans_id)
if existing is not None:
return existing
@@ -224,6 +288,7 @@ def record_external_reward(
reward_scene=reward_scene,
ad_session_id=ad_session_id,
ecpm_raw=ecpm,
boost_round_id=boost_round_id,
)
db.add(rec)
if commit:
+1 -80
View File
@@ -11,12 +11,11 @@ from dataclasses import dataclass
from datetime import timedelta
from sqlalchemy import select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.core import rewards
from app.core.rewards import SIGNIN_CYCLE_LEN, cn_today
from app.models.signin import SigninBoostRecord, SigninRecord
from app.models.signin import SigninRecord
from app.repositories import wallet as crud_wallet
@@ -24,18 +23,6 @@ class AlreadySignedError(Exception):
"""今天已经签过了。"""
class NotSignedTodayError(Exception):
"""今天尚未签到,不能膨胀。"""
class AlreadyBoostedError(Exception):
"""今天签到奖励已经膨胀过。"""
class LastCycleDayBoostBlockedError(Exception):
"""循环最后一天(第 SIGNIN_CYCLE_LEN 天)不允许签到膨胀。"""
@dataclass
class SigninStep:
day: int # 1..14
@@ -141,69 +128,3 @@ def do_signin(db: Session, user_id: int) -> tuple[SigninRecord, int]:
db.commit()
db.refresh(record)
return record, acc.coin_balance
def _today_record(db: Session, user_id: int) -> SigninRecord | None:
today = cn_today()
return db.execute(
select(SigninRecord).where(
SigninRecord.user_id == user_id,
SigninRecord.signin_date == today,
)
).scalar_one_or_none()
def boost_by_ad_ref(
db: Session, user_id: int, ad_ref_id: str
) -> SigninBoostRecord | None:
"""按广告交易号查签到膨胀记录。S2S 发奖后客户端确认用。"""
return db.execute(
select(SigninBoostRecord).where(
SigninBoostRecord.user_id == user_id,
SigninBoostRecord.ad_ref_id == ad_ref_id,
)
).scalar_one_or_none()
def boost_today_signin(
db: Session, user_id: int, *, ad_ref_id: str | None = None, commit: bool = True
) -> tuple[SigninBoostRecord, int]:
"""签到后看广告膨胀:固定补发配置金币。返回 (膨胀记录, 补发后余额)。"""
record = _today_record(db, user_id)
if record is None:
raise NotSignedTodayError
if record.cycle_day == SIGNIN_CYCLE_LEN:
raise LastCycleDayBoostBlockedError
today = record.signin_date
existing = db.execute(
select(SigninBoostRecord).where(
SigninBoostRecord.user_id == user_id,
SigninBoostRecord.signin_date == today,
)
).scalar_one_or_none()
if existing is not None:
raise AlreadyBoostedError
boost = SigninBoostRecord(
user_id=user_id,
signin_date=today,
coin_awarded=rewards.get_signin_boost_coin(db),
ad_ref_id=ad_ref_id,
)
db.add(boost)
try:
acc, _ = crud_wallet.grant_coins(
db, user_id, boost.coin_awarded,
biz_type="signin_boost", ref_id=ad_ref_id or today.isoformat(),
remark=f"签到膨胀 第{record.cycle_day}",
)
if commit:
db.commit()
else:
db.flush()
except IntegrityError as e:
db.rollback()
raise AlreadyBoostedError from e
db.refresh(boost)
return boost, acc.coin_balance
+38 -2
View File
@@ -43,6 +43,35 @@ class AdRewardStatusOut(BaseModel):
watch_seconds_remaining: int = Field(0, description="今日剩余可观看秒数;limit=0 时客户端不据此拦截")
class AdRewardResultOut(BaseModel):
"""按 ad_session_id 查本次广告的**权威发奖结果**(福利页看完视频的弹窗金额只认它)。
S2S 回调是异步的,客户端看完广告立刻轮询多半还查不到记录 这种还没到账返回
200 + status='pending' 让客户端重试,**不返 404**:404 只应表示路由不存在,两者混在
一起客户端没法区分后端没部署再等等纯只读,不产生任何奖励
"""
ad_session_id: str = Field(..., description="回显请求的广告会话 id")
status: str = Field(
...,
description="pending(S2S 未到账,客户端应继续轮询) / granted(已发奖) / capped(当日超限未发) / "
"ecpm_missing(缺 eCPM 未发) / closed_early(提前关闭未发);其余同 AdRewardRecord.status。"
"客户端只在 granted 且 coin>0 时弹窗,其它一律不弹(不显示假数字)",
)
coin: int | None = Field(
None,
description="本次实发金币:granted 为真实到账额;未发奖的状态为 0;pending 为 null",
)
round_coin: int | None = Field(
None,
description="**本轮膨胀累计已发金币**(含本条)——客户端「恭喜累计获得奖励」弹窗显示的就是它。"
"轮 = 用户点「去膨胀」到点「放弃赚钱」之间连看的若干条广告,边界由客户端的 boost_round_id 定。"
"本条不是 granted(capped/closed_early/…)时**仍返本轮累计**,只是这条按 0 计。"
"pending(没记录,取不到轮 id)、或该记录没有 boost_round_id(老客户端 / extra 丢失)时为 null,"
"客户端见 null 退回只显示单条 coin",
)
class EcpmReportIn(BaseModel):
"""客户端上报一次广告展示的 eCPM(内部收益统计/对账)。
@@ -105,13 +134,20 @@ class TestGrantIn(BaseModel):
reward_scene: str = Field(
"reward_video",
description="模拟发奖场景:reward_video(普通激励视频) / signin_boost(签到膨胀)",
description="模拟发奖场景。当前只支持 reward_video(普通激励视频);signin_boost(签到膨胀)"
"已于 2026-07 下线,传它会 422",
)
ad_session_id: str | None = Field(
None, min_length=8, max_length=64,
description="本次广告会话 id(与 ecpm-report 同值)。reward_video 场景下据此查回客户端"
"已上报的真实 eCPM 来按公式发奖;查不到或 eCPM≤0 时兜底 200,保证本地联调仍出非零金币",
)
boost_round_id: str | None = Field(
None, max_length=64,
description="本次广告属于哪一轮膨胀。正式链路走穿山甲 S2S 的 mediaExtra,test-grant 不经 S2S、"
"拿不到 extra,故在 body 里补一个——不传的话 debug 包 reward-result 的 round_coin 恒为 null,"
"「弹窗 40 → 60 → toast +60」那套验收在本地跑不起来",
)
class TestGrantOut(BaseModel):
@@ -119,7 +155,7 @@ class TestGrantOut(BaseModel):
granted: bool = Field(..., description="本次是否真的发了金币(达每日上限则 False)")
status: str = Field(
..., description="granted / capped / not_signed / already_boosted / last_day / unknown_scene"
..., description="granted / capped / ecpm_missing / unknown_scene"
)
coin: int = Field(..., description="本次发放金币(capped 时为 0)")
used_today: int = Field(..., description="今日已成功发奖次数")
-14
View File
@@ -36,20 +36,6 @@ class OkResponse(BaseModel):
ok: bool = True
class PendingCompare(BaseModel):
"""PoC:后端通过心跳下发的"用户要从某源平台比价"信号。前端据此弹选平台窗 + launch 源平台。"""
source_platform: str # pricebot 源平台代号(PoC 写死 "meituan")
source_package: str # 源平台 Android 包名(前端 launch 用)
class HeartbeatResponse(BaseModel):
"""心跳响应。常规只回 ok;PoC 期带回 pending_compare 触发比价(exclude_none 省略 null)。"""
model_config = ConfigDict()
ok: bool = True
pending_compare: PendingCompare | None = None
class LivenessOut(BaseModel):
"""本机掉线告警状态(后置检测 pull 版)。客户端只需这一个布尔判断要不要弹「开启自启动」引导,
故只返回 kill_alert_pending(不暴露设备详情 / 内部 liveness_state )从未注册过 默认 False(无告警)"""
-10
View File
@@ -223,16 +223,6 @@ class SigninResultOut(BaseModel):
coin_balance: int = Field(..., description="签到后金币余额")
class SigninBoostRequest(BaseModel):
ad_ref_id: str | None = Field(None, description="广告会话/交易号。当前开发期可空,后续接 S2S 时回填")
class SigninBoostResultOut(BaseModel):
coin_awarded: int = Field(..., description="本次膨胀补发金币")
coin_balance: int = Field(..., description="膨胀补发后金币余额")
signin_date: str = Field(..., description="被膨胀的签到日期 YYYY-MM-DD")
# ===== 任务 =====
class TaskOut(BaseModel):
-13
View File
@@ -1,13 +0,0 @@
"""会话存档独立 poller → app-server 内部信号端点的 schema。"""
from __future__ import annotations
from pydantic import BaseModel
class WxFinancePendingIn(BaseModel):
source: str = "meituan" # 源平台代号(pop 后前端 launch 用;PoC 写死 meituan)
kind: str = "" # 触发物类型 weapp / image(仅日志用)
class WxFinancePendingOut(BaseModel):
ok: bool = True
-18
View File
@@ -1,18 +0,0 @@
[Unit]
# 会话存档轮询 poller —— 独立进程跑 WeWorkFinanceSdk(Go c-shared .so),与 app-server 隔离:
# 嵌进 app-server 会 segfault 把主进程带崩,拆出来后崩了 systemd 只重启本服务、不影响 app-server。
Description=Shaguabijia 会话存档轮询 poller (WeWorkFinanceSdk, isolated from app-server)
After=network.target shaguabijia-app-server.service
[Service]
Type=simple
User=root
WorkingDirectory=/opt/shaguabijia-app-server
Environment="PATH=/opt/shaguabijia-app-server/.venv/bin:/usr/bin:/bin"
EnvironmentFile=/opt/shaguabijia-app-server/.env
ExecStart=/opt/shaguabijia-app-server/.venv/bin/python scripts/wx_finance_poller.py
Restart=on-failure
RestartSec=5
[Install]
WantedBy=multi-user.target
+1 -1
View File
@@ -78,7 +78,6 @@
| **签到**(前缀 `/api/v1/signin` |||
| 25 | `GET /api/v1/signin/status` | Bearer | [详情](./signin/signin-status.md) |
| 26 | `POST /api/v1/signin` | Bearer | [详情](./signin/signin-do.md) |
| 26a | `POST /api/v1/signin/boost` | Bearer | [详情](./signin/signin-boost.md) |
| **任务**(前缀 `/api/v1/tasks` |||
| 27 | `GET /api/v1/tasks` | Bearer | [详情](./tasks/tasks-list.md) |
| 28 | `POST /api/v1/tasks/{task_key}/claim` | Bearer | [详情](./tasks/tasks-claim.md) |
@@ -89,6 +88,7 @@
| **看广告发奖**(前缀 `/api/v1/ad` |||
| 32 | `GET /api/v1/ad/pangle-callback` | 验签 | [详情](./ad/ad-pangle-callback.md) |
| 33 | `GET /api/v1/ad/reward-status` | Bearer | [详情](./ad/ad-reward-status.md) |
| 33a | `GET /api/v1/ad/reward-result/{ad_session_id}` | Bearer | [详情](./ad/ad-reward-result.md)(本次实发金币 + 本轮膨胀累计 `round_coin`,弹窗数字用它) |
| 34 | `POST /api/v1/ad/test-grant` | Bearer | [详情](./ad/ad-test-grant.md) |
| 35 | `POST /api/v1/ad/ecpm-report` | Bearer | [详情](./ad/ad-ecpm-report.md) |
| 35a | `POST /api/v1/ad/feed-reward` | Bearer | [详情](./ad/ad-feed-reward.md) |
+10 -2
View File
@@ -15,7 +15,15 @@ GroMore 以 GET 回调,关键参数:
| `trans_id` | string | 交易号(**幂等键** + **唯一参与签名的字段** |
| `reward_name` | string | 奖励名(广告位配置,入库备注) |
| `ecpm` | string\|null | GroMore 回调携带的 eCPM。普通激励视频优先用它计算金币 |
| `extra` / `gromoreExtra` / `gromore_extra` | string | 客户端透传 JSON。支持 `ad_session_id``reward_scene`;`reward_scene=signin_boost` 表示签到膨胀 |
| `extra` / `gromoreExtra` / `gromore_extra` | string | 客户端透传 JSON。支持 `ad_session_id``reward_scene``srv_env``boost_round_id` |
### `extra` 里的 `boost_round_id`
客户端生成的「这条广告属于哪一轮膨胀」标签(32 位十六进制,同 `ad_session_id` 格式),随发奖记录存进 `ad_reward_record.boost_round_id`
**它不参与任何发奖判定** —— 发多少、发不发完全不受影响,只是让 [`/ad/reward-result`](./ad-reward-result.md) 能把同一轮的 granted 记录求和成 `round_coin`(客户端「恭喜累计获得奖励」弹窗显示的数)。
轮次边界由客户端定(只有它知道用户点没点「放弃赚钱」):点「去膨胀」新生成一个 → 点「继续看视频膨胀」复用同一个 → 点「放弃赚钱」/ ✕ / 返回 / 到每日上限 / 跨天 则丢弃。不带此字段(老客户端 / GroMore 偶发丢 extra)时存 NULL,`round_coin``null`
| `mediation_rit` | string | 代码位 IDGroMore 带,目前仅入 raw 备查) |
| `prime_rit` | string | 广告位 ID(同上) |
| `adn_name` | string | 实际出广告的 ADN 名(同上,可用于收益分析) |
@@ -40,6 +48,6 @@ GroMore 以 GET 回调,关键参数:
**发奖唯一可信入口**:验签 → 取 `user_id`/`extra` → 按 `reward_scene` 分流 → 幂等处理(按 `trans_id` 去重)。客户端不直接发奖,被破解也刷不到钱。
- `reward_scene=reward_video` 或缺省:普通激励视频。金币按 `eCPM / 1000 * eCPM因子 * 当日次数因子 * 10000` 计算;若回调没有 `ecpm`,会按 `extra.ad_session_id` 查客户端 `/ad/ecpm-report` 的上报值;两边都没有 eCPM 时不发币,记录 `status=ecpm_missing`
- `reward_scene=signin_boost`:签到膨胀。要求用户当天已签到且不是 Day14;看完视频固定发 `2000` 金币,写 `signin_boost_record``coin_transaction.biz_type=signin_boost`
- ~~`reward_scene=signin_boost`~~(签到膨胀):**2026-07 已下线**。它按固定 3000 金币发、与广告实际收益脱钩,产品确认非设计内口径。签到弹窗的「看广告膨胀」现与福利页看视频同走 `reward_video`。现在传 `signin_boost` 会落到「未知场景」分支(不发币,`status=unknown_scene`)
- 未知 `reward_scene`:不发币,记录 `status=unknown_scene`,返回 `is_verify=false/reason=1`
- 验签过但参数缺/坏或 user 不存在 → 不发(`is_verify=false` + `reason`;granted / capped / ecpm_missing / 业务不满足已记录 → `is_verify=true` + `reason=0`
+66
View File
@@ -0,0 +1,66 @@
# GET /api/v1/ad/reward-result/{ad_session_id} — 查本次广告的权威发奖结果 + 本轮累计
客户端看完激励视频后轮询本接口,拿**本次实发金币**和**本轮累计**用于「恭喜累计获得奖励」弹窗。不再用余额差 / `coin_per_ad` 估算。
**纯只读**:发奖仍只由验签过的 S2S 回调完成,本接口不写库、不产生任何奖励。按 `user_id` 收窄,被刷也只能查到自己的记录。
## 鉴权
需要 Bearer token。
## 路径参数
| 参数 | 类型 | 约束 | 说明 |
|---|---|---:|---|
| `ad_session_id` | string | 长度 8~64 | 本次广告会话 id,客户端生成,与 `mediaExtra` / `ecpm-report` 同值 |
## 响应
| 字段 | 类型 | 说明 |
|---|---|---|
| `ad_session_id` | string | 回显请求值 |
| `status` | string | `pending`(S2S 未到账,继续轮询) / `granted` / `capped`(当日超限) / `ecpm_missing` / `closed_early`(提前关闭) |
| `coin` | int \| null | **本条**实发金币。granted 为真实到账额;未发奖的状态为 0;pending 为 `null` |
| `round_coin` | int \| null | **本轮累计已发金币**(含本条) ← 弹窗显示的就是它 |
```json
{ "ad_session_id": "3f2a9c1b7e4d8a60", "status": "granted", "coin": 20, "round_coin": 60 }
```
### `round_coin` 的口径
「轮」= 用户点「去膨胀」到点「放弃赚钱」之间连看的若干条广告,边界由客户端的 `boost_round_id` 定(见 [ad-pangle-callback](./ad-pangle-callback.md))。
```sql
SELECT COALESCE(SUM(coin), 0) FROM ad_reward_record
WHERE user_id = :user_id -- 恒带,轮 id 是客户端生成的不可跨用户信任
AND boost_round_id = :该会话记录的 boost_round_id
AND status = 'granted'
```
由服务端求和而非客户端自己累加:客户端进程被杀 / 低内存重建后本地累计会丢,发奖记录不会。
**要守住的不变量:弹窗数字 == 本轮实际到账之和 == 用户看到的余额涨幅。** 三者对不上,用户就会认为少发了钱。
| 情形 | `round_coin` |
|---|---|
| 本条 `granted` | 本轮累计(含本条) |
| 本条 `capped` / `closed_early` / `ecpm_missing` | **仍返本轮累计**,该条按 0 计(撞上限那下的 toast 要能显示前几条的总额,不能是空) |
| `status=pending`(没记录) | `null` —— 连属于哪一轮都不知道。**不是 0**,0 会被读成「本轮没赚到」 |
| 该记录没有 `boost_round_id`(老客户端 / extra 丢失) | `null`,客户端退回只显示单条 `coin` |
## 错误
- `401`: 未登录
- `422`: `ad_session_id` 长度不在 8~64
**查不到记录不返 404**,而是 200 + `status="pending"`。404 只应表示路由不存在;两者混在一起客户端没法区分「后端没部署」和「再等等」。
## 实现注意
同一 `ad_session_id` 可能有多条记录,取值时**显式优先 `granted`**,不能只取最近一条:
- 客户端先报 `closed_early`、S2S 随后姗姗来迟 → 两条,`granted` 反而是后写的
- 本地联调重复调 `test-grant` → 同 session 多条 `granted``trans_id` 各不相同)
都没有 `granted` 才取最近一条,让客户端知道没发的原因。
+3 -2
View File
@@ -9,7 +9,8 @@
| 字段 | 类型 | 必填 | 默认 | 说明 |
|---|---|---|---|---|
| `reward_scene` | string | 否 | `reward_video` | 模拟发奖场景。`reward_video`=普通激励视频`signin_boost`=签到膨胀 |
| `reward_scene` | string | 否 | `reward_video` | 模拟发奖场景。当前**只支持** `reward_video``signin_boost`签到膨胀)已于 2026-07 下线,传它返 `422` |
| `boost_round_id` | string | 否 | `null` | 本次广告属于哪一轮膨胀。正式链路走 S2S 的 `mediaExtra`,本接口不经 S2S 拿不到 extra,故由 body 补。**不传的话 debug 包 `/ad/reward-result``round_coin` 恒为 `null`**,「弹窗 40 → 60 → toast +60」那套累计验收在本地跑不起来 |
| `ad_session_id` | string(8~64) \| null | 否 | null | 本次广告会话 id(与 [ecpm-report](./ad-ecpm-report.md) 同值)。**仅 `reward_video` 场景生效**:据此查回客户端已上报的真实 eCPM,走与正式发奖相同的公式发奖;查不到或 eCPM≤0(测试应用常返 0/假值)时兜底 200,保证本地联调仍出非零金币 |
## 出参
@@ -33,4 +34,4 @@
`reward_scene=reward_video` 时按上面 `ad_session_id` 查回的真实 eCPM 走金币公式发奖(取不到兜底 200)——便于本地用 [admin 金币审计](./admin-ad-coin-audit.md) 核对「看广告→金币」是否按公式计算。
`reward_scene=signin_boost` 时复用签到膨胀业务规则:必须当天已签到、非第 14 天、当天未膨胀过,成功后写入 `signin_boost` 金币流水。它让已登录客户端能自助发奖 = 绕过反作弊,**严禁在生产开启**。
它让已登录客户端能自助发奖 = 绕过反作弊,**严禁在生产开启**。
+2 -2
View File
@@ -37,8 +37,8 @@
| `feed_ad_watch_count` | int | 信息流广告有效完成视频数(`ad_feed_reward_record.status=granted`) |
| `signin_coin_total` | int | 签到累计发放金币(`biz_type=signin`) |
| `signin_count` | int | 签到次数(`signin_record`) |
| `signin_boost_coin_total` | int | 签到膨胀累计发放金币(`biz_type=signin_boost`) |
| `signin_boost_watch_count` | int | 签到膨胀有效视频数(`signin_boost_record`) |
| `signin_boost_coin_total` | int | **历史口径**:签到膨胀累计发放金币(`biz_type=signin_boost`)。功能已下线,数字不再增长,保留供对账 |
| `signin_boost_watch_count` | int | **历史口径**:签到膨胀次数。膨胀 2026-07 已下线、`signin_boost_record` 表已 drop,改数 `coin_transaction.biz_type='signin_boost'` 的入账笔数(一次膨胀 = 一笔,与原口径等价),只会停在历史值不再增长 |
**DashboardCash**
| 字段 | 类型 | 说明 |
-33
View File
@@ -1,33 +0,0 @@
# POST /api/v1/signin/boost — 签到后看广告膨胀金币
用户 Day1-Day13 当天已签到后,看完一条激励视频,由穿山甲 S2S 回调固定补发 2000 金币。本接口只用于 S2S 发奖后的确认。
## 鉴权
需要 Bearer token。
## 请求体
| 字段 | 类型 | 必填 | 说明 |
|---|---|---:|---|
| `ad_ref_id` | string | 是 | 穿山甲 S2S 回调的 `trans_id`。回调需先以 `extra.reward_scene=signin_boost` 完成发奖 |
## 响应
| 字段 | 类型 | 说明 |
|---|---|---|
| `coin_awarded` | int | 本次膨胀补发金币 |
| `coin_balance` | int | 补发后的金币余额 |
| `signin_date` | string | 被膨胀的签到日期,格式 `YYYY-MM-DD` |
## 错误
- `401`: 未登录
- `409`: 缺少/无效广告回调记录,非本人广告,回调未发奖,当天未签到,Day14,或当天已经膨胀过
## 数据写入
- 本接口不直接发奖;实际写入发生在 `/ad/pangle-callback``reward_scene=signin_boost` 分支。
- 回调写 `signin_boost_record` 新增一行,用 `(user_id, signin_date)` 唯一约束防重复。
- 回调使 `coin_account` 增加固定 `2000` 金币。
- 回调写入 `coin_transaction.biz_type=signin_boost`
+6 -7
View File
@@ -35,7 +35,7 @@
| 资产卡 / 钱包余额 | [`coin_account`](./coin_account.md) | 一用户一行的金币+现金余额快照 |
| 金币明细 | [`coin_transaction`](./coin_transaction.md) | 每次金币变动一笔流水 |
| 现金明细 | [`cash_transaction`](./cash_transaction.md) | 每次现金变动一笔流水(分) |
| 每日签到 | [`signin_record`](./signin_record.md) + [`signin_boost_record`](./signin_boost_record.md) | 7 天循环发币;签到后看广告可膨胀一次 |
| 每日签到 | [`signin_record`](./signin_record.md) | 7 天循环发币。签到弹窗的「看广告膨胀」2026-07 起走 `reward_video`(按 eCPM 发,记 `ad_reward_record`),不再有独立的膨胀表 |
| 一次性任务(开消息提醒等) | [`user_task`](./user_task.md) | 领一次发币 |
| 看激励视频赚金币 | [`ad_reward_record`](./ad_reward_record.md) + [`ad_watch_log`](./ad_watch_log.md) + [`ad_ecpm_record`](./ad_ecpm_record.md) | 独立数据流:发奖 / 旧版观看时长 / 收益对账 |
| 信息流/Draw 广告结算 | [`ad_feed_reward_record`](./ad_feed_reward_record.md) | 每展示满 10 秒累计一份奖励,完成后一次性入账;`ad_type`(feed/draw)+`feed_scene`(compare/coupon)分形态/场景 |
@@ -102,12 +102,11 @@
| 注销 `DELETE /user` | `user` | U(软删:`phone→deleted_<id>``status=deleted`) |
| 绑/解绑微信 `POST /wallet/bind-wechat``/unbind-wechat` | `user`.wechat_* | U |
| 签到 `POST /signin/do` | `signin_record`(C) + `coin_account`(U) + `coin_transaction`(C `signin`) | 同事务 |
| 签到膨胀 `POST /signin/boost` | `signin_boost_record`(C) + `coin_account`(U) + `coin_transaction`(C `signin_boost`) | 同事务;同日一次 |
| 领任务 `POST /tasks/claim` | `user_task`(C) + `coin_account`(U) + `coin_transaction`(C `task_<key>`) | 同事务 |
| 金币兑现金 `POST /wallet/exchange` | `coin_account`(U) + `coin_transaction`(C `exchange_out` ) + `cash_transaction`(C `exchange_in` +) | 同事务 |
| 发起提现 `POST /wallet/withdraw` | `withdraw_order`(C `reviewing`,记 `source`) + `coin_account`(U 按 source 扣对应余额) + 流水(C :`cash_transaction.withdraw``invite_cash_transaction.invite_withdraw`) | 同事务,**不打款**;#121`source` 分账 |
| 查提现状态 / 用户取消 `GET /wallet/withdraw/status` | `withdraw_order`(U) + 失败→对应账本退款流水(C `withdraw_refund` / `invite_withdraw_refund` +) | |
| 穿山甲发奖 S2S 回调 `POST /ad/pangle-callback` | `ad_reward_record`(C)+ granted→`coin_account`(U)+`coin_transaction`(C `reward_video`/`signin_boost`) | `trans_id` 幂等 |
| 穿山甲发奖 S2S 回调 `POST /ad/pangle-callback` | `ad_reward_record`(C)+ granted→`coin_account`(U)+`coin_transaction`(C `reward_video`) | `trans_id` 幂等 |
| 看广告时长上报 `POST /ad/watch-report` | `ad_watch_log`(C) | |
| 广告 eCPM 上报 `POST /ad/ecpm-report` | `ad_ecpm_record`(C) | |
| 信息流广告结算 `POST /ad/feed-reward` | `ad_feed_reward_record`(C)+ granted→`coin_account`(U)+`coin_transaction`(C `feed_ad_reward`) | `client_event_id` 幂等 |
@@ -175,7 +174,7 @@
## 三、表间关系 & Join Key
### 硬外键(数据库 FK 约束)
- **19 张用户维度表 `.user_id` → `user.id`**:`coin_account`(同时是 PK)、`coin_transaction``cash_transaction``invite_cash_transaction``withdraw_order``wechat_transfer_authorization`(同时是 PK)、`signin_record``signin_boost_record``user_task``comparison_record`(2026-07 起 `user_id` **可空**——harvest 帧0 建行时软鉴权可能拿不到)、`comparison_milestone_claim``savings_record``ad_reward_record``ad_watch_log``ad_ecpm_record``ad_feed_reward_record``price_report``feedback``device_liveness`
- **18 张用户维度表 `.user_id` → `user.id`**:`coin_account`(同时是 PK)、`coin_transaction``cash_transaction``invite_cash_transaction``withdraw_order``wechat_transfer_authorization`(同时是 PK)、`signin_record``user_task``comparison_record`(2026-07 起 `user_id` **可空**——harvest 帧0 建行时软鉴权可能拿不到)、`comparison_milestone_claim``savings_record``ad_reward_record``ad_watch_log``ad_ecpm_record``ad_feed_reward_record``price_report``feedback``device_liveness`
- `admin_audit_log.admin_id``admin_user.id`
- `price_report.comparison_record_id``comparison_record.id`(可空:关联记录被删后仍留上报历史)。
- **邀请两表**`user.id`:`invite_relation.inviter_user_id``invite_relation.invitee_user_id`(唯一)、`invite_fingerprint.inviter_user_id`——注意 FK 列名是 `inviter`/`invitee_user_id`,不是 `user_id`
@@ -187,7 +186,7 @@
| biz_type | ref_id 指向 | amount 符号 |
|---|---|---|
| `signin` | 当天日期串(= `signin_record.signin_date` 的 ISO `YYYY-MM-DD`) | + |
| `signin_boost` | 当天日期串(= `signin_boost_record.signin_date` 的 ISO `YYYY-MM-DD`) | + |
| `signin_boost`(**历史,2026-07 已下线**) | 当时的广告 `trans_id`,无则当天日期 ISO 串。不再产生新行,存量保留供对账 | + |
| `task_<key>` | `user_task.task_key` | + |
| `reward_video` / `ad_reward`(历史) | `ad_reward_record.trans_id` | + |
| `feed_ad_reward` | `ad_feed_reward_record.client_event_id` | + |
@@ -225,7 +224,7 @@
user ─1:1─ coin_account
user ─1:1─ wechat_transfer_authorization
user ─1:N─ { coin_transaction, cash_transaction, invite_cash_transaction, withdraw_order,
signin_record, signin_boost_record, user_task, comparison_record(user_id 可空),
signin_record, user_task, comparison_record(user_id 可空),
comparison_milestone_claim, savings_record, ad_reward_record, ad_watch_log,
ad_ecpm_record, ad_feed_reward_record, price_report, feedback, device_liveness }
(device_liveness 硬 FK; (user_id,device_id) 唯一)
@@ -255,7 +254,7 @@ launch_confirm_sample (独立, 无硬 FK; 都上报不去
1. **余额快照** `coin_account`:`coin_balance`(金币个数)+ `cash_balance_cents`(现金分)+ `invite_cash_balance_cents`(邀请奖励金分,#82),一用户一行,读取展示用。
2. **流水账本** `coin_transaction` / `cash_transaction` / `invite_cash_transaction`:每次变动写一笔,`balance_after*` 记变动后余额,可逐笔回溯对账。**现金与邀请奖励金是两本物理隔离的账**——发放口径与提现对账各自独立。
3. **唯一变动入口**:金币走 `repositories/wallet.grant_coins`,邀请奖励金走 `grant_invite_cash`——都是「更新快照 + 写流水,**不 commit**,由调用方同一事务 commit」。signin / signin_boost / task / ad_reward / feed_ad_reward / exchange / admin 走 `grant_coins`;`invite_reward` / admin 调整走 `grant_invite_cash`,靠 `biz_type` 区分来源。
3. **唯一变动入口**:金币走 `repositories/wallet.grant_coins`,邀请奖励金走 `grant_invite_cash`——都是「更新快照 + 写流水,**不 commit**,由调用方同一事务 commit」。signin / task / ad_reward / feed_ad_reward / exchange / admin 走 `grant_coins`(`signin_boost` 2026-07 已下线,存量流水保留);`invite_reward` / admin 调整走 `grant_invite_cash`,靠 `biz_type` 区分来源。
- **汇率**:`10000 金币 = 1 元 = 100 分`(`rewards.COIN_PER_YUAN`);兑换额必须是整分倍数。
- **提现状态机**:`reviewing`(发起即原子扣款、待人工审核、**不打款**)→ 审核通过 `pending`(微信转账在途)→ `success` / `failed`(失败自动退款);审核拒绝 `rejected`(退款)。**按 `withdraw_order.source` 分账**(#121):`coin_cash` 单的扣款/退款写 `cash_transaction`,`invite_cash` 单写 `invite_cash_transaction`;`out_bill_no` 幂等,孤儿 pending 单由 `reconcile_pending_withdraws` 对账兜底,admin `withdraws/ledger-check` 分账校验「单 ↔ 流水」。
-1
View File
@@ -35,7 +35,6 @@
| `withdraw_order` | 提现单(现金→微信零钱,含人工审核态;`source` 分账 coin_cash/invite_cash) | `models/wallet.py` | [详情](./withdraw_order.md) |
| `wechat_transfer_authorization` | 微信免确认转账授权(一用户一行) | `models/wallet.py` | [详情](./wechat_transfer_authorization.md) |
| `signin_record` | 签到记录(7 天循环) | `models/signin.py` | [详情](./signin_record.md) |
| `signin_boost_record` | 签到后看广告膨胀记录 | `models/signin.py` | [详情](./signin_boost_record.md) |
| `user_task` | 一次性任务领取去重 | `models/task.py` | [详情](./user_task.md) |
| `ad_reward_record` | 看激励视频发奖记录(S2S 回调,trans_id 幂等) | `models/ad_reward.py` | [详情](./ad_reward_record.md) |
| `ad_watch_log` | 看广告观看时长(旧版兼容字段) | `models/ad_watch_log.py` | [详情](./ad_watch_log.md) |
+8 -5
View File
@@ -2,7 +2,7 @@
> 模型 `app/models/ad_reward.py` · 仓库 `app/repositories/ad_reward.py` · 接口 [ad-pangle-callback](../api/ad-pangle-callback.md) / [ad-reward-status](../api/ad-reward-status.md) / [ad-test-grant](../api/ad-test-grant.md) · [← 索引](./README.md) · [总览](./OVERVIEW.md)
每条 = 穿山甲一次**服务端激励回调**。`trans_id` 唯一做幂等键(穿山甲会重试,同号只处理一次)。`reward_scene` 区分普通激励视频、签到膨胀等场景;`reward_date`(北京时间日期串)给普通激励视频"每日上限"计数用。
每条 = 穿山甲一次**服务端激励回调**。`trans_id` 唯一做幂等键(穿山甲会重试,同号只处理一次)。`reward_scene` 区分普通激励视频、提现看视频等场景;`reward_date`(北京时间日期串)给普通激励视频"每日上限"计数用。
## 用在哪 / 增删改查
- **C(插入)**:`POST /ad/pangle-callback`(穿山甲 S2S,经 SHA256 验签;`grant_ad_reward` 或场景业务处理)或 `POST /ad/test-grant`(本地联调)。普通激励视频三道闸:① 验签不过 → API 层 403,不进库;② `trans_id` 已存在 → 原样返回不重复发;③ **当日发奖次数(`DAILY_AD_REWARD_LIMIT`,默认 500)到顶** → 记一行 `status='capped'``coin=0`、不发币。否则按 eCPM 公式发币。另:`POST /ad/reward-noshow`(`record_reward_noshow`,Bearer)在用户提前关/未发奖时记一行 `status='closed_early'``coin=0` 留痕(同 session 已 granted 则跳过)。
@@ -13,15 +13,16 @@
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|---|---|---|---|
| `id` | Integer | PK, autoincrement | |
| `trans_id` | String(64) | UNIQUE, index, NOT NULL | 穿山甲交易号(幂等键)。**被 `coin_transaction.ref_id` 引用**(biz_type=reward_video/signin_boost 等)。`closed_early` 留痕记录无 S2S 交易号,用合成键 `noreward:{ad_session_id}` |
| `trans_id` | String(64) | UNIQUE, index, NOT NULL | 穿山甲交易号(幂等键)。**被 `coin_transaction.ref_id` 引用**(biz_type=reward_video 等)。`closed_early` 留痕记录无 S2S 交易号,用合成键 `noreward:{ad_session_id}` |
| `user_id` | Integer | FK→user.id, index, NOT NULL | 归属用户(回调 media_extra 带回;不存在抛 UnknownUserError) |
| `reward_scene` | String(32) | NOT NULL, default `reward_video` | 奖励场景:`reward_video` 普通激励视频;`signin_boost` 签到膨胀 |
| `reward_scene` | String(32) | NOT NULL, default `reward_video` | 奖励场景:`reward_video` 普通激励视频(当前唯一发币场景);`withdrawal_ad` 提现门槛视频(只留痕不发币);`signin_boost` **历史值,2026-07 已下线** |
| `ad_session_id` | String(64) | index, nullable | 客户端广告会话 ID,来自 `extra.ad_session_id`;用于匹配 `ad_ecpm_record` |
| `boost_round_id` | String(64) | nullable | 「这条广告属于哪一轮膨胀」,来自 `extra.boost_round_id`。一轮 = 用户点「去膨胀」到点「放弃赚钱」之间连看的若干条。**纯标签,不参与发奖判定**;仅供 `/ad/reward-result` 求和出 `round_coin`(弹窗显示的累计值)。老客户端 / extra 丢失时 NULL |
| `ecpm_raw` | String(32) | nullable | 本次发奖采用的 eCPM 原始值;可来自 S2S `ecpm` 或客户端上报 |
| `app_env` | String(16) | nullable | 来源应用 `prod`(傻瓜比价)/`test`(测试);S2S 不带,发奖时按 `ad_session_id` 匹配 `ad_ecpm_record` 回填,查不到 NULL。广告收益报表金币侧按它聚合 |
| `our_code_id` | String(64) | nullable | 我们配置的代码位 104xxx(同上回填) |
| `coin` | Integer | NOT NULL, default 0 | 实发金币;`capped`/`ecpm_missing`/`closed_early`/业务不满足时为 0 |
| `status` | String(16) | NOT NULL, default `granted` | 取值:`granted`(已发)/ `capped`(当日次数超限)/ `ecpm_missing`(缺 eCPM)/ `closed_early`(展示了但用户提前关/跳过,未发奖,客户端 reward-noshow 留痕)/ `not_signed`/`already_boosted`/`last_day`/ `unknown_scene`(回调 `reward_scene` 不在支持集合,只留痕不发) |
| `status` | String(16) | NOT NULL, default `granted` | 取值:`granted`(已发)/ `capped`(当日次数超限)/ `ecpm_missing`(缺 eCPM)/ `closed_early`(展示了但用户提前关/跳过,未发奖,客户端 reward-noshow 留痕)/ `unknown_scene`(回调 `reward_scene` 不在支持集合,只留痕不发) |
| `reward_date` | String(10) | index, NOT NULL | 北京时间日期串 `YYYY-MM-DD`,按它等值统计当日发奖次数 |
| `reward_name` | String(64) | nullable | 穿山甲上报奖励名(参考,不作发奖依据) |
| `raw` | String(1024) | nullable | 回调原始参数(审计排查) |
@@ -34,8 +35,10 @@
## 索引与约束
- PK `id`;UNIQUE+index `trans_id`;index `user_id``reward_date``created_at``ad_session_id`
- 复合 index `ix_ad_reward_user_boost_round` = `(user_id, boost_round_id)`:算「本轮累计已发」用。**求和恒带 `user_id`** —— `boost_round_id` 是客户端生成的,不带 `user_id` 等于让任何人拿别人的轮 id 查别人发了多少。
## 注意
- 普通激励视频按 eCPM 公式发奖;若 S2S 与客户端会话上报都缺 eCPM,记录 `status='ecpm_missing'``coin=0`,不发币。
- 签到膨胀复用本表记录 S2S 幂等,实发固定 `2000` 金币由 `signin_boost_record`/`coin_transaction.biz_type=signin_boost` 承载
- 签到膨胀(`reward_scene=signin_boost`)2026-07 已下线,存量行保留供对账;签到弹窗的「看广告膨胀」现与福利页看视频同走 `reward_video`(按 eCPM 公式发)
- **膨胀轮累计**:`SUM(coin) WHERE user_id=? AND boost_round_id=? AND status='granted'`,由 `/ad/reward-result` 返回为 `round_coin`。客户端就算一直复用同一个轮 id,也只是把展示数字滚大 —— 求和的是**已发生**的发奖记录,不产生任何新入账,无资损风险。
- 并发同 `trans_id` 撞唯一约束 → catch IntegrityError 回滚返回已存在那条(幂等兜底)。
+1 -1
View File
@@ -14,7 +14,7 @@
## 字段
| 列 | 类型 | 约束 / 默认 | 说明(取值 / join) |
|---|---|---|---|
| `key` | String(64) | **PK** | 配置标识,取值见 `config_schema.CONFIG_DEFS`:`signin_rewards` / `min_exchange_coin` / `withdraw_min_cents` / `withdraw_max_cents` / `task_rewards` / `record_milestones` / `ad_reward_coin` / `ad_daily_limit` / `ad_max_coin` / `ad_round_count` / `ad_cooldown_sec` / `signin_boost_coin` / `withdraw_auto_reconcile_enabled` / `comparing_ad_enabled` |
| `key` | String(64) | **PK** | 配置标识,取值见 `config_schema.CONFIG_DEFS`:`signin_rewards` / `min_exchange_coin` / `withdraw_min_cents` / `withdraw_max_cents` / `task_rewards` / `record_milestones` / `ad_reward_coin` / `ad_daily_limit` / `ad_max_coin` / `ad_round_count` / `ad_cooldown_sec` / `withdraw_auto_reconcile_enabled` / `comparing_ad_enabled` |
| `value` | JSON(PG: JSONB) | NOT NULL | 配置值,类型随 key(`int` / `int_list` 如签到 14 档 / `dict_str_int` 如 task_rewards / `bool` 如 withdraw_auto_reconcile_enabled / comparing_ad_enabled) |
| `updated_by_admin_id` | Integer | nullable | 最后修改的管理员 id(= `admin_user.id`,软引用,无 FK) |
| `updated_at` | DateTime(tz) | server_default now(), onupdate now() | 最后修改时间 |
+2 -2
View File
@@ -10,7 +10,7 @@
| 动作 / endpoint | `biz_type` | `amount` | `ref_id` 指向 |
|---|---|---|---|
| 签到 `POST /signin/do` | `signin` | + | 当天日期串(= `signin_record.signin_date` ISO) |
| 签到后看广告膨胀 `POST /signin/boost` | `signin_boost` | + | 广告 `trans_id`(= `signin_boost_record.ad_ref_id`);无 ad_ref_id 时回退当天日期 ISO 串 |
| ~~签到后看广告膨胀~~(**2026-07 已下线**) | `signin_boost` | + | 历史行:当时的广告 `trans_id`,无则当天日期 ISO 串。不再产生新行;签到弹窗的看广告改走 `reward_video` |
| 领任务 `POST /tasks/claim` | `task_<key>`(如 `task_enable_notification`) | + | 一次性任务=`user_task.task_key`;可重复任务(`enable_notification`)=带序号 `task_key:N` |
| 普通激励视频 S2S 回调 `POST /ad/pangle-callback` | `reward_video`(历史兼容:`ad_reward`) | + | `ad_reward_record.trans_id` |
| 信息流广告结算 `POST /ad/feed-reward` | `feed_ad_reward` | + | `ad_feed_reward_record.client_event_id` |
@@ -36,7 +36,7 @@
## 关系 / Join Key
- `user_id``user.id`(多对一)。
- `ref_id` 是**软关联**(无 FK),目标随 `biz_type`:`signin`→签到日(`signin_record.signin_date` ISO) / `signin_boost``signin_boost_record.ad_ref_id`(无则当天日期) / `task_<key>`→一次性任务=`user_task.task_key`、可重复任务=`task_key:N` / `reward_video`/`ad_reward``ad_reward_record.trans_id` / `feed_ad_reward``ad_feed_reward_record.client_event_id` / 其余 null。
- `ref_id` 是**软关联**(无 FK),目标随 `biz_type`:`signin`→签到日(`signin_record.signin_date` ISO) / `signin_boost`(历史)→当时的广告 `trans_id`(无则当天日期) / `task_<key>`→一次性任务=`user_task.task_key`、可重复任务=`task_key:N` / `reward_video`/`ad_reward``ad_reward_record.trans_id` / `feed_ad_reward``ad_feed_reward_record.client_event_id` / 其余 null。
## 索引与约束
- PK `id`;index `user_id``created_at`
-22
View File
@@ -1,22 +0,0 @@
# signin_boost_record — 签到膨胀记录
App 用户当天签到后,看完激励视频可固定膨胀一次(默认 3000 金币,`rewards.SIGNIN_BOOST_COIN`,运营后台 `app_config.signin_boost_coin` 可改)。循环最后一天(`cycle_day == SIGNIN_CYCLE_LEN`,即 7 天循环的第 7 天)不展示也不允许膨胀。本表记录膨胀动作,并用唯一约束防重复补发。
## 字段
| 字段 | 类型 | 约束 | 说明 |
|---|---|---|---|
| `id` | Integer | PK | 自增主键 |
| `user_id` | Integer | FK → `user.id`, index, NOT NULL | 用户 |
| `signin_date` | Date | NOT NULL | 被膨胀的签到日期,北京时间 |
| `coin_awarded` | Integer | NOT NULL | 本次补发金币,默认固定 3000(`rewards.get_signin_boost_coin`) |
| `ad_ref_id` | String(64) | nullable | 穿山甲 S2S 回调 `trans_id` |
| `created_at` | DateTime(tz) | server_default now(), NOT NULL | 创建时间 |
## 约束
- `UNIQUE(user_id, signin_date)` = `uq_signin_boost_user_date`:同一用户同一天只能膨胀一次。
## 关联
- 膨胀成功时写 `coin_transaction.biz_type=signin_boost``ref_id = ad_ref_id`(无 ad_ref_id 时回退当天日期 ISO 串)。
+1 -1
View File
@@ -29,7 +29,7 @@
| `last_login_at` | DateTime(tz) | 应用层 default utcnow | 最近登录时间(每次登录更新) |
## 关系 / Join Key
- **被引用方(本表是 1,对方是 N/1)**:`coin_account``coin_transaction``cash_transaction``withdraw_order``wechat_transfer_authorization``signin_record``signin_boost_record``user_task``comparison_record``comparison_milestone_claim``savings_record``ad_reward_record``ad_watch_log``ad_ecpm_record``ad_feed_reward_record``price_report``feedback``user_id` 均 → `user.id`;`invite_relation``inviter_user_id` / `invitee_user_id` 均 → `user.id`
- **被引用方(本表是 1,对方是 N/1)**:`coin_account``coin_transaction``cash_transaction``withdraw_order``wechat_transfer_authorization``signin_record``user_task``comparison_record``comparison_milestone_claim``savings_record``ad_reward_record``ad_watch_log``ad_ecpm_record``ad_feed_reward_record``price_report``feedback``user_id` 均 → `user.id`;`invite_relation``inviter_user_id` / `invitee_user_id` 均 → `user.id`
- 与 `admin_user` **无任何关联**(C 端用户 vs 后台管理员,两套体系)。
## 索引与约束
+1 -14
View File
@@ -27,7 +27,6 @@
| 钱包 | `withdraw_order` | 提现单 | 现金→微信零钱提现 |
| 钱包 | `wechat_transfer_authorization` | 微信转账免确认授权 | 用户授权后转账免逐笔确认 |
| 激励 | `signin_record` | 签到记录 | 每日签到 |
| 激励 | `signin_boost_record` | 签到膨胀记录 | 签到后看广告翻倍补发 |
| 激励 | `user_task` | 一次性任务完成 | 只能领一次的任务 |
| 激励 | `comparison_milestone_claim` | 比价战绩领取 | 比价次数里程碑奖励 |
| 比价 | `comparison_record` | 比价记录 | 用户视角「我的比价记录」 |
@@ -181,18 +180,6 @@ App 用户主表。两种登录(极光一键 / 短信验证码)都映射到
| coin_awarded | 整数 | 本次发放金币 |
| created_at | 时间 | 时间 |
## `signin_boost_record` — 签到膨胀记录
签到后看广告「膨胀」翻倍,一天最多一次,补发金额=当天签到原始奖励。
| 字段 | 类型 | 说明 |
|---|---|---|
| id | 整数 | **主键** |
| user_id | 整数 | **外键→user** |
| signin_date | 日期 | 与 user_id 组成**唯一**(防并发重复补发) |
| coin_awarded | 整数 | 补发金币 |
| ad_ref_id | 字符串 | 广告会话/交易号,可空 |
| created_at | 时间 | 时间 |
## `user_task` — 一次性任务完成记录
@@ -379,7 +366,7 @@ profile「累计帮你省了」「省钱战绩」的唯一数据源。(user_id,
| user_id | 整数 | **外键→user** |
| coin | 整数 | 实发金币(超限为 0) |
| status | 字符串 | granted(已发)/ capped(当日超限)/ ecpm_missing(缺 eCPM |
| reward_scene | 字符串 | reward_video(福利页看视频/ signin_boost(签到膨胀 |
| reward_scene | 字符串 | reward_video(福利页看视频,当前唯一发币场景)/ withdrawal_ad(提现门槛视频,不发币)/ signin_boost**历史值,2026-07 已下线** |
| ad_session_id | 字符串 | 广告会话 id,可空 |
| ecpm_raw | 字符串 | 本次发奖采用的 eCPM 原始值,可空 |
| app_env | 字符串 | 应用环境 prod/test(回填),可空 |
@@ -70,7 +70,7 @@
## C. 部署 + 包名
- [ ] **后端部署到公网**(由服务器管理员;`/opt/shaguabijia-app-server`,uvicorn 127.0.0.1:8770,nginx 反代)
- [ ] **跑迁移**:`alembic upgrade head`(包含 `ad_reward_record``signin_boost_record``ad_feed_reward_record` 等表)
- [ ] **跑迁移**:`alembic upgrade head`(包含 `ad_reward_record``ad_feed_reward_record` 等表)
- [ ] **包名定稿**:当前 `com.jishisongfu.shaguabijia`。穿山甲(APP_ID 5830519)、极光、微信都绑"包名 + 签名",定了再上,别再换
- 微信提现链路当前因复用 elderhelper 的 appid + 包名切换已 dead,要恢复需申请傻瓜比价自己的微信 appid(另见客户端 build.gradle 注释)
- [ ] (可选,提升真实填充)集成 **MSA OAID SDK**:申请证书(绑包名、审核几天)。App 侧当前 `getDevOaid=null`,有 OAID 后投放匹配 + 填充会明显改善
+138
View File
@@ -0,0 +1,138 @@
"""重置指定用户的新手引导完成标记,让这个账号重新进新手引导页,方便反复测试引导流程。
原理:是否跳过引导只由 onboarding_completion 表里 (user_id, device_id) 那一行决定
( app/models/onboarding.py)删掉该用户的行 登录响应 onboarding_completed=false
GET /api/v1/user/onboarding/status 也返 false 客户端下次登录/启动重走引导
本地 SharedPreferences 标记卸载即丢以后端为准,所以删这一行就够,不用重装 App
默认删该用户**所有设备**的记录(换机/多设备一起放开);只想放开某一台用 --device-id
(device_id = 客户端硬件级 ANDROID_ID,与登录 / onboarding/complete 传的是同一个值)
与已有两个入口的分工:
- admin设备维度引导管理**设备**重置(该设备上所有账号一起),本脚本按**账号**;
- POST /api/v1/user/onboarding/reset 要客户端自己带 device_id ,本脚本从库里反查设备
用法(在项目根 pip install -e . 的环境里跑):
python scripts/reset_onboarding.py # 默认测试号 11111111111
python scripts/reset_onboarding.py 13800138000 # 指定手机号
python scripts/reset_onboarding.py --user-id 5 # 直接指定 user_id
python scripts/reset_onboarding.py --dry-run # 预览(照常执行再回滚),不落库
python scripts/reset_onboarding.py --device-id abc123 # 只放开这一台设备,其余设备照旧跳过
SessionLocal DATABASE_URL(SQLite / Postgres 都行),因此**默认只允许 APP_ENV=dev 改库**
(--dry-run 只读,任何环境都能跑)线上确实要给某个用户开引导时加 --force 这张表只存
"引导走过没"的标记,删了最坏结果是用户多看一次引导,不涉及金额/账目
"""
from __future__ import annotations
import argparse
import sys
from sqlalchemy import delete, select
from app.core.config import settings
from app.db.session import SessionLocal, engine
from app.models.onboarding import OnboardingCompletion
from app.models.user import User
# Windows 控制台默认 GBK,强制 UTF-8 否则中文输出乱码。stderr 也要设:
# SystemExit(如"用户不存在")的中文提示走的是 stderr。
for _stream in (sys.stdout, sys.stderr):
if hasattr(_stream, "reconfigure"):
_stream.reconfigure(encoding="utf-8")
# dev 下 engine 是 echo=True(APP_DEBUG),几十行 SQL 会把前后对比刷没。echo 走 SQLAlchemy 自己的
# InstanceLogger,不吃 logging.setLevel,只能改 engine.echo。
engine.echo = False
DEFAULT_PHONE = "11111111111"
def resolve_user(db, phone: str, user_id: int | None) -> User:
if user_id is not None:
user = db.get(User, user_id)
if user is None:
raise SystemExit(f"user_id={user_id} 不存在")
return user
user = db.execute(select(User).where(User.phone == phone)).scalar_one_or_none()
if user is None:
raise SystemExit(f"手机号 {phone} 没有对应用户(注意 phone 才是登录账号,username 是展示 ID)")
return user
def print_state(db, user: User, device_id: str | None, label: str) -> None:
"""打印该用户当前的引导完成标记。--device-id 时只看那一台,便于确认没误伤别的设备。"""
stmt = (
select(OnboardingCompletion.device_id, OnboardingCompletion.completed_at)
.where(OnboardingCompletion.user_id == user.id)
.order_by(OnboardingCompletion.completed_at.desc())
)
if device_id:
stmt = stmt.where(OnboardingCompletion.device_id == device_id)
rows = db.execute(stmt).all()
print(f"--- {label} ---")
scope = f"device_id={device_id}" if device_id else "全部设备"
if not rows:
print(f" onboarding_completion({scope}): (无) → 该用户会走引导")
return
print(f" onboarding_completion({scope}): {len(rows)} 条 → 这些设备上会跳过引导")
for did, at in rows:
print(f" device_id={did} 完成于 {at}")
def main() -> None:
parser = argparse.ArgumentParser(description="重置指定用户的新手引导,让其重新进引导页")
parser.add_argument("phone", nargs="?", default=DEFAULT_PHONE,
help=f"手机号(默认 {DEFAULT_PHONE})")
parser.add_argument("--user-id", type=int, default=None, help="直接按 user_id 定位,优先于 phone")
parser.add_argument("--device-id", default=None,
help="只重置这一台设备(硬件级 ANDROID_ID);默认重置该用户所有设备")
parser.add_argument("--dry-run", action="store_true", help="预览,最后回滚不落库")
parser.add_argument("--force", action="store_true", help="非 dev 环境也允许改库(仅删引导标记,不涉及账目)")
args = parser.parse_args()
if not args.dry_run and settings.APP_ENV != "dev" and not args.force:
raise SystemExit(
f"APP_ENV={settings.APP_ENV},默认只有 dev 能改库。确认要在该环境重置请加 --force"
"(--dry-run 只读,任意环境可跑)"
)
db = SessionLocal()
try:
user = resolve_user(db, args.phone, args.user_id)
print(f"DB: {settings.DATABASE_URL} APP_ENV: {settings.APP_ENV}")
print(f"用户: id={user.id} phone={user.phone} username={user.username}")
print_state(db, user, args.device_id, "before")
stmt = delete(OnboardingCompletion).where(OnboardingCompletion.user_id == user.id)
if args.device_id:
stmt = stmt.where(OnboardingCompletion.device_id == args.device_id)
deleted = db.execute(stmt).rowcount
if not deleted:
# 没记录本来就会走引导 —— 常见于:换了新设备、或运营/客户端已经重置过一次。
print("该用户(该范围内)本来就没有完成标记,已经会走引导了,无需处理。")
db.rollback()
else:
print_state(db, user, args.device_id, "after")
if args.dry_run:
db.rollback()
print(f"(dry-run:以上 after 为预览,已回滚,库没动;真跑会删 {deleted} 条)")
return
db.commit()
print(f"完成:删掉 {deleted} 条完成标记,{user.phone} 下次登录会重走新手引导。")
# 测试号无论库里有没有记录都恒走引导(见 app/core/test_account.py),提醒一句免得白跑
if settings.test_account_phone and user.phone == settings.test_account_phone:
print(f"提示:{user.phone} 是配置的测试账号(TEST_ACCOUNT_PHONE),"
"登录响应 onboarding_completed 恒为 false、本就每次都走引导,无需重置。")
print("提醒:客户端是在登录响应 / 启动时查 onboarding/status 的,已经在首页的 App 不会自动跳转,"
"退出登录重进(或杀掉重开)才会看到引导页。")
finally:
db.close()
if __name__ == "__main__":
main()
+218
View File
@@ -0,0 +1,218 @@
"""撤销指定用户「今天已签到」的状态,让今天可以重新签到,方便反复测试签到流程。
reset_signin.py 的区别:那个删**全部**签到历史(连续天数从头再来);本脚本只精确撤销
**今天**这一次,昨天及以前的记录原样保留 所以重签后 cycle_day / streak 会接着昨天继续,
7 天循环的档位不会被打乱,可以连着好几天测 N 的奖励
默认是**完整撤销**(等于今天这次签到从没发生过):
1. signin_record 今天这行 今天变回未签到
2. 删今天的 signin 金币流水,并把金币从 coin_account 余额 / 累计收益里扣回
金币默认要退:签到流水**没有**唯一索引拦重复(ux_coin_transaction_task_ref 只覆盖
biz_type LIKE 'task%'),不退的话每测一轮余额就白涨一次奖励,coin_transaction 里还会堆出
同一 ref_id(日期)的重复流水,收益明细页会看到两条今天的签到真想留着奖励用 --keep-coins
例外:签到的金币若已被兑换成现金(余额已不够退),**自动跳过退款**并保留今天的签到流水
因为 coin_balance 必须恒等于流水总和,硬退会把余额退成负数 夹到 0 又会吃掉别处赚的金币,
两种做法都会让账对不上这时重签会再发一次奖励,余额多涨一档,属可接受的测试噪音
用法(在项目根 pip install -e . 的环境里跑):
python scripts/reset_signin_today.py # 默认测试号 11111111111
python scripts/reset_signin_today.py 13800138000 # 指定手机号
python scripts/reset_signin_today.py --user-id 5 # 直接指定 user_id
python scripts/reset_signin_today.py --dry-run # 预览(照常执行再回滚),不落库
python scripts/reset_signin_today.py --keep-coins # 只删签到记录,保留已发金币
今天直接复用 app.core.rewards.cn_today(北京时间),与签到判重同源,不自己算时区
SessionLocal DATABASE_URL(SQLite / Postgres 都行),因此**只允许 APP_ENV=dev 时改库**
(--dry-run 只读,任何环境都能跑)
"""
from __future__ import annotations
import argparse
import sys
from sqlalchemy import select
from app.core.config import settings
from app.core.rewards import cn_today
from app.db.session import SessionLocal, engine
from app.models.signin import SigninRecord
from app.models.user import User
from app.models.wallet import CoinAccount, CoinTransaction
from app.repositories import signin as crud_signin
# Windows 控制台默认 GBK,强制 UTF-8 否则中文输出乱码
if hasattr(sys.stdout, "reconfigure"):
sys.stdout.reconfigure(encoding="utf-8")
# dev 下 engine 是 echo=True(APP_DEBUG),几十行 SQL 会把前后对比刷没。echo 走 SQLAlchemy 自己的
# InstanceLogger,不吃 logging.setLevel,只能改 engine.echo。
engine.echo = False
DEFAULT_PHONE = "11111111111"
def resolve_user(db, phone: str, user_id: int | None) -> User:
if user_id is not None:
user = db.get(User, user_id)
if user is None:
raise SystemExit(f"user_id={user_id} 不存在")
return user
user = db.execute(select(User).where(User.phone == phone)).scalar_one_or_none()
if user is None:
raise SystemExit(f"手机号 {phone} 没有对应用户(注意 phone 才是登录账号,username 是展示 ID)")
return user
def print_state(db, user: User, today, label: str) -> None:
print(f"--- {label} ---")
rec = db.execute(
select(SigninRecord).where(
SigninRecord.user_id == user.id, SigninRecord.signin_date == today
)
).scalar_one_or_none()
if rec is None:
print(f" signin_record {today}: (无)")
else:
print(f" signin_record {today}: 第{rec.cycle_day}档 连续{rec.streak}天 +{rec.coin_awarded}金币")
last = db.execute(
select(SigninRecord.signin_date)
.where(SigninRecord.user_id == user.id)
.order_by(SigninRecord.signin_date.desc())
.limit(1)
).scalar_one_or_none()
print(f" 最近一次签到: {last or '(从未签到)'}")
rows = db.execute(
select(CoinTransaction).where(
CoinTransaction.user_id == user.id,
CoinTransaction.biz_type == "signin",
CoinTransaction.ref_id == today.isoformat(),
)
).scalars().all()
print(f" coin_transaction(signin, 今天): {len(rows)} 条 / {sum(r.amount for r in rows)} 金币")
acc = db.get(CoinAccount, user.id)
if acc is None:
print(" coin_account: (无)")
else:
print(f" coin_account: coin={acc.coin_balance} earned={acc.total_coin_earned}")
# 用 App 自己的 get_status 复核,而不是脚本里重算一遍规则 —— 这行就是客户端会看到的
st = crud_signin.get_status(db, user.id)
print(f" [签到接口] can_claim={st.can_claim} today_signed={st.today_signed} "
f"今天第{st.today_cycle_day}档({st.today_coin}金币) 已连续{st.consecutive_days}")
def refund_today(db, user_id: int, today) -> None:
"""退回今天签到(含膨胀)发的金币:删流水 + 扣余额。
不变量:coin_balance 必须恒等于流水总和所以余额不够退时**整笔跳过**,而不是硬退成
负数或夹到 0 夹到 0 会吃掉用户在别处赚的金币,两种做法都会让余额和流水对不上
"""
rows = list(db.execute(
select(CoinTransaction).where(
CoinTransaction.user_id == user_id,
CoinTransaction.biz_type == "signin",
CoinTransaction.ref_id == today.isoformat(),
)
).scalars().all())
if not rows:
return
acc = db.get(CoinAccount, user_id)
if acc is None:
return
# 从最近一笔往回退,退到余额兜不住为止:正常情况下今天只有一笔,整笔退掉 = 干净的撤销。
# 少数情况今天堆了多笔(上一轮测试时金币已被兑换、退不掉而留下的),这样也能保证
# 「本轮新发的那笔」一定被退掉 —— 否则每测一轮余额就永久多涨一档。
rows.sort(key=lambda r: r.id, reverse=True)
refundable: list[CoinTransaction] = []
total = 0
for r in rows:
if total + r.amount > acc.coin_balance:
break
refundable.append(r)
total += r.amount
for r in refundable:
db.delete(r)
if total:
acc.coin_balance -= total
acc.total_coin_earned = max(0, acc.total_coin_earned - total)
print(f" 已退回 {total} 金币({len(refundable)}/{len(rows)} 笔)")
stuck = len(rows) - len(refundable)
if stuck:
# 典型场景:签完就把金币兑换成现金了(exchange_out),这笔奖励已经变成 cash_balance_cents,
# 余额里已经没有它了。硬退会把余额退成负数 / 夹到 0 又会吃掉别处赚的金币,两者都会让账对不上。
print(f" ⚠️ 还有 {stuck} 笔今天的签到流水退不掉(金币已被兑换/花掉,余额 {acc.coin_balance} 兜不住),"
f"原样保留 —— 硬退会让余额和流水总和对不上。")
print(" → 收益明细今天会多出几条签到记录,不影响签到功能测试;想彻底清干净用 reset_signin.py --with-coins。")
def main() -> None:
parser = argparse.ArgumentParser(description="撤销用户今天的签到,让今天能重新签")
parser.add_argument("phone", nargs="?", default=DEFAULT_PHONE,
help=f"手机号(默认 {DEFAULT_PHONE})")
parser.add_argument("--user-id", type=int, default=None, help="直接按 user_id 定位,优先于 phone")
parser.add_argument("--keep-coins", action="store_true",
help="不退已发金币(余额会越测越高,且留下重复流水)")
parser.add_argument("--dry-run", action="store_true", help="预览,最后回滚不落库")
args = parser.parse_args()
if not args.dry_run and settings.APP_ENV != "dev":
raise SystemExit(f"APP_ENV={settings.APP_ENV},拒绝改库(只有 dev 能改;--dry-run 可任意环境)")
today = cn_today()
db = SessionLocal()
try:
user = resolve_user(db, args.phone, args.user_id)
print(f"DB: {settings.DATABASE_URL}")
print(f"用户: id={user.id} phone={user.phone} 今天(北京): {today} keep_coins: {args.keep_coins}")
print_state(db, user, today, "before")
# 今天的签到记录 —— 只删今天,昨天及以前保留,重签后 streak 接着涨
rec = db.execute(
select(SigninRecord).where(
SigninRecord.user_id == user.id, SigninRecord.signin_date == today
)
).scalar_one_or_none()
if rec is not None:
db.delete(rec)
if rec is None:
print("今天本来就没签到,无需处理。")
db.rollback()
return
# 退金币
if args.keep_coins:
print("(--keep-coins:保留已发金币,流水和余额不动)")
else:
refund_today(db, user.id, today)
# 注:更早流水的 balance_after 是当时的快照,不回改 —— 收益明细里历史行的
# 余额列会与现余额对不上,dev 测试库无妨。
# SessionLocal 是 autoflush=False,不 flush 的话下面 print_state 的 select
# 读到的还是删之前的旧行,"after" 会骗人
db.flush()
print_state(db, user, today, "after")
if args.dry_run:
db.rollback()
print("(dry-run:以上 after 为预览,已回滚,库没动)")
return
db.commit()
print(f"完成:{user.phone} 今天({today})可以重新签到了。"
f"提醒:App 内存状态不会自动同步,杀掉重进福利页(当天未签到)会重新自动弹签到弹窗。")
finally:
db.close()
if __name__ == "__main__":
main()
-160
View File
@@ -1,160 +0,0 @@
"""会话存档【独立进程】轮询 poller —— 单独 systemd 服务跑, 与 app-server 隔离。
为什么独立进程:WeWorkFinanceSdk Go 编译的 c-shared , 自带 Go runtime嵌进 app-server(uvicorn)
后台线程里跑会 segfault把主进程带崩(2026-07-17 实测 core-dump 5s 循环)独立进程让 Go runtime
独占进程 + 主线程, 崩了 systemd 只重启本 poller不动 app-server拉到美团卡片(weapp)/截图(image)
POST app-server 内部端点 /internal/wx-finance/pending 打比价信号(下游心跳弹窗比价 完全复用)
手动跑: cd /opt/shaguabijia-app-server && .venv/bin/python scripts/wx_finance_poller.py
生产: deploy/wx-finance-poller.service 常驻(systemctl enable --now wx-finance-poller)
"""
from __future__ import annotations
import json
import logging
import signal
import time
from pathlib import Path
import httpx
from app.core.config import settings
from app.core.logging import setup_logging
from app.integrations.wx_finance_sdk import WxFinanceError, WxFinanceSdk
logger = logging.getLogger("shagua.wx_finance_poller")
_stop = False
def _on_signal(_signum: int, _frame: object) -> None:
global _stop
_stop = True
def _load_seq() -> int:
try:
return int(Path(settings.WX_FINANCE_SEQ_FILE).read_text().strip() or "0")
except Exception:
return 0
def _save_seq(seq: int) -> None:
try:
p = Path(settings.WX_FINANCE_SEQ_FILE)
p.parent.mkdir(parents=True, exist_ok=True)
p.write_text(str(seq))
except Exception:
logger.exception("seq 持久化失败")
def _notify(client: httpx.Client, source: str, kind: str) -> None:
"""POST 内部端点让 app-server 打信号(set_pending 在 app-server 进程内存)。"""
url = f"{settings.WX_FINANCE_INTERNAL_URL.rstrip('/')}/internal/wx-finance/pending"
try:
resp = client.post(
url,
json={"source": source, "kind": kind},
headers={"X-Internal-Secret": settings.INTERNAL_API_SECRET},
timeout=10,
)
logger.info("notify app-server (%s) → %s", kind, resp.status_code)
except Exception:
logger.exception("notify app-server 失败")
def _resolve_source(msg: dict) -> str:
"""从 weapp 小程序卡判源平台代号(= app-server _SOURCE_PACKAGES 的 key)。
会话存档 weapp 结构里 username(gh_id)/appid/pagepath/title 任一含平台特征即判定:
先按关键词启发式(覆盖大多数);认不出 默认 meituan(过渡期不阻断)+ WARNING,
并把整个 weapp 结构打进日志,便于照真实 username/appid 精确补规则
"""
weapp = msg.get("weapp") or {}
blob = json.dumps(weapp, ensure_ascii=False).lower()
# 排障日志:打全 weapp 结构(顺带 msg 顶层 key,防结构不在 weapp 下),便于补/改判定规则
logger.info("weapp 判源: msg_keys=%s weapp=%s", list(msg.keys()), blob[:1500])
if any(k in blob for k in ("jingdong", "jd.com", "jddj", "京东")):
return "jd"
if any(k in blob for k in ("meituan", "sankuai", "美团")):
return "meituan"
logger.warning("weapp 源平台未识别, 暂默认 meituan;见上一行结构日志补 _resolve_source 规则")
return "meituan"
def _handle(client: httpx.Client, msg: dict) -> None:
"""明文消息:外部用户(from≠接收成员)发来的 weapp/image → 通知 app-server 打信号。"""
msgtype = msg.get("msgtype") or ""
frm = msg.get("from") or ""
receiver = settings.WX_FINANCE_RECEIVER_USERID
if receiver and frm == receiver:
return # 成员自己发的, 跳过
if msgtype not in ("weapp", "image"):
return
# weapp 卡片按小程序判源平台;image(截图)判不出平台,沿用 meituan。
source = _resolve_source(msg) if msgtype == "weapp" else "meituan"
logger.info("命中触发 from=%s type=%s source=%s", frm, msgtype, source)
# 归档延迟排查: msgtime 是腾讯给的消息发送时刻(毫秒), 与 poller 此刻之差
# ≈ 企业微信会话存档的归档延迟(+ 本地 ≤0.5s 轮询间隔), 用来隔离"是我们慢还是归档慢"。
try:
_msgtime_ms = float(msg.get("msgtime") or 0)
if _msgtime_ms > 0:
logger.info("归档延迟(msgtime→poller取到) ≈ %.2fs", time.time() - _msgtime_ms / 1000.0)
except Exception:
pass
_notify(client, source, msgtype)
def main() -> None:
setup_logging(debug=settings.APP_DEBUG)
if not settings.wx_finance_configured:
logger.warning("WX_FINANCE 未配齐(ENABLED/corpid/存档Secret), poller 退出")
return
if not settings.INTERNAL_API_SECRET:
logger.error("INTERNAL_API_SECRET 未配置, poller 无法通知 app-server, 退出")
return
signal.signal(signal.SIGTERM, _on_signal)
signal.signal(signal.SIGINT, _on_signal)
sdk = WxFinanceSdk(
settings.WX_FINANCE_SDK_PATH,
settings.WX_KF_CORP_ID,
settings.WX_FINANCE_SECRET,
Path(settings.WX_FINANCE_PRIVATE_KEY_PATH).read_bytes(),
)
seq = _load_seq()
logger.info(
"wx_finance poller 启动, 从 seq=%d 轮询, 间隔 %ss", seq, settings.WX_FINANCE_POLL_INTERVAL_SEC
)
with httpx.Client() as client:
while not _stop:
try:
msgs = sdk.get_chat_data(seq, limit=1000)
if msgs:
logger.info("拉到 %d 条(seq>%d)", len(msgs), seq)
for item in msgs:
try:
plain = sdk.decrypt(item["encrypt_random_key"], item["encrypt_chat_msg"])
_handle(client, plain)
except Exception:
logger.exception("解密/处理单条失败 seq=%s", item.get("seq"))
seq = max(seq, int(item.get("seq", seq)))
if msgs:
_save_seq(seq)
except WxFinanceError:
logger.exception("GetChatData 出错")
except Exception:
logger.exception("轮询异常")
# 可打断的 sleep:按 ≤0.5s 分片, 支持亚秒间隔(如 0.5s), 收到 SIGTERM 尽快退
interval = max(0.1, settings.WX_FINANCE_POLL_INTERVAL_SEC)
step = min(0.5, interval)
slept = 0.0
while slept < interval and not _stop:
time.sleep(step)
slept += step
sdk.close()
logger.info("wx_finance poller 已停止")
if __name__ == "__main__":
main()
-45
View File
@@ -1,45 +0,0 @@
"""会话存档 SDK 连通性探针 —— 上 worker 前,在【云服务器】单独验证整条链路是否打通。
/opt/shaguabijia-app-server 下跑: .venv/bin/python scripts/wx_finance_probe.py
依次验证:加载 .so NewSdk/Init(corpid+存档Secret+可信IP) GetChatData(0) 试解第一条
只读不写任何比价信号任一步报错都会明确指出卡在哪(便于逐项排:.so 缺依赖 / Init 失败 /
可信IP 没放行 / 私钥不匹配)全绿了再把 WX_FINANCE_ENABLED true worker
"""
from __future__ import annotations
from pathlib import Path
from app.core.config import settings
from app.integrations.wx_finance_sdk import WxFinanceSdk
def main() -> None:
print(f"[cfg] sdk_path = {settings.WX_FINANCE_SDK_PATH}")
print(f"[cfg] corpid = {settings.WX_KF_CORP_ID}")
print(f"[cfg] secret_set = {bool(settings.WX_FINANCE_SECRET)}")
print(f"[cfg] privkey = {settings.WX_FINANCE_PRIVATE_KEY_PATH}")
sdk = WxFinanceSdk(
settings.WX_FINANCE_SDK_PATH,
settings.WX_KF_CORP_ID,
settings.WX_FINANCE_SECRET,
Path(settings.WX_FINANCE_PRIVATE_KEY_PATH).read_bytes(),
)
print("[1] 加载 .so + Init 成功")
msgs = sdk.get_chat_data(0, limit=100)
print(f"[2] GetChatData(0) 成功, 拉到 {len(msgs)}")
if msgs:
first = msgs[0]
print(f" 第一条 seq={first.get('seq')} publickey_ver={first.get('publickey_ver')}")
plain = sdk.decrypt(first["encrypt_random_key"], first["encrypt_chat_msg"])
print(f"[3] 解密成功, msgtype={plain.get('msgtype')} from={plain.get('from')}")
else:
print("[3] 暂无历史消息(先让用户加接收成员好友、发条消息, 再跑一次)")
sdk.close()
print("探针完成:链路全通 ✓")
if __name__ == "__main__":
main()
+269
View File
@@ -356,3 +356,272 @@ def test_callback_disabled_returns_503(client, monkeypatch) -> None:
# 503 发生在验签/发奖之前,不需要真实用户
r = _callback(client, _signed(1, "trans_disabled"))
assert r.status_code == 503, r.text
# ===== 按 ad_session_id 查权威发奖结果(GET /reward-result/{ad_session_id})=====
# 客户端看完广告轮询它拿弹窗金额:只认 status='granted' 且 coin>0,其余一律不弹。
def _reward_result(client, token: str, session_id: str):
return client.get(f"/api/v1/ad/reward-result/{session_id}", headers=_auth(token))
def _session_extra(session_id: str, **kv: str) -> str:
return json.dumps({"ad_session_id": session_id, **kv})
def test_reward_result_pending_when_s2s_not_arrived(client) -> None:
"""S2S 还没回调 → 200 + pending(**不是 404**),客户端据此继续轮询。"""
token = _login(client, "13800003601")
r = _reward_result(client, token, "sess-not-yet-arrived")
assert r.status_code == 200, r.text
assert r.json() == {
"ad_session_id": "sess-not-yet-arrived",
"status": "pending",
"coin": None,
# 没记录 → 连属于哪一轮都不知道,累计值一并为 null(不是 0,0 会被读成"本轮没赚到")
"round_coin": None,
}
def test_reward_result_returns_granted_coin(client) -> None:
"""S2S 发奖后按会话查 → granted + 本次真实到账额(与钱包入账一致)。"""
phone = "13800003602"
token = _login(client, phone)
uid = _user_id(phone)
session_id = "sess-granted-1"
r = _callback(
client,
_signed(uid, "trans_rr_1", ecpm="200", extra=_session_extra(session_id)),
)
assert r.json() == {"is_verify": True, "reason": 0}
expected = calculate_ad_reward_coin("200", 1)
body = _reward_result(client, token, session_id).json()
assert body["status"] == "granted"
assert body["coin"] == expected
# 弹窗金额必须等于真实入账,这正是本接口存在的意义(不用余额差估算)
assert _coin_balance(client, token) == expected
def test_reward_result_prefers_granted_over_earlier_noshow(client) -> None:
"""竞态:客户端先报 closed_early、S2S 随后才到 → 同一会话两条记录,必须返回 granted 那条。
只按 created_at 取最近一条是不够的(SQLite 下两条可能同一时间戳),故仓储层显式优先 granted
"""
phone = "13800003603"
token = _login(client, phone)
uid = _user_id(phone)
session_id = "sess-race-noshow"
# 1) 客户端以为没发奖,先留痕
r = client.post(
"/api/v1/ad/reward-noshow",
json={"ad_session_id": session_id, "watched_seconds": 3},
headers=_auth(token),
)
assert r.status_code == 200, r.text
assert r.json()["status"] == "closed_early"
assert _reward_result(client, token, session_id).json()["status"] == "closed_early"
# 2) S2S 姗姗来迟,真发了钱
_callback(client, _signed(uid, "trans_rr_race", ecpm="200", extra=_session_extra(session_id)))
body = _reward_result(client, token, session_id).json()
assert body["status"] == "granted"
assert body["coin"] == calculate_ad_reward_coin("200", 1)
def test_reward_result_capped_reports_zero_not_popup(client) -> None:
"""达每日上限 → capped + coin=0;客户端不弹「获得 0 金币」。"""
phone = "13800003604"
token = _login(client, phone)
uid = _user_id(phone)
session_id = "sess-capped-1"
db = SessionLocal()
try:
db.add(
AdRewardRecord(
trans_id="trans_rr_capped", user_id=_user_id(phone), coin=0, status="capped",
reward_scene="reward_video", ad_session_id=session_id, reward_date="2026-07-17",
)
)
db.commit()
finally:
db.close()
assert uid # 记录挂在该用户名下
body = _reward_result(client, token, session_id).json()
assert body["status"] == "capped"
assert body["coin"] == 0
def test_reward_result_scoped_to_owner(client) -> None:
"""别人的会话查不到(按 user_id 收窄)→ pending,不泄漏他人发奖结果。"""
phone_a = "13800003605"
token_a = _login(client, phone_a)
uid_a = _user_id(phone_a)
token_b = _login(client, "13800003606")
session_id = "sess-owner-only"
_callback(client, _signed(uid_a, "trans_rr_owner", ecpm="200", extra=_session_extra(session_id)))
assert _reward_result(client, token_a, session_id).json()["status"] == "granted"
assert _reward_result(client, token_b, session_id).json()["status"] == "pending"
def test_reward_result_requires_auth(client) -> None:
"""无 Bearer → 401,不裸奔。"""
assert client.get("/api/v1/ad/reward-result/sess-anon-1").status_code == 401
# ===== 膨胀轮累计(boost_round_id → reward-result.round_coin)=====
# 不变量:弹窗数字 == 本轮实际到账之和 == 余额涨幅。三者对不上用户就认为少发了钱。
def _round_extra(session_id: str, round_id: str | None = None, **kv: str) -> str:
data = {"ad_session_id": session_id, **kv}
if round_id is not None:
data["boost_round_id"] = round_id
return json.dumps(data)
def test_round_coin_accumulates_across_ads_in_same_round(client) -> None:
"""一轮连看两条 → round_coin 逐条累计,且等于余额涨幅(第七节验收 1、2 步)。"""
phone = "13800003701"
token = _login(client, phone)
uid = _user_id(phone)
round_id = "b7e1c93a4f6d802b"
_callback(client, _signed(
uid, "trans_round_1", ecpm="200", extra=_round_extra("sess-r1-a", round_id)))
first = calculate_ad_reward_coin("200", 1)
body = _reward_result(client, token, "sess-r1-a").json()
assert body["coin"] == first
assert body["round_coin"] == first # 第 1 条:本轮累计 == 本条
_callback(client, _signed(
uid, "trans_round_2", ecpm="200", extra=_round_extra("sess-r1-b", round_id)))
second = calculate_ad_reward_coin("200", 2) # LT 因子递减,第 2 条比第 1 条少
body = _reward_result(client, token, "sess-r1-b").json()
assert body["coin"] == second
assert body["round_coin"] == first + second # 累计 = 两条之和
# 弹窗数字必须等于真实余额涨幅 —— 这条不变量是整个方案的目的
assert _coin_balance(client, token) == first + second
def test_new_round_restarts_accumulation(client) -> None:
"""换新轮 id → round_coin 从头累计,不接着上一轮往上加(第七节最后一句验收)。"""
phone = "13800003702"
token = _login(client, phone)
uid = _user_id(phone)
_callback(client, _signed(uid, "trans_r2_old", ecpm="200", extra=_round_extra("sess-r2-a", "round-old")))
old = _reward_result(client, token, "sess-r2-a").json()["round_coin"]
assert old > 0
_callback(client, _signed(uid, "trans_r2_new", ecpm="200", extra=_round_extra("sess-r2-b", "round-new")))
body = _reward_result(client, token, "sess-r2-b").json()
assert body["round_coin"] == body["coin"] # 新轮 = 只有本条
assert body["round_coin"] != old + body["coin"]
def test_round_coin_null_without_round_id(client) -> None:
"""extra 没带 boost_round_id(老客户端 / GroMore 丢字段)→ round_coin=null,客户端退回显示单条。"""
phone = "13800003703"
token = _login(client, phone)
uid = _user_id(phone)
_callback(client, _signed(uid, "trans_r3", ecpm="200", extra=_round_extra("sess-r3-noround", None)))
body = _reward_result(client, token, "sess-r3-noround").json()
assert body["coin"] == calculate_ad_reward_coin("200", 1)
assert body["round_coin"] is None
def test_round_coin_null_when_pending(client) -> None:
"""S2S 未到账 → 没有记录 → 连轮 id 都不知道,round_coin 也是 null(不是 0)。"""
token = _login(client, "13800003704")
body = _reward_result(client, token, "sess-r4-pending").json()
assert body == {
"ad_session_id": "sess-r4-pending",
"status": "pending",
"coin": None,
"round_coin": None,
}
def test_round_coin_returned_on_capped(client) -> None:
"""撞每日上限那条不是 granted,但 round_coin **仍返本轮累计**(该条按 0 计)。
客户端的限额 toast 要显示前面几条已到账的总额,不能是空
"""
phone = "13800003705"
token = _login(client, phone)
uid = _user_id(phone)
round_id = "round-capped"
_callback(client, _signed(uid, "trans_cap_ok", ecpm="200", extra=_round_extra("sess-cap-a", round_id)))
earned = _reward_result(client, token, "sess-cap-a").json()["round_coin"]
assert earned > 0
# 手插一条同轮的 capped 记录(跑满 500 次太慢),模拟第 N 条撞上限
db = SessionLocal()
try:
db.add(AdRewardRecord(
trans_id="trans_cap_hit", user_id=uid, coin=0, status="capped",
reward_scene="reward_video", ad_session_id="sess-cap-b",
reward_date="2026-07-20", boost_round_id=round_id,
))
db.commit()
finally:
db.close()
body = _reward_result(client, token, "sess-cap-b").json()
assert body["status"] == "capped"
assert body["coin"] == 0 # 这条没发钱
assert body["round_coin"] == earned # 但本轮累计照常返回
def test_round_coin_scoped_to_owner(client) -> None:
"""轮 id 是客户端生成的,不能跨用户信任:拿别人的轮 id 查不到别人的金币。"""
phone_a = "13800003706"
token_a = _login(client, phone_a)
uid_a = _user_id(phone_a)
phone_b = "13800003707"
token_b = _login(client, phone_b)
uid_b = _user_id(phone_b)
shared_round = "round-collision"
_callback(client, _signed(uid_a, "trans_own_a", ecpm="200", extra=_round_extra("sess-own-a", shared_round)))
a_total = _reward_result(client, token_a, "sess-own-a").json()["round_coin"]
# B 用同一个轮 id(伪造或碰撞)看一条:B 的累计里不能混进 A 的钱
_callback(client, _signed(uid_b, "trans_own_b", ecpm="200", extra=_round_extra("sess-own-b", shared_round)))
b_body = _reward_result(client, token_b, "sess-own-b").json()
assert b_body["round_coin"] == b_body["coin"]
assert b_body["round_coin"] < a_total + b_body["coin"]
def test_test_grant_accepts_boost_round_id(client, monkeypatch) -> None:
"""debug 的 test-grant 不经 S2S、拿不到 mediaExtra,轮 id 由 body 补 → 本地也能验累计。"""
monkeypatch.setattr(settings, "AD_REWARD_TEST_GRANT_ENABLED", True)
token = _login(client, "13800003708")
round_id = "round-testgrant"
coins = []
for i in range(2):
r = client.post(
"/api/v1/ad/test-grant",
json={"reward_scene": "reward_video", "boost_round_id": round_id,
"ad_session_id": f"sess-tg-{i}-padding"},
headers=_auth(token),
)
assert r.status_code == 200, r.text
coins.append(r.json()["coin"])
body = _reward_result(client, token, "sess-tg-1-padding").json()
assert body["round_coin"] == sum(coins)
+61
View File
@@ -107,6 +107,67 @@ def test_sms_send_device_ip_rate_limit(client, monkeypatch) -> None:
assert r.status_code == 200, r.text
def test_sms_send_cooldown_reject_not_counted(client, monkeypatch) -> None:
"""发码额度只算「成功发码」:被单号 60s 冷却挡下的重发(429)不占设备额度。
做法:同号狂发只成功 1 其余被冷却挡下;把小时额度设 2,证明换号后仍能再成功发 1
若冷却重发也计数,额度早被那几次耗尽"""
from app.api.v1 import auth
from app.core import ratelimit
monkeypatch.setattr(ratelimit.settings, "RATE_LIMIT_ENABLED", True)
monkeypatch.setattr(auth, "SMS_SEND_MAX_PER_HOUR_PER_DEVICE", 2)
ratelimit._buckets.clear()
device = "dev-cooldown"
phone_a = "13710137000"
# 首发成功(小时闸计 1/2)
assert client.post(
"/api/v1/auth/sms/send", json={"phone": phone_a, "device_id": device}
).status_code == 200
# 同号连发 3 次:都被单号 60s 冷却挡下 → 429,且**不占**设备额度
for _ in range(3):
r = client.post(
"/api/v1/auth/sms/send", json={"phone": phone_a, "device_id": device}
)
assert r.status_code == 429, r.text
# 换号再发:设备额度只用了 1/2(冷却那几次没算)→ 仍放行(计到 2/2)
assert client.post(
"/api/v1/auth/sms/send", json={"phone": "13710137001", "device_id": device}
).status_code == 200
# 又换号:此时小时闸已 2/2 → 429(反证成功发码确实各计了 1)
r = client.post(
"/api/v1/auth/sms/send", json={"phone": "13710137002", "device_id": device}
)
assert r.status_code == 429, r.text
def test_sms_send_daily_cap(client, monkeypatch) -> None:
"""每天发码上限(设备 + IP):成功发码累计到日上限即 429(用不同手机号绕开单号冷却)。
抬高小时闸单独测日闸;超限文案含今日以便前端提示明天再来"""
from app.api.v1 import auth
from app.core import ratelimit
monkeypatch.setattr(ratelimit.settings, "RATE_LIMIT_ENABLED", True)
monkeypatch.setattr(auth, "SMS_SEND_MAX_PER_HOUR_PER_DEVICE", 100) # 抬高小时闸,不干扰
monkeypatch.setattr(auth, "SMS_SEND_MAX_PER_DAY_PER_DEVICE", 3)
ratelimit._buckets.clear()
device = "dev-daily"
for i in range(3):
r = client.post(
"/api/v1/auth/sms/send",
json={"phone": f"13720137{i:03d}", "device_id": device},
)
assert r.status_code == 200, f"{i + 1} 次应放行: {r.text}"
# 第 4 次:同设备同 IP 当日超限 → 429
r = client.post(
"/api/v1/auth/sms/send",
json={"phone": "13720137999", "device_id": device},
)
assert r.status_code == 429, r.text
assert "今日" in r.json()["detail"]
def test_sms_login_device_ip_rate_limit(client, monkeypatch) -> None:
"""防刷:同一设备(device_id) + 同一 IP 每小时最多 SMS_LOGIN_MAX_PER_HOUR 次登录尝试,超出 429。
conftest 默认 RATE_LIMIT_ENABLED=false(内存计数跨用例累加),本用例临时打开并清空计数隔离"""
+57
View File
@@ -0,0 +1,57 @@
"""ratelimit 内存桶过期清理(GC)测试。
回归重点:_buckets **全局共享**混着不同窗口(60s 广告 / 3600s 登录 / 86400s 日闸) key
GC 必须按每个 key 自己存的 window_sec判过期,而不是当前调用方的窗口 否则高频的 60s 端点
触发 GC 时会把本该存活更久的 3600s/86400s 计数(如短信日闸)一并删掉,使其被反复清零限流失效
monkeypatch _GC_THRESHOLD 0 强制每次都扫,免造上万条(仿 test_auth 里对 sms._GC_THRESHOLD 的做法)
"""
from __future__ import annotations
from app.core import ratelimit
def test_purge_expired_respects_each_key_own_window(monkeypatch) -> None:
"""短窗口(60s)触发的 GC 只删真正过期的 key,不得删掉仍在自身窗口内的长窗口 key。"""
monkeypatch.setattr(ratelimit, "_GC_THRESHOLD", 0) # 强制每次都扫
ratelimit._buckets.clear()
now = 1_000_000.0
# 日闸:100s 前开窗、window=86400 → 远未过期,必须保留
ratelimit._buckets["sms-send-device-daily:D:IP"] = (now - 100, 7, 86400.0)
# 登录:1800s、window=3600 → 未过期,保留
ratelimit._buckets["sms-login-device:D:IP"] = (now - 1800, 2, 3600.0)
# 广告:120s、window=60 → 已过期,应删
ratelimit._buckets["ad-watch-report:IP2"] = (now - 120, 3, 60.0)
ratelimit._purge_expired(now)
assert "sms-send-device-daily:D:IP" in ratelimit._buckets
assert "sms-login-device:D:IP" in ratelimit._buckets
assert "ad-watch-report:IP2" not in ratelimit._buckets
def test_purge_expired_keeps_long_window_key_older_than_short_window(monkeypatch) -> None:
"""反证旧 bug:日闸 key 已老于 3600s,旧代码在 60s/3600s 端点触发 GC 时会误删它;
现在按自身 86400s 窗口判 未过期 必须保留"""
monkeypatch.setattr(ratelimit, "_GC_THRESHOLD", 0)
ratelimit._buckets.clear()
now = 2_000_000.0
# 3700s 前开窗(> 1 小时),但 window=86400 → 未过期
ratelimit._buckets["sms-send-device-daily:D:IP"] = (now - 3700, 20, 86400.0)
ratelimit._purge_expired(now)
assert "sms-send-device-daily:D:IP" in ratelimit._buckets
def test_purge_expired_noop_below_threshold(monkeypatch) -> None:
"""未超阈值时不扫(即便有过期 key 也不动),避免每次请求都 O(n) 扫全表。"""
monkeypatch.setattr(ratelimit, "_GC_THRESHOLD", 10)
ratelimit._buckets.clear()
now = 3_000_000.0
ratelimit._buckets["stale:IP"] = (now - 999, 1, 60.0) # 早过期,但没超阈值
ratelimit._purge_expired(now)
assert "stale:IP" in ratelimit._buckets # 桶数没超阈值 → 不清理
-43
View File
@@ -11,7 +11,6 @@ from app.core.rewards import (
COIN_PER_CENT,
COIN_PER_YUAN,
MIN_EXCHANGE_COIN,
SIGNIN_BOOST_COIN,
SIGNIN_REWARDS,
TASK_ENABLE_NOTIFICATION,
TASK_REWARDS,
@@ -106,48 +105,6 @@ def test_signin_flow(client) -> None:
assert txn["balance_after"] == SIGNIN_REWARDS[0]
def test_signin_boost_flow(client) -> None:
"""签到后看广告膨胀 → S2S 固定补发 2000 金币,每天只能膨胀一次。"""
phone = "13800001011"
token = _login(client, phone)
r = client.post("/api/v1/signin/boost", json={}, headers=_auth(token))
assert r.status_code == 409
r = client.post("/api/v1/signin", headers=_auth(token))
assert r.status_code == 200, r.text
first_coin = r.json()["coin_awarded"]
with SessionLocal() as db:
user = get_user_by_phone(db, phone)
assert user is not None
uid = user.id
extra = json.dumps({"reward_scene": "signin_boost", "ad_session_id": "signin-session-1"})
r = client.get(
"/api/v1/ad/pangle-callback",
params=_signed_ad_callback(uid, "signin-boost-trans-1", extra=extra, ecpm="200"),
)
assert r.status_code == 200, r.text
assert r.json() == {"is_verify": True, "reason": 0}
r = client.post("/api/v1/signin/boost", json={"ad_ref_id": "signin-boost-trans-1"}, headers=_auth(token))
assert r.status_code == 200, r.text
body = r.json()
assert body["coin_awarded"] == SIGNIN_BOOST_COIN
assert body["coin_balance"] == first_coin + SIGNIN_BOOST_COIN
r = client.get(
"/api/v1/ad/pangle-callback",
params=_signed_ad_callback(uid, "signin-boost-trans-2", extra=extra, ecpm="200"),
)
assert r.status_code == 200, r.text
r = client.get("/api/v1/wallet/coin-transactions", headers=_auth(token))
types = [t["biz_type"] for t in r.json()["items"]]
assert "signin" in types
assert "signin_boost" in types
def test_task_claim_flow(client) -> None:
"""打开消息提醒=可重复任务:每次领取金额减半(750/375/188),claimed 恒 False,余额累加。"""