Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3250a465c7 | |||
| de1fd58749 | |||
| e9fd51d119 |
@@ -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
|
||||
)
|
||||
@@ -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(
|
||||
|
||||
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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(),
|
||||
)
|
||||
|
||||
@@ -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)
|
||||
@@ -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 标记
|
||||
→ 心跳带回前端弹窗,以及 openid↔device 绑定。任何异常都返回 "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)
|
||||
@@ -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 一致。
|
||||
|
||||
@@ -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
@@ -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
@@ -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"))
|
||||
|
||||
@@ -1,27 +0,0 @@
|
||||
"""微信截图比价 PoC:进程内存的"待比价"信号(单 worker 够用, 重启即失效, PoC 可接受)。
|
||||
|
||||
收图端点 set_pending(device_id, source) → 该设备下次心跳 pop_pending 取走并清除 →
|
||||
心跳响应带回 → 前端弹选平台窗。只对 settings.WX_POC_TEST_DEVICE_ID 写入, 不碰其他设备。
|
||||
后续接真识别 / openid↔device 绑定时整体替换本模块。
|
||||
"""
|
||||
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)
|
||||
@@ -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 定:登录风控只留单号冷却 + 单设备频控);
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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]
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
@@ -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}>"
|
||||
|
||||
@@ -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_early、S2S 随后才姗姗来迟 → 两条,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:
|
||||
|
||||
@@ -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
@@ -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="今日已成功发奖次数")
|
||||
|
||||
@@ -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(无告警)。"""
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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
|
||||
@@ -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
@@ -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) |
|
||||
|
||||
@@ -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 | 代码位 ID(GroMore 带,目前仅入 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`。
|
||||
|
||||
@@ -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` 才取最近一条,让客户端知道没发的原因。
|
||||
@@ -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` 金币流水。它让已登录客户端能自助发奖 = 绕过反作弊,**严禁在生产开启**。
|
||||
它让已登录客户端能自助发奖 = 绕过反作弊,**严禁在生产开启**。
|
||||
|
||||
@@ -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**
|
||||
| 字段 | 类型 | 说明 |
|
||||
|
||||
@@ -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`。
|
||||
@@ -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` 分账校验「单 ↔ 流水」。
|
||||
|
||||
@@ -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) |
|
||||
|
||||
@@ -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 回滚返回已存在那条(幂等兜底)。
|
||||
|
||||
@@ -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() | 最后修改时间 |
|
||||
|
||||
@@ -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`。
|
||||
|
||||
@@ -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 串)。
|
||||
@@ -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 后台管理员,两套体系)。
|
||||
|
||||
## 索引与约束
|
||||
|
||||
@@ -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 后投放匹配 + 填充会明显改善
|
||||
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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()
|
||||
@@ -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)
|
||||
|
||||
@@ -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(内存计数跨用例累加),本用例临时打开并清空计数隔离。"""
|
||||
|
||||
@@ -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 # 桶数没超阈值 → 不清理
|
||||
@@ -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,余额累加。"""
|
||||
|
||||
Reference in New Issue
Block a user