Compare commits

..

5 Commits

Author SHA1 Message Date
exinglang d484bd5801 Merge remote-tracking branch 'origin/main' into fix-mate20PageError
# Conflicts:
#	app/integrations/vendor_push.py
2026-07-30 10:44:15 +08:00
exinglang 7474c2bebc fix(push): 补齐厂商请求日志参数 2026-07-30 10:26:26 +08:00
exinglang 7bfe703760 chore(tests): 取消推送测试文件改动 2026-07-30 10:25:05 +08:00
exinglang 69eeb43fe2 fix(push): 完善厂商推送参数与消息通知 2026-07-30 10:25:05 +08:00
exinglang a0196b8d64 fix(push): 完善厂商推送排障日志 2026-07-29 19:52:06 +08:00
32 changed files with 422 additions and 1607 deletions
+1 -10
View File
@@ -69,10 +69,8 @@ CHUANGLAN_SMS_TIMEOUT_SEC=10
# 配一个固定测试手机号,专供无 SIM 卡 / 不走一键登录时打通全流程:该号登录【免短信验证码】
# (real 模式下也跳过校验)、每次登录【都重走新手引导】,并有【每日登录上限】防被人猜到号后脚本刷。
# 逻辑见 app/core/test_account.py,与其他业务解耦。
# ⚠️ 两项都留空 = 关闭整功能(生产默认)。多账号用英文逗号分隔,改完重启生效
# ⚠️ 留空 = 关闭整功能(生产默认);要启用才填号(如 11111111111)。改完重启生效,随时可清空停用
TEST_ACCOUNT_PHONE=
# 推荐新配置;例如联调环境填 11111111111,22222222222。
TEST_ACCOUNT_PHONES=
# 该测试号每日最多登录次数,当日超过即拒绝(429),次日归零。
TEST_ACCOUNT_DAILY_LIMIT=500
@@ -178,13 +176,6 @@ PANGLE_REPORT_SITE_ID_TEST=5832303
# APPLOG_MAX_BODY_BYTES=2097152 # 请求体上限 2MB(超 → 413;运行期可调)
# APPLOG_MAX_MSG_BYTES=8192 # 单条 msg 超此字节数截断
# ===== 荣耀 Push =====
HONOR_PUSH_APP_ID=
HONOR_PUSH_CLIENT_ID=
HONOR_PUSH_CLIENT_SECRET=
# 0=正式消息(默认);1=测试消息(仅开发联调,生产必须保持 0)
HONOR_PUSH_TARGET_USER_TYPE=0
# ===== 华为 Push =====
HUAWEI_PUSH_APP_ID=
HUAWEI_PUSH_APP_SECRET=
@@ -1,164 +0,0 @@
"""guide video prepare/start/ten-circle reward state machine
Revision ID: guide_video_ten_circle_v2
Revises: guide_video_scene_unique
"""
from __future__ import annotations
import json
from pathlib import Path
import sqlalchemy as sa
from alembic import op
revision = "guide_video_ten_circle_v2"
down_revision = "guide_video_scene_unique"
branch_labels = None
depends_on = None
def _backfill_configs(connection) -> None:
from app.core import media
from app.core.config import settings
for key in ("coupon_guide_video", "comparison_guide_video"):
row = connection.execute(
sa.text("SELECT value FROM app_config WHERE key = :key"), {"key": key}
).first()
if not row:
continue
raw = row[0]
value = json.loads(raw) if isinstance(raw, str) else dict(raw or {})
if value.get("guide_video_v2_migrated") is True:
continue
maximum = int(value.get("max_plays", 3) or 0)
reward = int(value.get("reward_coin", 100) or 0)
video_url = str(value.get("video_url") or "")
error = None
analysis = None
if not 1 <= maximum <= 50:
error = "旧播放次数不符合 150 约束"
elif not 10 <= reward <= 10_000 or reward % 10:
error = "旧金币总价不符合 10~10000 且为 10 的倍数约束"
elif not video_url:
error = "尚未上传视频"
else:
prefix = f"{settings.MEDIA_URL_PREFIX}/guide_video/"
if not video_url.startswith(prefix):
error = "旧视频不是服务端托管文件,请重新上传"
else:
path = Path(settings.MEDIA_ROOT) / "guide_video" / video_url[len(prefix):]
if not path.is_file():
error = "旧视频文件缺失,请重新上传"
else:
try:
analysis = media._probe_guide_video(path)
except media.MediaError as exc:
error = str(exc)
value["config_version"] = int(value.get("config_version", 0) or 0) + 1
value["guide_video_v2_migrated"] = True
if analysis:
value.update(analysis)
else:
value.update(
enabled=False,
duration_ms=None,
video_codec=None,
audio_codec=None,
analysis_status="invalid" if video_url else "missing",
analysis_error=error,
)
statement = sa.text(
"UPDATE app_config SET value = :value WHERE key = :key"
).bindparams(sa.bindparam("value", type_=sa.JSON()))
connection.execute(statement, {"key": key, "value": value})
def _seal_legacy_plays(connection) -> None:
"""旧整笔发奖 token 永久封口;旧未发奖会话终止但仍保留计次行。"""
connection.execute(
sa.text(
"""
UPDATE guide_video_play
SET status = CASE WHEN status = 'granted' THEN 'legacy_completed' ELSE 'legacy_closed' END,
settled_circles = CASE WHEN status = 'granted' THEN 10 ELSE 0 END,
prepared_at = started_at,
expires_at = started_at
WHERE status IN ('granted', 'playing')
"""
)
)
def _prepare_downgrade(connection) -> None:
"""清除未起播计划,并把所有保留 token 封成旧版不可再领奖的 granted。"""
connection.execute(
sa.text("DELETE FROM guide_video_play WHERE status = 'prepared'")
)
connection.execute(
sa.text(
"""
UPDATE guide_video_play
SET status = 'granted',
completed = CASE WHEN settled_circles >= 10 THEN 1 ELSE completed END,
granted_at = CASE
WHEN settled_circles >= 10 AND granted_at IS NULL THEN started_at
ELSE granted_at
END
WHERE status IN (
'started', 'completed', 'legacy_completed', 'legacy_closed'
)
"""
)
)
# 防御未知/中间状态:宁可删除,也绝不让旧应用把它当 playing 整笔发奖。
connection.execute(
sa.text(
"""
DELETE FROM guide_video_play
WHERE seq IS NULL OR started_at IS NULL OR status != 'granted'
"""
)
)
def upgrade() -> None:
with op.batch_alter_table("guide_video_play") as batch:
batch.alter_column("seq", existing_type=sa.Integer(), nullable=True)
batch.alter_column("started_at", existing_type=sa.DateTime(timezone=True), nullable=True)
batch.add_column(sa.Column("duration_ms", sa.Integer(), nullable=False, server_default="0"))
batch.add_column(sa.Column("config_version", sa.Integer(), nullable=False, server_default="0"))
batch.add_column(sa.Column("settled_circles", sa.Integer(), nullable=False, server_default="0"))
batch.add_column(
sa.Column("prepared_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now())
)
batch.add_column(sa.Column("expires_at", sa.DateTime(timezone=True), nullable=True))
batch.alter_column("status", existing_type=sa.String(length=16), type_=sa.String(length=24))
op.create_index("ix_guide_video_play_expires_at", "guide_video_play", ["expires_at"])
op.create_index(
"ux_coin_transaction_guide_video_ref",
"coin_transaction",
["user_id", "biz_type", "ref_id"],
unique=True,
sqlite_where=sa.text("biz_type = 'guide_video' AND ref_id IS NOT NULL"),
postgresql_where=sa.text("biz_type = 'guide_video' AND ref_id IS NOT NULL"),
)
connection = op.get_bind()
_seal_legacy_plays(connection)
_backfill_configs(connection)
def downgrade() -> None:
_prepare_downgrade(op.get_bind())
op.drop_index("ux_coin_transaction_guide_video_ref", table_name="coin_transaction")
op.drop_index("ix_guide_video_play_expires_at", table_name="guide_video_play")
with op.batch_alter_table("guide_video_play") as batch:
batch.alter_column("status", existing_type=sa.String(length=24), type_=sa.String(length=16))
batch.drop_column("expires_at")
batch.drop_column("prepared_at")
batch.drop_column("settled_circles")
batch.drop_column("config_version")
batch.drop_column("duration_ms")
batch.alter_column("started_at", existing_type=sa.DateTime(timezone=True), nullable=False)
batch.alter_column("seq", existing_type=sa.Integer(), nullable=False)
@@ -1,89 +0,0 @@
"""isolate vendor push binding across user accounts
Revision ID: push_binding_isolation
Revises: guide_video_ten_circle_v2
"""
from __future__ import annotations
import sqlalchemy as sa
from alembic import op
revision = "push_binding_isolation"
down_revision = "guide_video_ten_circle_v2"
branch_labels = None
depends_on = None
def upgrade() -> None:
with op.batch_alter_table("device_liveness") as batch_op:
batch_op.add_column(sa.Column("push_binding_id", sa.String(length=128), nullable=True))
batch_op.add_column(
sa.Column(
"push_binding_revoked",
sa.Boolean(),
nullable=False,
server_default=sa.false(),
)
)
batch_op.create_index(
"ix_device_liveness_push_binding_id",
["push_binding_id"],
unique=False,
)
connection = op.get_bind()
device = sa.table(
"device_liveness",
sa.column("id", sa.Integer),
sa.column("push_vendor", sa.String),
sa.column("push_token", sa.String),
sa.column("updated_at", sa.DateTime),
)
connection.execute(
device.update()
.where(
sa.or_(
device.c.push_vendor == "",
device.c.push_token == "",
)
)
.values(push_vendor=None, push_token=None)
)
rows = connection.execute(
sa.select(device.c.id, device.c.push_vendor, device.c.push_token)
.where(
device.c.push_vendor.is_not(None),
device.c.push_token.is_not(None),
)
.order_by(device.c.updated_at.desc(), device.c.id.desc())
)
seen: set[tuple[str, str]] = set()
duplicate_ids: list[int] = []
for row in rows:
key = (row.push_vendor, row.push_token)
if key in seen:
duplicate_ids.append(row.id)
else:
seen.add(key)
if duplicate_ids:
connection.execute(
device.update()
.where(device.c.id.in_(duplicate_ids))
.values(push_vendor=None, push_token=None)
)
with op.batch_alter_table("device_liveness") as batch_op:
batch_op.create_unique_constraint(
"uq_device_liveness_vendor_token",
["push_vendor", "push_token"],
)
def downgrade() -> None:
with op.batch_alter_table("device_liveness") as batch_op:
batch_op.drop_constraint("uq_device_liveness_vendor_token", type_="unique")
batch_op.drop_index("ix_device_liveness_push_binding_id")
batch_op.drop_column("push_binding_revoked")
batch_op.drop_column("push_binding_id")
+1 -1
View File
@@ -1,6 +1,6 @@
"""Admin 后台 FastAPI app(独立进程)。
启动:uvicorn app.admin.main:admin_app --host 127.0.0.1 --port 8773
启动:uvicorn app.admin.main:admin_app --host 127.0.0.1 --port 8771
复用 App 的 DB/models/repositories/integrations;鉴权独立(admin JWT,见 app/admin/security.py)。
现有 app.main:app 不 import 本模块,两进程互不影响。
"""
+6 -109
View File
@@ -1,7 +1,6 @@
"""admin 反馈工单:列表筛选 + 审核采纳/拒绝(带金币发放与审计)。"""
from __future__ import annotations
import logging
from datetime import datetime
from typing import Annotated
@@ -26,8 +25,6 @@ from app.models.feedback import Feedback
from app.repositories import wallet as wallet_repo
from app.services import notification_events
logger = logging.getLogger("shagua.admin.feedback")
router = APIRouter(
prefix="/admin/api/feedbacks",
tags=["admin-feedback"],
@@ -49,24 +46,10 @@ def _approve_feedback(
*,
bulk: bool = False,
) -> FeedbackOut:
logger.info(
"feedback approve started feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
fb = db.get(Feedback, feedback_id, with_for_update=True)
if fb is None:
logger.warning(
"feedback approve rejected not found feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
raise HTTPException(status_code=404, detail="反馈不存在")
if fb.status not in {"pending", "new"}:
logger.warning(
"feedback approve rejected invalid status feedback_id=%s user_id=%s "
"admin_id=%s status=%s bulk=%s",
feedback_id, fb.user_id, admin.id, fb.status, bulk,
)
_ensure_pending(fb)
_ensure_pending(fb)
before = fb.status
mutations.review_feedback(
@@ -108,18 +91,8 @@ def _approve_feedback(
)
db.commit()
db.refresh(fb)
logger.info(
"feedback approve committed feedback_id=%s user_id=%s admin_id=%s "
"before=%s after=%s reward_coins=%s bulk=%s",
feedback_id, fb.user_id, admin.id, before, fb.status, payload.reward_coins, bulk,
)
out = FeedbackOut.model_validate(fb)
notification_events.notify_feedback_reward(db, fb)
logger.info(
"feedback approve notification dispatch returned feedback_id=%s user_id=%s "
"admin_id=%s notification_type=feedback_reward bulk=%s",
feedback_id, fb.user_id, admin.id, bulk,
)
return out
@@ -132,24 +105,10 @@ def _reject_feedback(
*,
bulk: bool = False,
) -> FeedbackOut:
logger.info(
"feedback reject started feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
fb = db.get(Feedback, feedback_id, with_for_update=True)
if fb is None:
logger.warning(
"feedback reject rejected not found feedback_id=%s admin_id=%s bulk=%s",
feedback_id, admin.id, bulk,
)
raise HTTPException(status_code=404, detail="反馈不存在")
if fb.status not in {"pending", "new"}:
logger.warning(
"feedback reject rejected invalid status feedback_id=%s user_id=%s "
"admin_id=%s status=%s bulk=%s",
feedback_id, fb.user_id, admin.id, fb.status, bulk,
)
_ensure_pending(fb)
_ensure_pending(fb)
before = fb.status
mutations.review_feedback(
@@ -183,18 +142,8 @@ def _reject_feedback(
)
db.commit()
db.refresh(fb)
logger.info(
"feedback reject committed feedback_id=%s user_id=%s admin_id=%s "
"before=%s after=%s bulk=%s",
feedback_id, fb.user_id, admin.id, before, fb.status, bulk,
)
out = FeedbackOut.model_validate(fb)
notification_events.notify_feedback_reply(db, fb)
logger.info(
"feedback reject notification dispatch returned feedback_id=%s user_id=%s "
"admin_id=%s notification_type=feedback_reply bulk=%s",
feedback_id, fb.user_id, admin.id, bulk,
)
return out
@@ -252,10 +201,6 @@ def bulk_approve_feedbacks(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> FeedbackBulkResult:
logger.info(
"feedback bulk approve started admin_id=%s item_count=%s",
admin.id, len(body.ids),
)
results: list[FeedbackBulkItemResult] = []
ip = get_client_ip(request)
for feedback_id in body.ids:
@@ -264,24 +209,11 @@ def bulk_approve_feedbacks(
results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status))
except HTTPException as exc:
db.rollback()
logger.warning(
"feedback bulk approve item failed feedback_id=%s admin_id=%s error=%s",
feedback_id, admin.id, exc.detail,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail)))
except Exception: # noqa: BLE001 - 单笔失败不打断整批
db.rollback()
logger.exception(
"feedback bulk approve item failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常"))
result = _bulk_result(results)
logger.info(
"feedback bulk approve completed admin_id=%s total=%s success=%s failed=%s",
admin.id, result.total, result.success, result.failed,
)
return result
return _bulk_result(results)
@router.post("/bulk/reject", response_model=FeedbackBulkResult, summary="批量拒绝采纳反馈")
@@ -291,10 +223,6 @@ def bulk_reject_feedbacks(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> FeedbackBulkResult:
logger.info(
"feedback bulk reject started admin_id=%s item_count=%s",
admin.id, len(body.ids),
)
results: list[FeedbackBulkItemResult] = []
ip = get_client_ip(request)
for feedback_id in body.ids:
@@ -303,24 +231,11 @@ def bulk_reject_feedbacks(
results.append(FeedbackBulkItemResult(id=feedback_id, ok=True, status=out.status))
except HTTPException as exc:
db.rollback()
logger.warning(
"feedback bulk reject item failed feedback_id=%s admin_id=%s error=%s",
feedback_id, admin.id, exc.detail,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error=str(exc.detail)))
except Exception: # noqa: BLE001 - 单笔失败不打断整批
db.rollback()
logger.exception(
"feedback bulk reject item failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
results.append(FeedbackBulkItemResult(id=feedback_id, ok=False, error="系统异常"))
result = _bulk_result(results)
logger.info(
"feedback bulk reject completed admin_id=%s total=%s success=%s failed=%s",
admin.id, result.total, result.success, result.failed,
)
return result
return _bulk_result(results)
@router.post("/{feedback_id}/handle", response_model=OkResponse, summary="标记反馈已处理")
@@ -343,16 +258,7 @@ def approve_feedback(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> FeedbackOut:
try:
return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request))
except HTTPException:
raise
except Exception:
logger.exception(
"feedback approve failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
raise
return _approve_feedback(db, admin, feedback_id, payload, get_client_ip(request))
@router.post("/{feedback_id}/reject", response_model=FeedbackOut, summary="拒绝采纳反馈")
@@ -363,13 +269,4 @@ def reject_feedback(
admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb,
) -> FeedbackOut:
try:
return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request))
except HTTPException:
raise
except Exception:
logger.exception(
"feedback reject failed feedback_id=%s admin_id=%s",
feedback_id, admin.id,
)
raise
return _reject_feedback(db, admin, feedback_id, payload, get_client_ip(request))
+28 -41
View File
@@ -31,7 +31,12 @@ GuideScene = Literal["coupon", "comparison"]
def _out(db: AdminDb, scene: GuideScene) -> GuideVideoConfigOut:
return GuideVideoConfigOut(scene=scene, **guide_video.get_config(db, scene))
"""配置 + 播放统计合成响应(四个写接口都以最新状态返回,前端一次同步到位)。"""
return GuideVideoConfigOut(
scene=scene,
**guide_video.get_config(db, scene),
**guide_video.play_stats(db, scene),
)
@router.get("", response_model=GuideVideoConfigOut, summary="新手引导视频配置(领券浮层)")
@@ -47,18 +52,15 @@ def update_config(
db: AdminDb,
scene: GuideScene = "coupon",
) -> GuideVideoConfigOut:
try:
before, after = guide_video.update_config(
db,
enabled=body.enabled,
max_plays=body.max_plays,
reward_coin=body.reward_coin,
scene=scene,
admin_id=admin.id,
commit=False,
)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
before, after = guide_video.update_config(
db,
enabled=body.enabled,
max_plays=body.max_plays,
reward_coin=body.reward_coin,
scene=scene,
admin_id=admin.id,
commit=False,
)
write_audit(
db, admin, action="guide_video.update", target_type="guide_video", target_id=None,
detail={"scene": scene, "before": before, "after": after}, ip=get_client_ip(request), commit=False,
@@ -76,35 +78,21 @@ async def upload_video(
scene: GuideScene = "coupon",
) -> GuideVideoConfigOut:
data = await file.read()
url: str | None = None
try:
url, analysis = media.save_guide_video(data)
url = media.save_guide_video(data)
except media.MediaError as e:
raise HTTPException(status_code=400, detail=str(e)) from e
try:
before, after = guide_video.set_video(
db, url, analysis=analysis, scene=scene, admin_id=admin.id, commit=False
)
write_audit(
db, admin, action="guide_video.set_video", target_type="guide_video", target_id=None,
detail={
"scene": scene,
"before": before.get("video_url"),
"after": url,
"bytes": len(data),
"duration_ms": analysis["duration_ms"],
"video_codec": analysis["video_codec"],
"audio_codec": analysis["audio_codec"],
},
ip=get_client_ip(request), commit=False,
)
db.commit()
except Exception:
db.rollback()
media.delete_guide_video(url)
raise
# 旧文件仍可能被 prepare/已起播快照引用,不能在这里立即删除。
guide_video.cleanup_old_videos(db)
before, after = guide_video.set_video(
db, url, scene=scene, admin_id=admin.id, commit=False
)
write_audit(
db, admin, action="guide_video.set_video", target_type="guide_video", target_id=None,
detail={"scene": scene, "before": before.get("video_url"), "after": url, "bytes": len(data)},
ip=get_client_ip(request), commit=False,
)
db.commit()
# 提交成功后再删旧片,避免新片没落库就把旧片丢了
media.delete_guide_video(before.get("video_url"))
return _out(db, scene)
@@ -124,6 +112,5 @@ def delete_video(
detail={"scene": scene, "before": before.get("video_url")}, ip=get_client_ip(request), commit=False,
)
db.commit()
# 旧文件由保护期清理任务处理,避免破坏有效计划与已起播快照。
guide_video.cleanup_old_videos(db)
media.delete_guide_video(before.get("video_url"))
return _out(db, scene)
+11 -29
View File
@@ -1,44 +1,26 @@
"""后台引导视频配置协议"""
"""admin 新手引导视频配置 schemas(开关 / 视频地址 / 前几次 / 每次金币)"""
from __future__ import annotations
from pydantic import BaseModel, Field, field_validator
from pydantic import BaseModel, Field
from app.repositories.guide_video import (
MAX_PLAYS_LIMIT,
MIN_PLAYS,
MIN_REWARD_COIN,
REWARD_COIN_LIMIT,
)
from app.repositories.guide_video import MAX_PLAYS_LIMIT, REWARD_COIN_LIMIT
class GuideVideoConfigOut(BaseModel):
scene: str
enabled: bool
video_url: str | None = None
video_url: str | None = None # 相对地址 /media/guide_video/xxx.mp4;未配片 = None
max_plays: int
reward_coin: int
duration_ms: int | None = None
circle_count: int = 10
circle_duration_ms: float | None = None
reward_per_circle: int
video_codec: str | None = None
audio_codec: str | None = None
analysis_status: str
analysis_error: str | None = None
config_version: int
updated_at: str | None = None
# 只读统计,后台展示用:已有多少次播放、其中已发币多少次。
total_plays: int = 0
granted_plays: int = 0
class GuideVideoConfigUpdate(BaseModel):
enabled: bool | None = None
max_plays: int | None = Field(default=None, ge=MIN_PLAYS, le=MAX_PLAYS_LIMIT)
reward_coin: int | None = Field(
default=None, ge=MIN_REWARD_COIN, le=REWARD_COIN_LIMIT
)
"""部分更新:只改传入(非 None)字段。视频文件走 /video 上传接口。"""
@field_validator("reward_coin")
@classmethod
def reward_must_be_multiple_of_ten(cls, value: int | None) -> int | None:
if value is not None and value % 10:
raise ValueError("金币总价必须是 10 的倍数")
return value
enabled: bool | None = None
max_plays: int | None = Field(default=None, ge=0, le=MAX_PLAYS_LIMIT)
reward_coin: int | None = Field(default=None, ge=0, le=REWARD_COIN_LIMIT)
+1 -1
View File
@@ -256,7 +256,7 @@ def sms_login(req: SmsLoginRequest, request: Request, db: DbSession) -> TokenWit
# 放在最前面:命中即不校验验证码;先扣当日额度,超限直接拒,挡住有人猜到号后脚本刷。
# 测试账号走自己的每日额度、不受下面 (设备+IP) 每小时限流约束(QA 需在一小时内反复登录联调)。
if test_account.is_test_account(req.phone):
if not test_account.try_consume_quota(req.phone):
if not test_account.try_consume_quota():
raise HTTPException(status_code=429, detail="测试账号今日使用次数已达上限,请明天再试")
user = user_repo.upsert_user_for_login(db, phone=req.phone, register_channel="sms")
risk_repo.record_behavior_event(
-33
View File
@@ -19,10 +19,8 @@ from app.api.deps import CurrentUser, DbSession
from app.integrations import vendor_push
from app.repositories import device as device_repo
from app.schemas.device import (
DeferredDeviceUnregisterRequest,
DeviceOut,
DeviceRegisterRequest,
DeviceUnregisterRequest,
HeartbeatRequest,
LivenessAckRequest,
LivenessOut,
@@ -80,7 +78,6 @@ def register_device(
registration_id=req.registration_id,
push_vendor=req.push_vendor,
push_token=req.push_token,
push_binding_id=req.push_binding_id,
platform=req.platform,
app_version=req.app_version,
)
@@ -109,36 +106,6 @@ def report_heartbeat(
registration_id=req.registration_id,
push_vendor=req.push_vendor,
push_token=req.push_token,
push_binding_id=req.push_binding_id,
)
return OkResponse()
@router.post("/unregister", response_model=OkResponse, summary="解绑当前用户的本机推送目标")
def unregister_device(
req: DeviceUnregisterRequest,
user: CurrentUser,
db: DbSession,
) -> OkResponse:
device_repo.unregister_push_binding(db, user_id=user.id, device_id=req.device_id)
logger.info("device unregister user_id=%d device_id=%s", user.id, req.device_id)
return OkResponse()
@router.post(
"/unregister/deferred",
response_model=OkResponse,
summary="离线退出后按登录会话凭据幂等解绑",
)
def unregister_device_deferred(
req: DeferredDeviceUnregisterRequest,
db: DbSession,
) -> OkResponse:
# 不返回是否命中,避免公开接口泄露设备绑定状态。
device_repo.unregister_push_binding_deferred(
db,
device_id=req.device_id,
push_binding_id=req.push_binding_id,
)
return OkResponse()
+27 -42
View File
@@ -1,16 +1,22 @@
"""引导视频 prepare/start/reward 客户端 API。"""
"""新手引导视频(领券等候浮层前 N 次替代广告)。
路由前缀 `/api/v1/guide-video`(均需 Bearer):
POST /start 这次浮层放引导视频还是放广告?命中则**当场计次**并下发 play_token
POST /reward 播完 / 中途关闭都调,按 play_token 幂等发固定金币
发币额度以**服务端配置**为准(运营后台可改),客户端只报"播完/关闭",报不了金额,
所以被破解也刷不到超额金币;次数上限由 guide_video_play 行数(按账号)硬卡。
"""
from __future__ import annotations
import logging
from fastapi import APIRouter, Depends, HTTPException
from fastapi import APIRouter, Depends
from app.api.deps import CurrentUser, DbSession
from app.core.ratelimit import rate_limit
from app.repositories import guide_video as crud_guide
from app.schemas.guide_video import (
GuideVideoPrepareIn,
GuideVideoPrepareOut,
GuideVideoRewardIn,
GuideVideoRewardOut,
GuideVideoStartIn,
@@ -18,41 +24,25 @@ from app.schemas.guide_video import (
)
logger = logging.getLogger("shagua.guide_video")
router = APIRouter(prefix="/api/v1/guide-video", tags=["guide-video"])
@router.post(
"/prepare",
response_model=GuideVideoPrepareOut,
dependencies=[Depends(rate_limit(60, 60, "guide-video-prepare"))],
)
def prepare(
payload: GuideVideoPrepareIn, user: CurrentUser, db: DbSession
) -> GuideVideoPrepareOut:
result = crud_guide.prepare_play(db, user.id, scene=payload.scene)
logger.info(
"guide video prepare user_id=%d scene=%s should_play=%s reason=%s",
user.id, payload.scene, result["should_play"], result["reason"],
)
return GuideVideoPrepareOut(**result)
@router.post(
"/start",
response_model=GuideVideoStartOut,
summary="领券浮层是否放新手引导视频(命中即计次)",
dependencies=[Depends(rate_limit(60, 60, "guide-video-start"))],
)
def start(
payload: GuideVideoStartIn, user: CurrentUser, db: DbSession
) -> GuideVideoStartOut:
try:
result = crud_guide.start_play(db, user.id, play_token=payload.play_token)
except crud_guide.PlayStateError as exc:
status_code = 404 if exc.code == "play_not_found" else 409
raise HTTPException(status_code=status_code, detail=exc.detail()) from exc
def start(payload: GuideVideoStartIn, user: CurrentUser, db: DbSession) -> GuideVideoStartOut:
"""开播即计数:返回 should_play=True 时服务端已写下这一次,客户端必须真的播。
没配视频 / 开关关 / 次数用完 → should_play=False,客户端照旧走广告链路(行为不变)。
"""
result = crud_guide.start_play(db, user.id, scene=payload.scene or "coupon")
logger.info(
"guide video start user_id=%d token=%s status=%s seq=%d",
user.id, payload.play_token[:12], result["status"], result["seq"],
"guide video start user_id=%d scene=%s should_play=%s seq=%d remaining=%d",
user.id, payload.scene, result["should_play"], result["seq"], result["remaining"],
)
return GuideVideoStartOut(**result)
@@ -60,20 +50,15 @@ def start(
@router.post(
"/reward",
response_model=GuideVideoRewardOut,
dependencies=[Depends(rate_limit(120, 60, "guide-video-reward"))],
summary="引导视频发金币(播完/中途关闭都发,play_token 幂等)",
dependencies=[Depends(rate_limit(60, 60, "guide-video-reward"))],
)
def reward(
payload: GuideVideoRewardIn, user: CurrentUser, db: DbSession
) -> GuideVideoRewardOut:
result = crud_guide.grant_circle(
db,
user.id,
play_token=payload.play_token,
circle=payload.circle,
def reward(payload: GuideVideoRewardIn, user: CurrentUser, db: DbSession) -> GuideVideoRewardOut:
result = crud_guide.grant_play(
db, user.id, play_token=payload.play_token, completed=payload.completed
)
logger.info(
"guide video reward user_id=%d token=%s circle=%d status=%s granted=%s",
user.id, payload.play_token[:12], payload.circle,
result["status"], result["granted"],
"guide video reward user_id=%d token=%s completed=%s granted=%s coin=%d",
user.id, payload.play_token[:12], payload.completed, result["granted"], result["coin"],
)
return GuideVideoRewardOut(**result)
+3 -19
View File
@@ -79,8 +79,6 @@ class Settings(BaseSettings):
HONOR_PUSH_APP_ID: str = ""
HONOR_PUSH_CLIENT_ID: str = ""
HONOR_PUSH_CLIENT_SECRET: str = ""
# 0=正式消息(默认);1=测试消息(仅开发联调,勿用于生产)。
HONOR_PUSH_TARGET_USER_TYPE: int = Field(default=0, ge=0, le=1)
HONOR_PUSH_TOKEN_ENDPOINT: str = "https://iam.developer.honor.com/auth/token"
HONOR_PUSH_SEND_ENDPOINT_TEMPLATE: str = (
"https://push-api.cloud.honor.com/api/v1/{app_id}/sendMessage"
@@ -193,27 +191,14 @@ class Settings(BaseSettings):
# (real 模式下也跳过校验)、每次登录【强制重走新手引导】,并设【每日使用次数上限】防被人
# 猜到号后脚本滥用。两个值都能随时改 .env。逻辑全在 app/core/test_account.py,与其他业务解耦。
# ⚠️ TEST_ACCOUNT_PHONE 留空 = 整个功能关闭(生产默认安全;要启用才显式填号)。
TEST_ACCOUNT_PHONE: str = "" # 兼容旧配置:单个测试手机号
TEST_ACCOUNT_PHONES: str = "" # 多个测试手机号,英文逗号分隔
TEST_ACCOUNT_PHONE: str = "" # 测试手机号(11 位,如 11111111111);空=关闭整功能
TEST_ACCOUNT_DAILY_LIMIT: int = 500 # 该测试号每日最多登录次数,当日超过即拒绝登录
@property
def test_account_phone(self) -> str:
"""兼容旧调用:规整后的单个测试手机号。"""
"""规整后的测试手机号(去空白);空串=功能关闭"""
return self.TEST_ACCOUNT_PHONE.strip()
@property
def test_account_phones(self) -> frozenset[str]:
"""所有测试手机号;新旧配置取并集,便于线上平滑迁移。"""
phones = {
phone.strip()
for phone in self.TEST_ACCOUNT_PHONES.split(",")
if phone.strip()
}
if self.test_account_phone:
phones.add(self.test_account_phone)
return frozenset(phones)
# ===== 美团联盟 CPS =====
# 未配置时所有 /api/v1/meituan/* 接口 200 返空(优雅降级),不影响登录/领券等其他业务。
MT_CPS_APP_KEY: str = ""
@@ -453,7 +438,6 @@ class Settings(BaseSettings):
# 运营后台上传的新手引导视频上限。视频比图片大一个量级,单独一档;
# ⚠️ 改大时同步放宽网关 client_max_body_size(实测 QA 4MiB / prod 32MiB),否则 nginx 先挡下。
GUIDE_VIDEO_MAX_BYTES: int = 100 * 1024 * 1024 # 引导视频最大 100MB
FFPROBE_BINARY: str = "ffprobe"
# ===== 邀请好友 =====
# 分享落地页(二维码 / 分享链接指向这里;扫码 → 落地页 → 引导浏览器下载 APK)。
@@ -506,7 +490,7 @@ class Settings(BaseSettings):
return self.APP_ENV == "prod"
@model_validator(mode="after")
def _enforce_prod_secrets(self) -> Settings:
def _enforce_prod_secrets(self) -> "Settings":
"""prod 下强校验 JWT secret,弱/默认/空即启动报错(fail-fast,挡住 token 被伪造)。
只校验两个签发凭证:App 用户的 JWT_SECRET_KEY、后台的 ADMIN_JWT_SECRET——它们沿用默认值
+9 -76
View File
@@ -10,13 +10,8 @@
"""
from __future__ import annotations
import json
import os
import secrets
import subprocess
from decimal import ROUND_HALF_UP, Decimal, InvalidOperation
from pathlib import Path
from typing import Any
from app.core.config import settings
@@ -86,91 +81,29 @@ def _sniff_video_ext(data: bytes) -> str | None:
"""按魔数判定视频类型,返回扩展名;非支持类型返回 None。
只认 MP4 家族(ISO BMFF):`....ftyp` 在偏移 4。Android ExoPlayer 与浏览器 <video>
支持 H.264 或 HEVC/H.265 视频及 AAC 音频;放开 mkv/avi 只会让端上放不出来,
不如在入口就挡掉。
都稳吃 H.264/AAC 的 mp4;放开 mkv/avi 只会让端上放不出来,不如在入口就挡掉。
"""
if len(data) >= 12 and data[4:8] == b"ftyp":
return ".mp4"
return None
def _probe_guide_video(path: Path) -> dict[str, Any]:
"""用 ffprobe 读取服务端可信的时长/编码信息。"""
try:
proc = subprocess.run(
[
settings.FFPROBE_BINARY,
"-v", "error",
"-show_format",
"-show_streams",
"-of", "json",
str(path),
],
capture_output=True,
text=True,
timeout=30,
check=False,
)
except FileNotFoundError as exc:
raise MediaError("服务器未安装 ffprobe,暂时无法分析视频") from exc
except (OSError, subprocess.TimeoutExpired) as exc:
raise MediaError("ffprobe 分析视频失败") from exc
if proc.returncode != 0:
raise MediaError("无法分析视频,请确认文件是有效的 MP4")
try:
payload = json.loads(proc.stdout)
streams = payload.get("streams") or []
video = next(s for s in streams if s.get("codec_type") == "video")
audio = next((s for s in streams if s.get("codec_type") == "audio"), None)
raw_duration = (
(payload.get("format") or {}).get("duration")
or video.get("duration")
)
duration_ms = int(
(Decimal(str(raw_duration)) * 1000).quantize(Decimal("1"), rounding=ROUND_HALF_UP)
)
except (StopIteration, TypeError, ValueError, InvalidOperation) as exc:
raise MediaError("视频缺少可识别的视频轨或时长") from exc
video_codec = str(video.get("codec_name") or "").lower()
audio_codec = str(audio.get("codec_name") or "").lower() if audio else None
if not 30_000 <= duration_ms <= 180_000:
raise MediaError("视频时长必须在 30180 秒之间")
if video_codec not in ("h264", "hevc", "h265"):
raise MediaError("视频编码必须为 H.264 或 HEVC/H.265")
if audio_codec not in (None, "aac"):
raise MediaError("音频编码必须为 AAC")
return {
"duration_ms": duration_ms,
"video_codec": video_codec,
"audio_codec": audio_codec,
"analysis_status": "valid",
"analysis_error": None,
}
def save_guide_video(data: bytes) -> str:
"""保存新手引导视频(运营后台上传的运营素材),返回相对 URL(`/media/guide_video/<file>`)。
def save_guide_video(data: bytes) -> tuple[str, dict[str, Any]]:
"""临时落盘、ffprobe 校验后原子发布,返回 URL 与分析结果。"""
与图片分开一套校验:体积上限走 [settings.GUIDE_VIDEO_MAX_BYTES],类型只认 MP4。
"""
if not data:
raise MediaError("空文件")
limit = settings.GUIDE_VIDEO_MAX_BYTES
if len(data) > limit:
raise MediaError(f"视频过大(上限 {limit // (1024 * 1024)}MB)")
if _sniff_video_ext(data) is None:
raise MediaError("仅支持 MP4 视频(H.264 或 HEVC/H.265 编码)")
raise MediaError("仅支持 MP4 视频(H.264 编码)")
directory = _media_dir("guide_video")
stem = f"guide_{secrets.token_hex(8)}"
temp_path = directory / f".{stem}.upload"
final_path = directory / f"{stem}.mp4"
try:
temp_path.write_bytes(data)
analysis = _probe_guide_video(temp_path)
os.replace(temp_path, final_path)
except Exception:
temp_path.unlink(missing_ok=True)
final_path.unlink(missing_ok=True)
raise
return f"{settings.MEDIA_URL_PREFIX}/guide_video/{final_path.name}", analysis
fname = f"guide_{secrets.token_hex(8)}.mp4"
(_media_dir("guide_video") / fname).write_bytes(data)
return f"{settings.MEDIA_URL_PREFIX}/guide_video/{fname}"
def save_cps_image(admin_id: int, data: bytes) -> str:
+16 -15
View File
@@ -9,8 +9,9 @@
3. **每日使用次数上限**:防被人猜到这个号后写脚本一直刷。当天登录数超过上限即拒绝(429),
次日自动归零。
手机号与上限都在 .env 配(`TEST_ACCOUNT_PHONES` / `TEST_ACCOUNT_DAILY_LIMIT`),随时可改。
兼容旧的单号配置 `TEST_ACCOUNT_PHONE`;两项都留空 = 整个功能关闭(生产默认态)
手机号与上限都在 .env 配(`TEST_ACCOUNT_PHONE` / `TEST_ACCOUNT_DAILY_LIMIT`),随时可改。
`TEST_ACCOUNT_PHONE` 留空 = 整个功能关闭(生产默认态),`is_test_account()` 对任何号都返回
False,登录/短信回到原逻辑,零影响。
计数存**进程内存**(单 worker 够用,与 sms.py 同款约定):重启清零、多 worker 不共享。作为
一个测试号的粗粒度防滥用闸够用;且因所有登录都落同一个 phone → 同一个 user,滥用面天然只
@@ -26,9 +27,9 @@ from app.core.config import settings
logger = logging.getLogger("shagua.test_account")
# 进程内每日计数:{手机号: (date_str, 当日已登录次数)}。单 worker 有效,重启清零。
# 进程内每日计数:(date_str, 当日已登录次数)。单 worker 有效,重启清零(见模块 docstring)
_lock = Lock()
_usage: dict[str, tuple[str, int]] = {}
_usage: tuple[str, int] = ("", 0)
def _today() -> str:
@@ -36,16 +37,16 @@ def _today() -> str:
def is_enabled() -> bool:
"""功能总开关:至少配置了一个测试手机号才启用"""
return bool(settings.test_account_phones)
"""功能总开关:配了 TEST_ACCOUNT_PHONE 才启用(空=关闭)"""
return bool(settings.test_account_phone)
def is_test_account(phone: str) -> bool:
"""该手机号是否在测试账号集合中"""
return phone in settings.test_account_phones
"""该手机号是否为配置的测试账号。功能关闭时对任何号都返回 False"""
return is_enabled() and phone == settings.test_account_phone
def try_consume_quota(phone: str) -> bool:
def try_consume_quota() -> bool:
"""测试账号登录时调:当日计数 +1。
Returns:
@@ -58,17 +59,17 @@ def try_consume_quota(phone: str) -> bool:
limit = settings.TEST_ACCOUNT_DAILY_LIMIT
with _lock:
today = _today()
day, cnt = _usage.get(phone, ("", 0))
day, cnt = _usage
if day != today: # 跨天归零
cnt = 0
if cnt >= limit:
_usage[phone] = (today, cnt) # 已满,保持不变
_usage = (today, cnt) # 已满,保持不变
logger.warning(
"测试账号 %s 今日登录数已达上限 %d,拒绝", phone, limit
"测试账号 %s 今日登录数已达上限 %d,拒绝", settings.test_account_phone, limit
)
return False
_usage[phone] = (today, cnt + 1)
logger.info("测试账号 %s%d/%d 次登录", phone, cnt + 1, limit)
_usage = (today, cnt + 1)
logger.info("测试账号 %s%d/%d 次登录", settings.test_account_phone, cnt + 1, limit)
return True
@@ -76,4 +77,4 @@ def _reset_for_test() -> None:
"""仅供单测:清空进程内计数,隔离用例间状态。"""
global _usage
with _lock:
_usage = {}
_usage = ("", 0)
+1 -1
View File
@@ -35,7 +35,7 @@ _GC_THRESHOLD = 10000 # 超此阈值,send 时顺手清老于
# 发码错误码 → (HTTP 码, 用户提示)。未列出的一律 503(供应商不可用)。
_SEND_ERRORS: dict[str, tuple[int, str]] = {
"MOBILE_NUMBER_ILLEGAL": (400, "请输入有效的手机号"),
"MOBILE_NUMBER_ILLEGAL": (400, "手机号无效"),
"BUSINESS_LIMIT_CONTROL": (429, "今日发送次数过多,请明天再试"),
"FREQUENCY_FAIL": (429, "发送过于频繁,请稍后再试"),
}
+1 -1
View File
@@ -53,7 +53,7 @@ _GC_THRESHOLD = 10000 # 任一内存 dict 超此阈值,send 时顺手清过期
# 发码错误码(创蓝 `code`)→ (HTTP 码, 用户提示)。未列出的一律 503(供应商不可用)。
_SEND_ERRORS: dict[str, tuple[int, str]] = {
"103": (429, "发送过于频繁,请稍后再试"), # 提交速度过快
"107": (400, "请输入有效的手机号"), # 手机号码错误
"107": (400, "手机号无效"), # 手机号码错误
}
# 需运维介入的配置/开通/余额类错误:打 critical 日志(仍归 503)。
_SEND_CRITICAL_CODES = frozenset({
+1 -1
View File
@@ -185,5 +185,5 @@ def _send_via_jiguang(phone: str, code: str) -> None:
if ecode == 50009: # 极光侧超频
raise SmsError("发送过于频繁,请稍后再试", status_code=429)
if ecode == 50006: # 手机号无效(schema 已挡格式,这里多是空号/停机)
raise SmsError("请输入有效的手机号", status_code=400)
raise SmsError("手机号无效", status_code=400)
raise SmsError(f"短信发送失败(code={ecode})", status_code=503)
+11 -39
View File
@@ -224,10 +224,6 @@ def send_data_event(
if vendor == "vivo":
return {"skipped": True, "reason": "foreground notification callback"}
# 小米服务端不再发送透传消息;通知栏消息仍照常发送,前台角标由客户端主动刷新。
if vendor == "xiaomi":
return {"skipped": True, "reason": "xiaomi data messages disabled"}
# OPush 当前只支持通知栏消息,没有服务端透传单推接口。
# 客户端在 OPPO 首页可见时轮询未读数;这里必须安全跳过,不能请求不存在的
# /message/transparent/unicast(该地址会稳定返回 HTTP 404)。
@@ -237,6 +233,7 @@ def send_data_event(
dispatch: dict[str, Callable[[str, dict[str, str]], dict[str, Any]]] = {
"honor": _send_honor_data,
"huawei": _send_huawei_data,
"xiaomi": _send_xiaomi_data,
}
return dispatch[vendor](token, payload)
@@ -350,25 +347,6 @@ def _response_summary(resp: Any, parsed: Any | None = None) -> str:
return _raw_log_summary(getattr(resp, "text", ""))
def _vendor_response_failed(vendor: str, data: dict[str, Any]) -> bool:
"""识别 HTTP 200 中明确的厂商业务失败,避免将失败请求记录成 success。"""
if data.get("error") or data.get("success") is False:
return True
if vendor == "xiaomi":
code = data.get("code")
result = str(data.get("result", "ok")).lower()
return code not in (0, "0", None) or result not in ("ok", "success")
if vendor == "oppo" and data.get("code") is not None:
return int(data["code"]) != 0
if vendor == "vivo" and data.get("result") is not None:
return int(data["result"]) != 0
if vendor == "honor" and data.get("code") is not None:
return int(data["code"]) != 200
if vendor == "huawei" and data.get("code") is not None:
return str(data["code"]) != "80000000"
return False
def _request_json(
method: str,
url: str,
@@ -432,16 +410,13 @@ def _request_json(
_elapsed_ms(started),
)
raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e
vendor_failed = _vendor_response_failed(vendor, data)
log = logger.warning if vendor_failed else logger.info
log(
logger.info(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=%s http_status=%s request=%s response=%s elapsed_ms=%.1f",
"outcome=success http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
"vendor_error" if vendor_failed else "success",
resp.status_code,
request_summary,
_response_summary(resp, data),
@@ -523,7 +498,7 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di
"notification": {"title": title, "body": body},
"android": {
"ttl": f"{settings.PUSH_TIME_TO_LIVE_SEC}s",
"targetUserType": settings.HONOR_PUSH_TARGET_USER_TYPE,
"targetUserType": 1,
"notification": {
"title": title,
"body": body,
@@ -817,21 +792,18 @@ def _send_xiaomi(token: str, title: str, body: str, extras: dict[str, str]) -> d
def _send_xiaomi_data(token: str, payload: dict[str, str]) -> dict[str, Any]:
app_secret = _require(settings.XIAOMI_PUSH_APP_SECRET, "XIAOMI_PUSH_APP_SECRET")
form = {
"registration_id": token,
"restricted_package_name": settings.ANDROID_PACKAGE_NAME,
"payload": json.dumps(payload, ensure_ascii=False),
"pass_through": "1",
"time_to_live": str(settings.PUSH_TIME_TO_LIVE_SEC * 1000),
}
if settings.XIAOMI_PUSH_CHANNEL_ID.strip():
form["extra.channel_id"] = settings.XIAOMI_PUSH_CHANNEL_ID.strip()
data = _request_form(
"POST",
settings.XIAOMI_PUSH_SEND_ENDPOINT,
vendor="xiaomi",
operation="send_data_event",
data=form,
data={
"registration_id": token,
"restricted_package_name": settings.ANDROID_PACKAGE_NAME,
"payload": json.dumps(payload, ensure_ascii=False),
"pass_through": "1",
"time_to_live": str(settings.PUSH_TIME_TO_LIVE_SEC * 1000),
},
headers={"Authorization": f"key={app_secret}"},
)
code = data.get("code")
-10
View File
@@ -34,11 +34,6 @@ class DeviceLiveness(Base):
__tablename__ = "device_liveness"
__table_args__ = (
UniqueConstraint("user_id", "device_id", name="uq_device_liveness_user_device"),
UniqueConstraint(
"push_vendor",
"push_token",
name="uq_device_liveness_vendor_token",
),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
@@ -53,11 +48,6 @@ class DeviceLiveness(Base):
push_vendor: Mapped[str | None] = mapped_column(String(32), nullable=True)
# 厂商 push token / regId / registration_id;不同厂商命名不同,后端统一存这里。
push_token: Mapped[str | None] = mapped_column(String(256), nullable=True)
# 客户端每次登录生成的高熵随机值。离线退出后只凭此值精确撤销旧会话绑定,
# 新账号/新会话注册时会替换它,故旧解绑任务绝不能误解绑新登录。
push_binding_id: Mapped[str | None] = mapped_column(String(128), index=True, nullable=True)
# 被退出/转移的 binding 留作墓碑,阻止旧的在途 register/heartbeat 把 token 抢回。
push_binding_revoked: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
platform: Mapped[str] = mapped_column(String(16), nullable=False, default="android")
app_version: Mapped[str | None] = mapped_column(String(32), nullable=True)
+41 -25
View File
@@ -1,4 +1,20 @@
"""引导视频播放计划、起播快照与逐圈结算状态。"""
"""新手引导视频播放记录(领券浮层前 N 次用它替代广告)。
产品规则(2026-07 拍板):新用户点一键自动领取后的等候浮层,** 3 **不放广告,
改放运营后台上传的引导视频;每次固定 120 金币,中途关闭也算看完照发
口径:
- **计次按账号**(user_id),与设备无关 换设备不重新送 3
- **开播即计数**:客户端每次要展示浮层时调 `/api/v1/guide-video/start`,服务端当场
写一行(status='playing')并返回 play_token;`COUNT(*)` 即已用次数用户中途 kill
App 也算用掉一次(产品选定口径,防反复进出刷金币)
- **发币幂等** play_token 定位 + `status='playing'` 条件更新:并发两次上报只有一次
改到行(另一次 rowcount=0),所以只发一次币光有 play_token 唯一键挡不住 发币走的是
UPDATE, INSERT,撞不到任何唯一键
- **次数上限** (user_id, seq) 唯一键兜底,防并发 /start 绕过 COUNT 判定(见下)
与广告收益(ad_feed_reward_record)彻底分离:引导视频不是广告,不该进广告收益报表
"""
from __future__ import annotations
from datetime import datetime
@@ -10,48 +26,48 @@ from app.db.base import Base
class GuideVideoPlay(Base):
"""一次引导视频播放一行。开播时建(status='playing'),发币后置 'granted'"""
__tablename__ = "guide_video_play"
__table_args__ = (
# 客户端幂等键:同一次播放重复上报奖励只发一次。
UniqueConstraint("play_token", name="uq_guide_video_play_token"),
Index(
"uq_guide_video_play_user_scene_seq",
"user_id",
"scene",
"seq",
unique=True,
),
# 次数上限的**硬约束**:start_play 是无锁 check-then-insert(读 COUNT 算 seq 再插),
# N 个并发 /start 会都读到同一个已用次数、算出同一个 seq,不拦就能各拿一个 token、
# 各发一次金币,3 次上限形同虚设(改包即可无限刷)。seq 唯一 → 并发同 seq 必撞,
# start_play 捕获 IntegrityError 降级成"这次不放视频"。
# 用 unique Index 而非 UniqueConstraint:与迁移里的 create_index 对齐(SQLite 加约束
# 要整表重建),autogenerate 才不会每次报一条假 diff。
Index("uq_guide_video_play_user_scene_seq", "user_id", "scene", "seq", unique=True),
)
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
)
# 服务端生成下发给客户端的幂等键(uuid hex)。
play_token: Mapped[str] = mapped_column(String(64), nullable=False)
# 触发场景:目前只有 coupon(领券等候浮层);留字段以便日后比价等场景复用。
scene: Mapped[str] = mapped_column(String(16), nullable=False, default="coupon")
# prepare 不占次数,seq=NULLstart 成功才写入 1-based seq
seq: Mapped[int | None] = mapped_column(Integer, nullable=True)
# 本账号第几次(1-based),= 建行时已有行数 + 1。日常判定仍以 COUNT 为准,但 (user_id, seq)
# 唯一键让并发 /start 只能成一个 —— 见 __table_args__。
seq: Mapped[int] = mapped_column(Integer, nullable=False, default=1)
# 当次下发的视频地址(运营换片后能回溯用户当时看的是哪支)。
video_url: Mapped[str | None] = mapped_column(String(512), nullable=True)
# 实发金币;未发时 0。
coin: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
duration_ms: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
config_version: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
settled_circles: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# prepared / started / completed / legacy_completed / legacy_closed
status: Mapped[str] = mapped_column(String(24), nullable=False, default="prepared")
# playing(已开播未发币) / granted(已发币)。
status: Mapped[str] = mapped_column(String(16), nullable=False, default="playing")
# 客户端上报时是否播完(true=自然播完 / false=中途关闭)。仅留痕:两者都发币。
completed: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
prepared_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
expires_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), index=True, nullable=True
)
started_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True), index=True, nullable=True
started_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), index=True, nullable=False
)
granted_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
def __repr__(self) -> str: # pragma: no cover
return (
f"<GuideVideoPlay user={self.user_id} scene={self.scene} seq={self.seq} "
f"{self.status} circles={self.settled_circles}>"
f"<GuideVideoPlay user={self.user_id} seq={self.seq} "
f"{self.status} coin={self.coin}>"
)
-9
View File
@@ -60,15 +60,6 @@ class CoinTransaction(Base):
sqlite_where=text("biz_type LIKE 'task%' AND ref_id IS NOT NULL"),
postgresql_where=text("biz_type LIKE 'task%' AND ref_id IS NOT NULL"),
),
Index(
"ux_coin_transaction_guide_video_ref",
"user_id",
"biz_type",
"ref_id",
unique=True,
sqlite_where=text("biz_type = 'guide_video' AND ref_id IS NOT NULL"),
postgresql_where=text("biz_type = 'guide_video' AND ref_id IS NOT NULL"),
),
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
-10
View File
@@ -375,16 +375,6 @@ def _derive_from_platforms(
rows = [p for p in (platforms or []) if isinstance(p, dict)]
src = next((p for p in rows if p.get("role") == "source"), None)
best = next((p for p in rows if p.get("is_best")), None)
if best is None:
# pricebot 没标 is_best(如全平台 has_dish_diff「相似替换/仅供参考」→ 不认定权威最低价)
# 但仍有有价行 → 兜底取有价行里最低价当参考 best,避免记录级 best_*/saved 整条落 NULL
# (否则首页价 0.00 / 记录页无最低红框 / 省额丢失 /「累计发现可省」漏计)。**候选含源**:
# 源常年全菜、价可信,源本身最便宜时 best 回落源(saved=0、is_source_best=True),与 _derive
# 「全目标缺菜回落源、不虚报省」同一语义;若排除源强选更贵目标,saved 会变负、倒扣「累计
# 发现可省」(get_stats 对 saved_amount_cents 求和不带 >0 过滤)。
priced = [p for p in rows if p.get("price") is not None]
if priced:
best = min(priced, key=lambda p: (p["price"], p.get("display_order") or 0))
source_price_cents = _yuan_to_cents(src.get("price")) if src else None
best_price_cents = _yuan_to_cents(best.get("price")) if best else None
saved_amount_cents = None
+14 -117
View File
@@ -1,9 +1,9 @@
"""device 表读写(设备注册 / 心跳 / 超时扫描)。"""
from __future__ import annotations
from datetime import UTC, datetime, timedelta
from datetime import datetime, timedelta, timezone
from sqlalchemy import and_, or_, select
from sqlalchemy import select
from sqlalchemy.orm import Session
from app.models.device import DeviceLiveness
@@ -24,29 +24,13 @@ def register_or_update(
registration_id: str | None = None,
push_vendor: str | None = None,
push_token: str | None = None,
push_binding_id: str | None = None,
platform: str = "android",
app_version: str | None = None,
) -> DeviceLiveness:
"""注册设备或更新推送目标;同一设备/token 的旧用户绑定在同一事务内失效"""
"""注册设备或更新其厂商 push token / 元信息。upsert by (user_id, device_id)"""
normalized_vendor = _normalize_push_vendor(push_vendor)
normalized_token = push_token.strip() if push_token else None
normalized_binding = push_binding_id.strip() if push_binding_id else None
device = _get(db, user_id=user_id, device_id=device_id)
stale_binding = bool(
device
and normalized_binding
and device.push_binding_id == normalized_binding
and device.push_binding_revoked
)
if not stale_binding:
_release_conflicting_push_bindings(
db,
user_id=user_id,
device_id=device_id,
push_vendor=normalized_vendor,
push_token=normalized_token,
)
if device is None:
device = DeviceLiveness(
user_id=user_id,
@@ -54,22 +38,17 @@ def register_or_update(
registration_id=registration_id,
push_vendor=normalized_vendor,
push_token=normalized_token,
push_binding_id=normalized_binding,
push_binding_revoked=False,
platform=platform or "android",
app_version=app_version,
)
db.add(device)
else:
if registration_id and not stale_binding:
if registration_id:
device.registration_id = registration_id
if normalized_vendor and not stale_binding:
if normalized_vendor:
device.push_vendor = normalized_vendor
if normalized_token and not stale_binding:
if normalized_token:
device.push_token = normalized_token
if normalized_binding and not stale_binding:
device.push_binding_id = normalized_binding
device.push_binding_revoked = False
if platform:
device.platform = platform
if app_version:
@@ -88,45 +67,26 @@ def touch_heartbeat(
registration_id: str | None = None,
push_vendor: str | None = None,
push_token: str | None = None,
push_binding_id: str | None = None,
) -> DeviceLiveness:
"""处理一次心跳(心跳也能自注册)。
service 心跳或 accessibility_enabled=true ,刷新存活并把状态机重置回 alive
清掉 notified_at(掉线恢复 下次再断才会再推一条)
"""
now = datetime.now(UTC)
now = datetime.now(timezone.utc)
device = _get(db, user_id=user_id, device_id=device_id)
normalized_vendor = _normalize_push_vendor(push_vendor)
normalized_token = push_token.strip() if push_token else None
normalized_binding = push_binding_id.strip() if push_binding_id else None
stale_binding = bool(
device
and normalized_binding
and device.push_binding_id == normalized_binding
and device.push_binding_revoked
)
if (normalized_token or normalized_binding) and not stale_binding:
_release_conflicting_push_bindings(
db,
user_id=user_id,
device_id=device_id,
push_vendor=normalized_vendor,
push_token=normalized_token,
)
if device is None:
device = DeviceLiveness(user_id=user_id, device_id=device_id)
db.add(device)
if registration_id and not stale_binding:
if registration_id:
device.registration_id = registration_id
if normalized_vendor and not stale_binding:
normalized_vendor = _normalize_push_vendor(push_vendor)
normalized_token = push_token.strip() if push_token else None
if normalized_vendor:
device.push_vendor = normalized_vendor
if normalized_token and not stale_binding:
if normalized_token:
device.push_token = normalized_token
if normalized_binding and not stale_binding:
device.push_binding_id = normalized_binding
device.push_binding_revoked = False
device.last_report_protection_on = accessibility_enabled
if accessibility_enabled:
@@ -147,7 +107,7 @@ def list_overdue(db: Session, *, timeout_minutes: int) -> list[DeviceLiveness]:
即使没有厂商 token 也要检出,后续由 kill_alert_pending 走客户端进 App 后兜底提醒
"""
cutoff = datetime.now(UTC) - timedelta(minutes=timeout_minutes)
cutoff = datetime.now(timezone.utc) - timedelta(minutes=timeout_minutes)
stmt = select(DeviceLiveness).where(
DeviceLiveness.ever_protected.is_(True),
DeviceLiveness.liveness_state == "alive",
@@ -166,7 +126,7 @@ def mark_notified(db: Session, *, device_id_pk: int) -> None:
device = db.get(DeviceLiveness, device_id_pk)
if device is not None:
device.liveness_state = "notified"
device.notified_at = datetime.now(UTC)
device.notified_at = datetime.now(timezone.utc)
device.kill_alert_pending = True
db.commit()
@@ -217,69 +177,6 @@ def list_push_targets(db: Session, *, user_id: int) -> list[DeviceLiveness]:
return targets
def unregister_push_binding(db: Session, *, user_id: int, device_id: str) -> None:
"""当前用户主动退出:清推送归属但保留设备存活历史。"""
device = _get(db, user_id=user_id, device_id=device_id)
if device is not None:
_clear_push_binding(device)
db.commit()
def unregister_push_binding_deferred(
db: Session, *, device_id: str, push_binding_id: str
) -> None:
"""离线退出补偿;只撤销完全匹配的旧登录会话,且始终幂等。"""
stmt = select(DeviceLiveness).where(
DeviceLiveness.device_id == device_id,
DeviceLiveness.push_binding_id == push_binding_id,
)
device = db.execute(stmt).scalar_one_or_none()
if device is not None:
_clear_push_binding(device)
db.commit()
def _release_conflicting_push_bindings(
db: Session,
*,
user_id: int,
device_id: str,
push_vendor: str | None,
push_token: str | None,
) -> None:
"""把设备或 token 从其他记录转移走,当前 (user,device) 行除外。"""
conflicts = [DeviceLiveness.device_id == device_id]
if push_vendor and push_token:
conflicts.append(
and_(
DeviceLiveness.push_vendor == push_vendor,
DeviceLiveness.push_token == push_token,
)
)
stmt = select(DeviceLiveness).where(
or_(*conflicts),
~and_(
DeviceLiveness.user_id == user_id,
DeviceLiveness.device_id == device_id,
),
)
released = False
for old in db.execute(stmt).scalars():
_clear_push_binding(old)
released = True
if released:
# 唯一约束下必须先落旧行清理,再把同一 token 写给新行。
db.flush()
def _clear_push_binding(device: DeviceLiveness) -> None:
device.registration_id = None
device.push_vendor = None
device.push_token = None
# binding_id 留作撤销墓碑;同一旧会话的在途请求不能重新认领 token。
device.push_binding_revoked = bool(device.push_binding_id)
def _normalize_push_vendor(push_vendor: str | None) -> str | None:
if not push_vendor:
return None
+187 -370
View File
@@ -1,20 +1,33 @@
"""双场景引导视频配置与 prepare/start/reward 状态机。"""
"""新手引导视频:运营配置读写 + 播放计次 + 发币。
**配置**(开关 / 视频地址 / 前几次 / 每次金币)整体作为一个 JSON 存进通用 app_config
(key=coupon_guide_video),写法完全对齐 feedback_qr 不进 CONFIG_DEFS,所以不会污染
系统配置页的通用列表,由本模块独占维护
**计次**按账号(user_id)**开播即计数**:客户端每次要展示领券等候浮层时调
`/api/v1/guide-video/start`,命中则当场写一行 guide_video_play(status='playing')
已用次数 = 该账号的行数,达到 max_plays 后不再下发,客户端改放广告(原逻辑)
COUNT 判定本身无锁,真正卡住次数上限的是 (user_id, seq) 唯一键:并发 /start 只能成一个
**发币**幂等键是 play_token,落地方式是 `status='playing' 'granted'` **条件更新**:
同一次播放重复上报只入账一次(网络重试 / 关闭与播完同时触发都靠它挡住)
中途关闭也照发 产品拍板中途关闭也算看完
两处都是直接铸币的路径,改动前先看 `start_play` / `grant_play` 上的并发注释
"""
from __future__ import annotations
import uuid
from datetime import datetime, timedelta
from pathlib import Path
from datetime import datetime
from typing import Any
from sqlalchemy import func, select, update
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from app.core import media, rewards
from app.core.config import settings
from app.core import rewards
from app.models.app_config import AppConfig
from app.models.guide_video import GuideVideoPlay
from app.models.wallet import CoinAccount
from app.repositories import wallet as crud_wallet
SCENES = ("coupon", "comparison")
@@ -22,47 +35,38 @@ _KEY_BY_SCENE = {
"coupon": "coupon_guide_video",
"comparison": "comparison_guide_video",
}
#: 金币流水 biz_type。客户端收益明细按它显示「新手引导视频奖励」。
BIZ_TYPE = "guide_video"
CIRCLE_COUNT = 10
PLAN_TTL = timedelta(minutes=10)
MIN_PLAYS = 1
# 默认值 = 「运营还没配」时的行为:video_url 为空 → 一律不下发引导视频,浮层维持现状(放广告)。
# 所以本功能上线后**不配视频就等于没上线**,不会影响存量用户。
_DEFAULTS: dict[str, Any] = {
"enabled": True,
"video_url": None, # None/空 = 未配片 → 不下发,浮层照旧放广告
"max_plays": 3, # 每个账号前 N 次浮层放引导视频
"reward_coin": 100, # 每次固定金币
}
_FIELDS = tuple(_DEFAULTS.keys())
# 后台可配范围的护栏:防手滑把次数/金币填成天文数字(配置直接决定发币)。
MAX_PLAYS_LIMIT = 50
MIN_REWARD_COIN = 10
REWARD_COIN_LIMIT = 10_000
_DEFAULTS: dict[str, Any] = {
"enabled": False,
"video_url": None,
"max_plays": 3,
"reward_coin": 100,
"duration_ms": None,
"video_codec": None,
"audio_codec": None,
"analysis_status": "missing",
"analysis_error": None,
"config_version": 0,
}
_FIELDS = tuple(_DEFAULTS)
# ===== 配置 =====
class PlayStateError(Exception):
def __init__(self, code: str, message: str, *, reprepare_required: bool = False):
super().__init__(message)
self.code = code
self.message = message
self.reprepare_required = reprepare_required
def detail(self) -> dict[str, Any]:
return {
"code": self.code,
"message": self.message,
"retryable": False,
"reprepare_required": self.reprepare_required,
}
def _now() -> datetime:
return datetime.now(rewards.CN_TZ).replace(tzinfo=None)
def _merge(raw: Any) -> dict[str, Any]:
"""DB 里(可能不全的)dict 叠加到默认上,得到完整配置(4 个字段,无 updated_at)。"""
out = dict(_DEFAULTS)
if isinstance(raw, dict):
for k in _FIELDS:
v = raw.get(k)
if v is not None:
out[k] = v
return out
def _config_key(scene: str) -> str:
@@ -71,60 +75,34 @@ def _config_key(scene: str) -> str:
return _KEY_BY_SCENE[scene]
def _merge(raw: Any) -> dict[str, Any]:
out = dict(_DEFAULTS)
if isinstance(raw, dict):
for key in _FIELDS:
if key in raw and raw[key] is not None:
out[key] = raw[key]
return out
def _public_config(cfg: dict[str, Any]) -> dict[str, Any]:
duration = int(cfg.get("duration_ms") or 0)
reward = int(cfg.get("reward_coin") or 0)
return {
**cfg,
"duration_ms": duration or None,
"circle_count": CIRCLE_COUNT,
"circle_duration_ms": duration / CIRCLE_COUNT if duration else None,
"reward_per_circle": reward // CIRCLE_COUNT,
}
def get_config(db: Session, scene: str = "coupon") -> dict[str, Any]:
"""完整配置 + updated_at(admin 读 / 业务读共用)。"""
row = db.get(AppConfig, _config_key(scene))
cfg = _public_config(_merge(row.value if row is not None else None))
cfg = _merge(row.value if row is not None else None)
cfg["updated_at"] = row.updated_at.isoformat() if row is not None and row.updated_at else None
return cfg
def _validate_settings(max_plays: int, reward_coin: int) -> None:
if not MIN_PLAYS <= max_plays <= MAX_PLAYS_LIMIT:
raise ValueError(f"播放次数必须在 {MIN_PLAYS}{MAX_PLAYS_LIMIT} 之间")
if not MIN_REWARD_COIN <= reward_coin <= REWARD_COIN_LIMIT:
raise ValueError(f"金币总价必须在 {MIN_REWARD_COIN}{REWARD_COIN_LIMIT} 之间")
if reward_coin % CIRCLE_COUNT:
raise ValueError("金币总价必须是 10 的倍数")
def _write(
db: Session, value: dict[str, Any], *, scene: str, admin_id: int, commit: bool
) -> dict[str, Any]:
"""整体覆写该行(value 须为完整字段 dict),返回合并后的完整配置(含 updated_at)。"""
key = _config_key(scene)
row = db.get(AppConfig, key)
if row is None:
row = AppConfig(key=key, value=value, updated_by_admin_id=admin_id)
db.add(row)
else:
row.value = value
row.value = value # 整体重新赋值,SQLAlchemy 才侦测得到变更
row.updated_by_admin_id = admin_id
if commit:
db.commit()
db.refresh(row)
else:
db.flush()
return get_config(db, scene)
out = _merge(row.value)
out["updated_at"] = row.updated_at.isoformat() if row.updated_at else None
return out
def update_config(
@@ -137,360 +115,199 @@ def update_config(
admin_id: int,
commit: bool = True,
) -> tuple[dict[str, Any], dict[str, Any]]:
before = get_config(db, scene)
raw = {key: before[key] for key in _FIELDS}
candidate_plays = int(max_plays if max_plays is not None else raw["max_plays"])
candidate_reward = int(reward_coin if reward_coin is not None else raw["reward_coin"])
_validate_settings(candidate_plays, candidate_reward)
if enabled is True and (
not raw.get("video_url")
or raw.get("analysis_status") != "valid"
or not raw.get("duration_ms")
):
raise ValueError("请先上传并通过分析的视频,再启用")
"""改开关 / 次数 / 金币(只改传了的字段;视频走 set_video)。返回 (before, after) 供审计。"""
row = db.get(AppConfig, _config_key(scene))
before = _merge(row.value if row is not None else None)
new_value = {k: before[k] for k in _FIELDS}
if enabled is not None:
raw["enabled"] = enabled
raw["max_plays"] = candidate_plays
raw["reward_coin"] = candidate_reward
raw["config_version"] = int(raw.get("config_version") or 0) + 1
after = _write(db, raw, scene=scene, admin_id=admin_id, commit=commit)
new_value["enabled"] = enabled
if max_plays is not None:
new_value["max_plays"] = max(0, min(int(max_plays), MAX_PLAYS_LIMIT))
if reward_coin is not None:
new_value["reward_coin"] = max(0, min(int(reward_coin), REWARD_COIN_LIMIT))
after = _write(db, new_value, scene=scene, admin_id=admin_id, commit=commit)
return before, after
def set_video(
db: Session,
video_url: str | None,
*,
analysis: dict[str, Any] | None = None,
scene: str = "coupon",
admin_id: int,
commit: bool = True,
db: Session, video_url: str | None, *, scene: str = "coupon",
admin_id: int, commit: bool = True
) -> tuple[dict[str, Any], dict[str, Any]]:
before = get_config(db, scene)
raw = {key: before[key] for key in _FIELDS}
raw["video_url"] = video_url
if video_url:
if not analysis or analysis.get("analysis_status") != "valid":
raise ValueError("视频必须先完成服务端分析")
for key in (
"duration_ms", "video_codec", "audio_codec",
"analysis_status", "analysis_error",
):
raw[key] = analysis.get(key)
else:
raw.update(
duration_ms=None,
video_codec=None,
audio_codec=None,
analysis_status="missing",
analysis_error=None,
)
raw["enabled"] = False
raw["config_version"] = int(raw.get("config_version") or 0) + 1
after = _write(db, raw, scene=scene, admin_id=admin_id, commit=commit)
"""设置/清空引导视频地址。返回 (before, after);before['video_url'] 供调用方删旧文件。"""
row = db.get(AppConfig, _config_key(scene))
before = _merge(row.value if row is not None else None)
new_value = {k: before[k] for k in _FIELDS}
new_value["video_url"] = video_url
after = _write(db, new_value, scene=scene, admin_id=admin_id, commit=commit)
return before, after
# ===== 播放计次 =====
def used_plays(db: Session, user_id: int, scene: str = "coupon") -> int:
"""该账号已用掉的引导视频次数(开播即算,含未发币的)。"""
return int(
db.execute(
select(func.count()).select_from(GuideVideoPlay).where(
GuideVideoPlay.user_id == user_id,
GuideVideoPlay.scene == scene,
GuideVideoPlay.status != "prepared",
)
).scalar_one()
)
def _prepare_miss(scene: str, reason: str, cfg: dict[str, Any], used: int) -> dict[str, Any]:
maximum = int(cfg.get("max_plays") or 0)
return {
"should_play": False,
"reason": reason,
"scene": scene,
"video_url": None,
"play_token": "",
"config_version": int(cfg.get("config_version") or 0),
"duration_ms": 0,
"circle_count": CIRCLE_COUNT,
"circle_duration_ms": 0,
"reward_coin": int(cfg.get("reward_coin") or 0),
"reward_per_circle": int(cfg.get("reward_coin") or 0) // CIRCLE_COUNT,
"seq": used,
"remaining": max(0, maximum - used),
"expires_at": None,
}
def play_stats(db: Session, scene: str = "coupon") -> dict[str, int]:
"""全站播放统计(admin 页展示):总播放次数 / 其中已发币次数。"""
total = int(
db.execute(
select(func.count()).select_from(GuideVideoPlay).where(
GuideVideoPlay.scene == scene
)
).scalar_one()
)
granted = int(
db.execute(
select(func.count()).select_from(GuideVideoPlay).where(
GuideVideoPlay.status == "granted",
GuideVideoPlay.scene == scene,
)
).scalar_one()
)
return {"total_plays": total, "granted_plays": granted}
def prepare_play(db: Session, user_id: int, *, scene: str = "coupon") -> dict[str, Any]:
def start_play(
db: Session, user_id: int, *, scene: str = "coupon", commit: bool = True
) -> dict[str, Any]:
"""决定这次浮层是否放引导视频;命中则**当场计次**并返回 play_token。
返回 dict:
should_play 是否放引导视频(False 客户端照旧放广告)
video_url 相对地址(/media/...);客户端自行拼 BASE_URL
play_token 发币幂等键(should_play=False 时为空串)
reward_coin 播完/中途关闭都发的固定金币
seq / remaining 第几次 / 发完这次还剩几次(仅展示与排查用)
"""
_config_key(scene)
cfg = get_config(db, scene)
video_url = (cfg.get("video_url") or "").strip()
max_plays = int(cfg.get("max_plays") or 0)
reward_coin = int(cfg.get("reward_coin") or 0)
used = used_plays(db, user_id, scene)
video_url = str(cfg.get("video_url") or "").strip()
duration = int(cfg.get("duration_ms") or 0)
maximum = int(cfg.get("max_plays") or 0)
if not cfg.get("enabled"):
return _prepare_miss(scene, "disabled", cfg, used)
if not video_url or cfg.get("analysis_status") != "valid" or duration <= 0:
return _prepare_miss(scene, "video_unavailable", cfg, used)
if used >= maximum:
return _prepare_miss(scene, "play_limit_reached", cfg, used)
now = _now()
def _miss(used_now: int) -> dict[str, Any]:
return {
"should_play": False,
"video_url": None,
"play_token": "",
"reward_coin": reward_coin,
"seq": used_now,
"remaining": max(0, max_plays - used_now),
}
if not cfg.get("enabled") or not video_url or max_plays <= 0 or used >= max_plays:
return _miss(used)
seq = used + 1
play = GuideVideoPlay(
user_id=user_id,
play_token=uuid.uuid4().hex,
scene=scene,
seq=None,
seq=seq,
video_url=video_url,
coin=int(cfg["reward_coin"]),
duration_ms=duration,
config_version=int(cfg["config_version"]),
settled_circles=0,
status="prepared",
coin=reward_coin,
status="playing",
completed=0,
prepared_at=now,
expires_at=now + PLAN_TTL,
started_at=None,
started_at=datetime.now(rewards.CN_TZ).replace(tzinfo=None),
)
db.add(play)
db.commit()
# 上面的 COUNT 判定是无锁 check-then-insert:并发 /start 会都算出同一个 seq。
# (user_id, seq) 唯一键让只有一个能落库,其余撞键 → 回滚后按"这次不放视频"降级,
# 客户端照旧走广告链路。没有它,并发就能绕过 max_plays 无限刷金币。
try:
if commit:
db.commit()
else:
db.flush()
except IntegrityError:
db.rollback()
return _miss(used_plays(db, user_id, scene))
return {
"should_play": True,
"reason": "ready",
"scene": scene,
"video_url": video_url,
"play_token": play.play_token,
"config_version": play.config_version,
"duration_ms": duration,
"circle_count": CIRCLE_COUNT,
"circle_duration_ms": duration / CIRCLE_COUNT,
"reward_coin": play.coin,
"reward_per_circle": play.coin // CIRCLE_COUNT,
"seq": used + 1,
"remaining": max(0, maximum - used - 1),
"expires_at": play.expires_at.isoformat(),
"reward_coin": reward_coin,
"seq": seq,
"remaining": max(0, max_plays - seq),
}
def _find_play(db: Session, user_id: int, token: str) -> GuideVideoPlay | None:
"""按 (play_token, user_id) 取播放行 —— 带 user_id 是防拿别人的 token 来兑。"""
return db.execute(
select(GuideVideoPlay).where(
GuideVideoPlay.play_token == token.strip(),
GuideVideoPlay.play_token == token,
GuideVideoPlay.user_id == user_id,
)
).scalar_one_or_none()
def _start_out(play: GuideVideoPlay, maximum: int, status: str) -> dict[str, Any]:
assert play.started_at is not None and play.seq is not None and play.video_url
return {
"started": True,
"status": status,
"play_token": play.play_token,
"scene": play.scene,
"video_url": play.video_url,
"config_version": play.config_version,
"duration_ms": play.duration_ms,
"circle_count": CIRCLE_COUNT,
"circle_duration_ms": play.duration_ms / CIRCLE_COUNT,
"reward_coin": play.coin,
"reward_per_circle": play.coin // CIRCLE_COUNT,
"seq": play.seq,
"remaining": max(0, maximum - play.seq),
"started_at": play.started_at.isoformat(),
}
def start_play(db: Session, user_id: int, *, play_token: str) -> dict[str, Any]:
play = _find_play(db, user_id, play_token)
if play is None:
raise PlayStateError("play_not_found", "播放计划不存在")
cfg = get_config(db, play.scene)
maximum = int(cfg["max_plays"])
if play.status in {"started", "completed"}:
return _start_out(play, maximum, "already_started")
if play.status != "prepared":
raise PlayStateError("play_not_found", "播放计划不可用")
now = _now()
if play.expires_at is None or now > play.expires_at:
raise PlayStateError(
"plan_expired", "播放计划已过期,请重新获取",
reprepare_required=True,
)
if int(cfg["config_version"]) != play.config_version:
raise PlayStateError(
"config_changed", "视频配置已变化,请重新获取",
reprepare_required=True,
)
used = used_plays(db, user_id, play.scene)
if used >= maximum:
raise PlayStateError("play_limit_reached", "播放次数已用完")
play.seq = used + 1
play.status = "started"
play.started_at = now
try:
db.commit()
db.refresh(play)
except IntegrityError as exc:
db.rollback()
raise PlayStateError("play_limit_reached", "并发起播冲突,请重新获取") from exc
return _start_out(play, maximum, "started")
def _coin_balance(db: Session, user_id: int) -> int:
account = db.get(CoinAccount, user_id)
return int(account.coin_balance if account else 0)
def grant_circle(
db: Session, user_id: int, *, play_token: str, circle: int
def grant_play(
db: Session, user_id: int, *, play_token: str, completed: bool
) -> dict[str, Any]:
play = _find_play(db, user_id, play_token)
balance = _coin_balance(db, user_id)
base = {
"granted": False,
"retryable": False,
"retry_after_ms": 0,
"circle": circle,
"granted_coin": 0,
"settled_circles": int(play.settled_circles if play else 0),
"coin_balance": balance,
}
if play is None:
return {**base, "status": "not_found"}
if play.status not in {"started", "completed"} or play.started_at is None:
return {**base, "status": "not_started"}
settled = int(play.settled_circles)
if circle <= settled:
return {
**base,
"status": "already_granted",
"granted_coin": play.coin // CIRCLE_COUNT,
"settled_circles": settled,
}
if settled >= CIRCLE_COUNT:
return {**base, "status": "finished", "settled_circles": settled}
if circle != settled + 1:
return {
**base,
"status": "out_of_order",
"retryable": True,
"settled_circles": settled,
}
# 向上取整到毫秒,绝不让第 N 圈早于精确 N/10 边界发奖。
boundary_ms = (circle * play.duration_ms + CIRCLE_COUNT - 1) // CIRCLE_COUNT
elapsed_ms = max(0, int((_now() - play.started_at).total_seconds() * 1000))
if elapsed_ms < boundary_ms:
return {
**base,
"status": "too_early",
"retryable": True,
"retry_after_ms": boundary_ms - elapsed_ms,
"settled_circles": settled,
}
"""按 play_token 发这次引导视频的金币(幂等)。播完 / 中途关闭都发。
final = circle == CIRCLE_COUNT
返回 {granted, coin, status}:granted=True 表示**本次调用真的入账了**;
重复上报返回 granted=False + 已发金币(客户端据此不重复累加 toast 金额)
"""
token = (play_token or "").strip()
# 金币额度以**服务端配置**为准,不信客户端(客户端只上报"播完/关闭")。
play = _find_play(db, user_id, token)
coin = int(play.coin if play is not None else 0)
# 幂等核心:把 status 放进 WHERE 做条件更新(compare-and-set),而不是"先读再判再写"。
# 「播完」与「✕ 关闭」抢跑、或客户端超时重试时,两个请求会都读到 status='playing',
# 无锁的话就都往下发币、都 commit,金币入账两次(不用恶意,重试就会中招)。改成条件更新后
# 并发里只有一条 rowcount=1,另一条拿 0 → 按已发返回,不二次铸币。
# (PG READ COMMITTED 下后到的 UPDATE 阻塞到对手提交,再按新版本重判 status;SQLite 写串行。)
#
# 别指望 IntegrityError 兜底:这里只 UPDATE 不 INSERT,撞不到 uq_guide_video_play_token;
# 而 biz_type='guide_video' 的金币流水也不在 ux_coin_transaction_task_ref 的谓词
# (biz_type LIKE 'task%')覆盖范围内 —— 两个唯一键在这条路径上都是不生效的。
won = db.execute(
update(GuideVideoPlay)
.where(
GuideVideoPlay.id == play.id,
GuideVideoPlay.status == "started",
GuideVideoPlay.settled_circles == settled,
GuideVideoPlay.play_token == token,
GuideVideoPlay.user_id == user_id,
GuideVideoPlay.status == "playing",
)
.values(
settled_circles=circle,
status="completed" if final else "started",
completed=1 if final else 0,
granted_at=_now() if final else None,
status="granted",
coin=coin,
completed=1 if completed else 0,
granted_at=datetime.now(rewards.CN_TZ).replace(tzinfo=None),
)
.execution_options(synchronize_session=False)
).rowcount
if not won:
# 没抢到:token 不存在 / 不是本人的 / 已被另一次上报发过。回滚拿干净快照再区分两者
# (对手此时必然已提交 —— 我们就是被它挡下的,所以读得到它写的 coin)。
db.rollback()
fresh = _find_play(db, user_id, play_token)
fresh_settled = int(fresh.settled_circles if fresh else 0)
return {
**base,
"status": "already_granted" if circle <= fresh_settled else "out_of_order",
"retryable": circle > fresh_settled,
"granted_coin": (
play.coin // CIRCLE_COUNT if circle <= fresh_settled else 0
),
"settled_circles": fresh_settled,
"coin_balance": _coin_balance(db, user_id),
}
play = _find_play(db, user_id, token)
if play is None:
return {"granted": False, "coin": 0, "status": "not_found"}
return {"granted": False, "coin": play.coin, "status": "already_granted"}
per_circle = play.coin // CIRCLE_COUNT
account, _ = crud_wallet.grant_coins(
db,
user_id,
per_circle,
biz_type=BIZ_TYPE,
ref_id=f"{play.play_token}:{circle}",
remark=f"新手引导视频{circle}奖励",
)
try:
db.commit()
except IntegrityError:
# 流水唯一键是第二道幂等保险;若撞键,回滚后返回权威状态。
db.rollback()
fresh = _find_play(db, user_id, play_token)
return {
**base,
"status": "already_granted",
"granted_coin": play.coin // CIRCLE_COUNT,
"settled_circles": int(fresh.settled_circles if fresh else settled),
"coin_balance": _coin_balance(db, user_id),
}
return {
**base,
"granted": True,
"status": "granted",
"granted_coin": per_circle,
"settled_circles": circle,
"coin_balance": int(account.coin_balance),
}
def cleanup_old_videos(db: Session, *, protection: timedelta = PLAN_TTL) -> int:
"""清理超过保护期且未被当前配置/有效计划/近期起播引用的托管旧片。"""
now = _now()
protected_urls: set[str] = set()
for scene in SCENES:
url = get_config(db, scene).get("video_url")
if url:
protected_urls.add(str(url))
protected_urls.update(
str(url)
for url in db.execute(
select(GuideVideoPlay.video_url).where(
GuideVideoPlay.video_url.is_not(None),
(
(
(GuideVideoPlay.status == "prepared")
& (GuideVideoPlay.expires_at >= now)
)
| (
(GuideVideoPlay.status == "started")
& (GuideVideoPlay.started_at >= now - protection)
)
),
)
).scalars()
if url
)
directory = Path(settings.MEDIA_ROOT) / "guide_video"
if not directory.is_dir():
return 0
cutoff = now.timestamp() - protection.total_seconds()
removed = 0
for path in directory.glob("guide_*.mp4"):
url = f"{settings.MEDIA_URL_PREFIX}/guide_video/{path.name}"
try:
if url not in protected_urls and path.stat().st_mtime <= cutoff:
media.delete_guide_video(url)
removed += 1
except OSError:
continue
return removed
if coin > 0:
crud_wallet.grant_coins(
db,
user_id,
coin,
biz_type=BIZ_TYPE,
ref_id=token,
remark="新手引导视频奖励",
)
db.commit()
return {"granted": True, "coin": coin, "status": "granted"}
+7 -20
View File
@@ -132,20 +132,10 @@ def get_or_create_account(
) -> CoinAccount:
"""取用户金币账户,不存在则建一个空账户。
lock=True 时对已存在的账户行加 SELECT FOR UPDATE并强制刷新 identity map 中可能
缓存的旧余额所有金币写入口都会使用它SQLite FOR UPDATE no-op
lock=True 时对已存在的账户行加 SELECT FOR UPDATE(--写余额的调用方串行化,防并发
双写余额错位, admin set 模式连点);默认 False 不改 C 端发奖行为SQLite no-op
"""
if lock:
# populate_existing 很关键:同一 Session 可能早已缓存旧余额。只加 FOR UPDATE 而
# 继续复用 identity map 里的旧对象,仍会在另一笔发奖提交后把余额覆盖回去。
acc = db.execute(
select(CoinAccount)
.where(CoinAccount.user_id == user_id)
.with_for_update()
.execution_options(populate_existing=True)
).scalar_one_or_none()
else:
acc = db.get(CoinAccount, user_id)
acc = db.get(CoinAccount, user_id, with_for_update=True) if lock else db.get(CoinAccount, user_id)
if acc is None:
acc = CoinAccount(
user_id=user_id,
@@ -175,9 +165,7 @@ def grant_coins(
返回 (account, transaction)调用方负责 commit
"""
# 所有金币来源统一遵守账户行锁协议;否则 guide 锁了账户,广告/签到等未锁路径
# 仍可拿旧余额做 ORM 读改写,最终余额会小于流水累计。
acc = get_or_create_account(db, user_id, commit=False, lock=True)
acc = get_or_create_account(db, user_id, commit=False)
acc.coin_balance += amount
if amount > 0:
acc.total_coin_earned += amount
@@ -210,7 +198,7 @@ def grant_cash(
[grant_coins] 同模式(运营手动调现金 / 测试发现金用)返回 (account, transaction),
调用方负责 commit不在此校验扣成负由调用方(admin router)按业务保护
"""
acc = get_or_create_account(db, user_id, commit=False, lock=True)
acc = get_or_create_account(db, user_id, commit=False)
acc.cash_balance_cents += amount_cents
txn = CashTransaction(
@@ -240,7 +228,7 @@ def grant_invite_cash(
invite_cash_transaction, commit与金币兑换的 cash_balance_cents **物理隔离**
(产品红线:邀请奖励金 金币现金,两本账不可累加)返回 (account, transaction),
调用方负责 commit不在此校验扣成负由调用方按业务保护"""
acc = get_or_create_account(db, user_id, commit=False, lock=True)
acc = get_or_create_account(db, user_id, commit=False)
acc.invite_cash_balance_cents += amount_cents
txn = InviteCashTransaction(
@@ -298,8 +286,7 @@ def exchange_coins_to_cash(
if coin_amount < floor_min or coin_amount % COIN_PER_CENT != 0:
raise InvalidExchangeAmountError
# 余额充足校验必须和扣减遵守同一账户锁,否则并发消费都可能通过旧余额校验。
acc = get_or_create_account(db, user_id, commit=False, lock=True)
acc = get_or_create_account(db, user_id, commit=False)
if acc.coin_balance < coin_amount:
raise InsufficientCoinError
+3 -15
View File
@@ -7,19 +7,9 @@
"""
from __future__ import annotations
import re
from datetime import datetime
from pydantic import BaseModel, ConfigDict, Field, field_validator
from app.core.config import settings
def _validate_login_phone(phone: str) -> str:
"""普通号保持大陆手机号格式;仅显式配置的测试号允许例外。"""
if re.fullmatch(r"^1\d{10}$", phone) or phone in settings.test_account_phones:
return phone
raise ValueError("invalid phone")
from pydantic import BaseModel, ConfigDict, Field
# ===== 用户对外信息 =====
@@ -82,7 +72,7 @@ class JverifyLoginRequest(BaseModel):
# ===== 短信验证码 =====
class SmsSendRequest(BaseModel):
phone: str = Field(..., min_length=11, max_length=11)
phone: str = Field(..., min_length=11, max_length=11, pattern=r"^1\d{10}$")
device_id: str = Field(
"", max_length=64,
description="硬件级设备标识(Android ANDROID_ID),用于发码防刷按 设备+IP 限流;空=按 IP 聚一桶",
@@ -90,7 +80,6 @@ class SmsSendRequest(BaseModel):
device_model: str = Field(
"", max_length=128, description="客户端设备型号快照,用于短信安全审计"
)
_valid_phone = field_validator("phone")(_validate_login_phone)
class SmsSendResponse(BaseModel):
@@ -100,7 +89,7 @@ class SmsSendResponse(BaseModel):
class SmsLoginRequest(BaseModel):
phone: str = Field(..., min_length=11, max_length=11)
phone: str = Field(..., min_length=11, max_length=11, pattern=r"^1\d{10}$")
code: str = Field(..., min_length=4, max_length=8)
device_id: str = Field(
"", max_length=64,
@@ -109,7 +98,6 @@ class SmsLoginRequest(BaseModel):
device_model: str = Field(
"", max_length=128, description="客户端设备型号快照,用于短信验证安全审计"
)
_valid_phone = field_validator("phone")(_validate_login_phone)
# ===== Refresh =====
-11
View File
@@ -12,7 +12,6 @@ class DeviceRegisterRequest(BaseModel):
registration_id: str | None = None
push_vendor: str | None = None
push_token: str | None = None
push_binding_id: str | None = Field(default=None, max_length=128)
platform: str = "android"
app_version: str | None = None
@@ -24,7 +23,6 @@ class HeartbeatRequest(BaseModel):
registration_id: str | None = None
push_vendor: str | None = None
push_token: str | None = None
push_binding_id: str | None = Field(default=None, max_length=128)
class DeviceOut(BaseModel):
@@ -35,7 +33,6 @@ class DeviceOut(BaseModel):
registration_id: str | None
push_vendor: str | None
push_token: str | None
push_binding_id: str | None
ever_protected: bool
liveness_state: str
last_heartbeat_at: datetime | None
@@ -46,14 +43,6 @@ class OkResponse(BaseModel):
ok: bool = True
class DeviceUnregisterRequest(BaseModel):
device_id: str = Field(min_length=1, max_length=128)
class DeferredDeviceUnregisterRequest(DeviceUnregisterRequest):
push_binding_id: str = Field(min_length=32, max_length=128)
class LivenessOut(BaseModel):
"""本机掉线告警状态(后置检测 pull 版)。客户端只需这一个布尔判断要不要弹「开启自启动」引导,
故只返回 kill_alert_pending(不暴露设备详情 / 内部 liveness_state )从未注册过 默认 False(无告警)"""
+21 -55
View File
@@ -1,72 +1,38 @@
"""客户端引导视频三阶段协议"""
"""新手引导视频(领券等候浮层前 N 次替代广告)的客户端请求/响应契约"""
from __future__ import annotations
from typing import Literal
from pydantic import BaseModel, Field
GuideScene = Literal["coupon", "comparison"]
class GuideVideoPrepareIn(BaseModel):
scene: GuideScene = "coupon"
class GuideVideoPrepareOut(BaseModel):
should_play: bool
reason: str
scene: GuideScene
video_url: str | None = None
play_token: str = ""
config_version: int = 0
duration_ms: int = 0
circle_count: int = 10
circle_duration_ms: float = 0
reward_coin: int = 0
reward_per_circle: int = 0
seq: int = 0
remaining: int = 0
expires_at: str | None = None
class GuideVideoStartIn(BaseModel):
play_token: str = Field(min_length=1, max_length=64)
"""开播询问。scene 目前只有 coupon(领券浮层);预留给日后比价等场景。"""
scene: Literal["coupon", "comparison"] = "coupon"
class GuideVideoStartOut(BaseModel):
started: bool
status: Literal["started", "already_started"]
play_token: str
scene: GuideScene
video_url: str
config_version: int
duration_ms: int
circle_count: int = 10
circle_duration_ms: float
reward_coin: int
reward_per_circle: int
seq: int
remaining: int
started_at: str
"""should_play=False 时客户端照旧走广告链路,其余字段无意义。"""
should_play: bool
video_url: str | None = None # 相对地址 /media/...;客户端自行拼 BASE_URL
play_token: str = "" # 发奖幂等键
reward_coin: int = 0 # 播完/中途关闭都发的固定金币
seq: int = 0 # 本账号第几次
remaining: int = 0 # 发完这次还剩几次
class GuideVideoRewardIn(BaseModel):
"""播完或中途关闭都调这个;completed 只做留痕,两者都发币。"""
play_token: str = Field(min_length=1, max_length=64)
circle: int = Field(ge=1, le=10)
completed: bool = False
class GuideVideoRewardOut(BaseModel):
"""granted=True 表示本次调用真的入账(重复上报为 False,coin 是已发金额)。"""
granted: bool
status: Literal[
"granted",
"already_granted",
"too_early",
"out_of_order",
"not_started",
"not_found",
"finished",
]
retryable: bool = False
retry_after_ms: int = 0
circle: int
granted_coin: int = 0
settled_circles: int = 0
coin_balance: int = 0
coin: int
status: str
+19 -122
View File
@@ -83,13 +83,8 @@ def _dispatch(
extra: dict[str, str] | None = None,
dedup_key: str | None = None,
push_vars: dict[str, str] | None = None,
push_token_contains: str | None = None,
) -> Notification | None:
"""落一条站内消息并向该用户设备直推。返回落库行;去重命中/失败返回 None。"""
logger.info(
"notification dispatch started user_id=%s type=%s dedup_key=%s",
user_id, type_key, dedup_key,
)
try:
row = notif_repo.create_notification(
db,
@@ -116,20 +111,11 @@ def _dispatch(
logger.exception("rollback after notification failure also failed")
return None
logger.info(
"notification created user_id=%s type=%s notification_id=%s dedup_key=%s",
user_id, type_key, row.id, dedup_key,
)
_push_to_user_devices(db, row, push_vars, push_token_contains)
_push_to_user_devices(db, row, push_vars)
return row
def _push_to_user_devices(
db: Session,
row: Notification,
push_vars: dict[str, str] | None,
push_token_contains: str | None = None,
) -> None:
def _push_to_user_devices(db: Session, row: Notification, push_vars: dict[str, str] | None) -> None:
"""向消息归属用户的全部厂商推送目标直推(best-effort,单设备失败不影响其余)。"""
try:
title, body = catalog.render_push(row.type, push_vars)
@@ -138,52 +124,17 @@ def _push_to_user_devices(
extras.update({str(k): str(v) for k, v in (row.extra or {}).items()})
extras["notificationId"] = str(row.id)
targets = device_repo.list_push_targets(db, user_id=row.user_id)
if push_token_contains:
normalized_filter = push_token_contains.casefold()
targets = [
target
for target in targets
if normalized_filter in (target.push_token or "").casefold()
]
logger.info(
"push targets resolved user_id=%s type=%s notification_id=%s "
"token_filter=%s target_count=%s",
row.user_id, row.type, row.id, bool(push_token_contains), len(targets),
)
if not targets:
logger.warning(
"push skipped no targets user_id=%s type=%s notification_id=%s",
row.user_id, row.type, row.id,
)
return
sent = failed = skipped = data_sent = data_failed = 0
for dev in targets:
for dev in device_repo.list_push_targets(db, user_id=row.user_id):
vendor = vendor_push.normalize_vendor(dev.push_vendor)
if not vendor or vendor not in vendor_push.SUPPORTED_VENDORS:
skipped += 1
logger.warning(
"push target skipped unsupported vendor user_id=%s type=%s "
"notification_id=%s device_id=%s raw_vendor=%s normalized_vendor=%s",
row.user_id, row.type, row.id, dev.device_id, dev.push_vendor, vendor,
)
continue
missing = vendor_push.missing_settings(vendor)
if missing:
skipped += 1
if vendor_push.missing_settings(vendor):
# 本地/测试环境凭据不齐 → 只落站内消息,不发真推送(与 push/vendors 的报缺口径一致)
logger.warning(
"push target skipped vendor not configured user_id=%s type=%s "
"notification_id=%s device_id=%s vendor=%s missing_settings=%s",
row.user_id, row.type, row.id, dev.device_id, vendor, missing,
logger.info(
"skip push (vendor %s not configured) user_id=%s type=%s",
vendor, row.user_id, row.type,
)
continue
logger.info(
"push send started user_id=%s type=%s notification_id=%s "
"device_id=%s vendor=%s",
row.user_id, row.type, row.id, dev.device_id, vendor,
)
try:
response = vendor_push.send_notification(
vendor, dev.push_token, title=title, body=body, extras=extras
@@ -197,44 +148,26 @@ def _push_to_user_devices(
row.user_id, row.type, row.id, response,
)
logger.info(
"push sent user_id=%s type=%s vendor=%s notification_id=%s device_id=%s",
row.user_id, row.type, vendor, row.id, dev.device_id,
"push sent user_id=%s type=%s vendor=%s notification_id=%s",
row.user_id, row.type, vendor, row.id,
)
sent += 1
except vendor_push.VendorPushError as e:
failed += 1
logger.warning(
"push failed user_id=%s type=%s vendor=%s notification_id=%s "
"device_id=%s error=%s",
row.user_id, row.type, vendor, row.id, dev.device_id, e,
"push failed user_id=%s type=%s vendor=%s: %s", row.user_id, row.type, vendor, e
)
try:
data_response = vendor_push.send_data_event(
vendor_push.send_data_event(
vendor,
dev.push_token,
event=vendor_push.DATA_EVENT_NOTIFICATION_CREATED,
notification_id=str(row.id),
)
data_sent += 1
logger.info(
"push data event completed user_id=%s type=%s vendor=%s "
"notification_id=%s device_id=%s response=%s",
row.user_id, row.type, vendor, row.id, dev.device_id, data_response,
)
except vendor_push.VendorPushError as e:
data_failed += 1
# 透传只负责前台铃铛实时刷新,失败不能影响通知栏消息或站内消息。
logger.warning(
"push data event failed user_id=%s type=%s vendor=%s notification_id=%s "
"device_id=%s error=%s",
row.user_id, row.type, vendor, row.id, dev.device_id, e,
"push data event failed user_id=%s type=%s vendor=%s: %s",
row.user_id, row.type, vendor, e,
)
logger.info(
"push dispatch completed user_id=%s type=%s notification_id=%s targets=%s "
"sent=%s failed=%s skipped=%s data_sent=%s data_failed=%s",
row.user_id, row.type, row.id, len(targets),
sent, failed, skipped, data_sent, data_failed,
)
except Exception: # noqa: BLE001 — 渲染/查设备等意外失败同样不外抛
logger.exception("push notification failed user_id=%s type=%s", row.user_id, row.type)
@@ -244,12 +177,7 @@ def _push_to_user_devices(
# ---------------------------------------------------------------------------
def notify_withdraw_success(
db: Session,
order: WithdrawOrder,
*,
push_token_contains: str | None = None,
) -> None:
def notify_withdraw_success(db: Session, order: WithdrawOrder) -> None:
"""#3 提现成功:款项已存入微信零钱。点击无跳转仅消红点(extra 空)。"""
_dispatch(
db,
@@ -263,16 +191,10 @@ def notify_withdraw_success(
extra={},
dedup_key=order.out_bill_no,
push_vars={"amount": notif_repo.cash_yuan(order.amount_cents)},
push_token_contains=push_token_contains,
)
def notify_withdraw_failed(
db: Session,
order: WithdrawOrder,
*,
push_token_contains: str | None = None,
) -> None:
def notify_withdraw_failed(db: Session, order: WithdrawOrder) -> None:
"""#4 提现失败/退回:含微信侧失败、审核拒绝、解绑退回。点击跳提现页重新提现。
失败原因用 order.fail_reason( /withdraw/status 下发的用户可读原因同源)
@@ -290,16 +212,10 @@ def notify_withdraw_failed(
extra={"withdrawId": order.out_bill_no},
dedup_key=order.out_bill_no,
push_vars={"amount": notif_repo.cash_yuan(order.amount_cents), "reason": reason},
push_token_contains=push_token_contains,
)
def notify_feedback_reply(
db: Session,
feedback: Feedback,
*,
push_token_contains: str | None = None,
) -> None:
def notify_feedback_reply(db: Session, feedback: Feedback) -> None:
"""#9 官方回复:运营审核了反馈且未采纳(用户可见原因/留言落在反馈记录上)。
点击跳反馈历史页滚动高亮该条(extra.feedbackId)"""
_dispatch(
@@ -309,16 +225,10 @@ def notify_feedback_reply(
info_rows=[{"label": "说明文案", "value": "快去看看官方给您的回复吧~"}],
extra={"feedbackId": str(feedback.id)},
dedup_key=str(feedback.id),
push_token_contains=push_token_contains,
)
def notify_feedback_reward(
db: Session,
feedback: Feedback,
*,
push_token_contains: str | None = None,
) -> None:
def notify_feedback_reward(db: Session, feedback: Feedback) -> None:
"""#10 反馈奖励:反馈被采纳,金币已到账。PRD 约定发奖必带官方留言(admin_reply);
运营漏填时省略该信息行,不硬造文案"""
coins = int(feedback.reward_coins or 0)
@@ -336,16 +246,10 @@ def notify_feedback_reward(
extra={"feedbackId": str(feedback.id)},
dedup_key=str(feedback.id),
push_vars={"coins": str(coins)},
push_token_contains=push_token_contains,
)
def notify_report_approved(
db: Session,
report: PriceReport,
*,
push_token_contains: str | None = None,
) -> None:
def notify_report_approved(db: Session, report: PriceReport) -> None:
"""#11 爆料审核通过:上报的更低价过审,金币已到账。点击跳爆料记录页高亮该条。"""
coins = int(report.reward_coins or 0)
store = (report.store_name or "").strip() or "该店铺"
@@ -361,17 +265,11 @@ def notify_report_approved(
extra={"reportId": str(report.id)},
dedup_key=str(report.id),
push_vars={"store": store, "coins": str(coins)},
push_token_contains=push_token_contains,
)
def notify_invite_order_reward(
db: Session,
*,
inviter_user_id: int,
invitee_user_id: int,
cash_cents: int,
push_token_contains: str | None = None,
db: Session, *, inviter_user_id: int, invitee_user_id: int, cash_cents: int
) -> None:
"""#12 好友下单到账:被邀请好友完成首次下单(比价),现金奖励已入邀请人账户。
通知发给邀请人;每个好友只发一次奖 dedup 按被邀请人"""
@@ -389,5 +287,4 @@ def notify_invite_order_reward(
extra={"inviteeNickname": nickname},
dedup_key=str(invitee_user_id),
push_vars={"nickname": nickname, "amount": _yuan_trim(cash_cents)},
push_token_contains=push_token_contains,
)
-1
View File
@@ -87,7 +87,6 @@ PRD §5 的 13 条 push 文案(标题固定 ≤11 字不带变量;正文 `{var}`
**真发注意**:
- 目标手机必须先装 App 且客户端已集成对应厂商 SDK、`/device/register` 上报过 token;
- 荣耀正式环境必须使用正式消息(`HONOR_PUSH_TARGET_USER_TYPE=0`);开发联调测试设备时才设为 `1`;
- vivo 未上架前走测试推送(`VIVO_PUSH_MODE=1`),目标手机需在 vivo 开放平台加入测试设备;
- 小米新设备需在开放平台把签名/包名配好,token 才有效。
+10 -66
View File
@@ -16,10 +16,6 @@ xiaomi/oppo/vivo)到该用户 device_liveness 里已注册的 push token。
.venv\\Scripts\\python.exe scripts\\fire_push_events.py --count 1 # 各 1 条(先小量验证通道)
.venv\\Scripts\\python.exe scripts\\fire_push_events.py --types withdraw_failed --count 3
.venv\\Scripts\\python.exe scripts\\fire_push_events.py --types all --count 2
.venv\\Scripts\\python.exe scripts\\fire_push_events.py --phone 22222222222 --token ABC123 --count 1
`--token` 是可选不区分大小写的 token 子串过滤器;会推送给该账号下所有 token 包含该
子串的设备没有匹配设备时脚本直接停止,不会生成站内消息
推送成败看输出里的 `shagua.vendor_push` 日志(push sent / push failed);2 台设备则每条各推 2
凭据缺失或 token 失效时 notification_events 只记日志不抛错(站内消息仍会落库)
@@ -62,52 +58,33 @@ _FAIL_REASONS = [
]
def _fire_one(
db,
uid: int,
type_key: str,
i: int,
*,
push_token_contains: str | None = None,
) -> None:
def _fire_one(db, uid: int, type_key: str, i: int) -> None:
"""构造一条该类型的瞬态业务对象(不落业务表,只为给 notify 函数读字段),触发真实推送。"""
if type_key == "withdraw_success":
order = WithdrawOrder(user_id=uid, out_bill_no=uuid.uuid4().hex, amount_cents=50, source="coin_cash")
notification_events.notify_withdraw_success(
db, order, push_token_contains=push_token_contains
)
notification_events.notify_withdraw_success(db, order)
elif type_key == "withdraw_failed":
order = WithdrawOrder(
user_id=uid, out_bill_no=uuid.uuid4().hex, amount_cents=350, source="coin_cash",
fail_reason=random.choice(_FAIL_REASONS),
)
notification_events.notify_withdraw_failed(
db, order, push_token_contains=push_token_contains
)
notification_events.notify_withdraw_failed(db, order)
elif type_key == "invite_order_reward":
# 假被邀请人 id(> 真实用户范围,避重):昵称回退「好友」。真实昵称请走 API 流程(见文末说明)。
fake_invitee = random.randint(900000, 999999)
notification_events.notify_invite_order_reward(
db,
inviter_user_id=uid,
invitee_user_id=fake_invitee,
cash_cents=INVITE_COMPARE_REWARD_CENTS,
push_token_contains=push_token_contains,
db, inviter_user_id=uid, invitee_user_id=fake_invitee, cash_cents=INVITE_COMPARE_REWARD_CENTS
)
elif type_key == "feedback_reward":
fb = Feedback(user_id=uid, content="(直发)", contact="", status="adopted",
reward_coins=300, admin_reply="感谢反馈,您说的问题已修复上线,金币请查收~")
fb.id = random.randint(900000, 999999)
notification_events.notify_feedback_reward(
db, fb, push_token_contains=push_token_contains
)
notification_events.notify_feedback_reward(db, fb)
elif type_key == "feedback_reply":
fb = Feedback(user_id=uid, content="(直发)", contact="", status="rejected",
admin_reply="您的建议我们记录啦,会在后续版本评估~")
fb.id = random.randint(900000, 999999)
notification_events.notify_feedback_reply(
db, fb, push_token_contains=push_token_contains
)
notification_events.notify_feedback_reply(db, fb)
elif type_key == "report_approved":
rep = PriceReport(
user_id=uid, reported_platform_id="jd", reported_platform_name="京东外卖",
@@ -115,9 +92,7 @@ def _fire_one(
reward_coins=PRICE_REPORT_REWARD_COINS, store_name=f"测试火锅店{i:02d}",
)
rep.id = random.randint(900000, 999999)
notification_events.notify_report_approved(
db, rep, push_token_contains=push_token_contains
)
notification_events.notify_report_approved(db, rep)
else:
raise SystemExit(f"未知类型: {type_key}(可选: {', '.join(ALL_TYPES)})")
@@ -125,11 +100,6 @@ def _fire_one(
def main() -> None:
parser = argparse.ArgumentParser(description="直接触发消息通知中心真实推送(后台驱动不了的事件用)")
parser.add_argument("--phone", default=DEFAULT_PHONE, help=f"目标用户手机号(默认 {DEFAULT_PHONE})")
parser.add_argument(
"--token",
default="",
help="可选的厂商 token 子串(不区分大小写);仅推送到该账号下所有匹配设备",
)
parser.add_argument("--count", type=int, default=10, help="每类发多少条(默认 10)")
parser.add_argument(
"--types", default=",".join(DEFAULT_TYPES),
@@ -152,40 +122,14 @@ def main() -> None:
uid = user.id
targets = device_repo.list_push_targets(db, user_id=uid)
token_filter = args.token.strip()
normalized_filter = token_filter.casefold()
matched_targets = [
target
for target in targets
if not token_filter
or normalized_filter in (target.push_token or "").casefold()
]
if token_filter and not matched_targets:
print(
f"❌ 用户 {args.phone}{len(targets)} 台有效推送设备中,"
f"没有 token 包含 {token_filter!r};未生成站内消息,也未发送推送。"
)
return
print(
f"目标用户 {args.phone}(id={uid});有效推送设备 {len(targets)} 台;"
f"本次匹配 {len(matched_targets)} 台:"
f"{[target.push_vendor for target in matched_targets] or '无(手机收不到!先在 App 上报 push token)'}"
)
if token_filter:
print(f"token 子串过滤(不区分大小写): {token_filter!r}")
print(f"目标用户 {args.phone}(id={uid});已注册推送设备 {len(targets)} 台:"
f"{[t.push_vendor for t in targets] or '无(手机收不到!先在 App 上报 push token)'}")
print(f"即将触发:{types},每类 {args.count} 条 → 共 {len(types) * args.count}\n")
for t in types:
print(f"── {t} ×{args.count} " + "" * 30)
for i in range(1, args.count + 1):
_fire_one(
db,
uid,
t,
i,
push_token_contains=token_filter or None,
)
_fire_one(db, uid, t, i)
print(f"\n✅ 已触发完。站内消息已落 notification 表(用 {args.phone} 登录 App 可在消息中心看到);"
"\n 手机推送成败见上方 `shagua.vendor_push` 日志(push sent=成功 / push failed=失败)。")
+3 -36
View File
@@ -15,10 +15,6 @@ services/notification_events.notify_invite_order_reward —— 与生产同一
.venv\\Scripts\\python.exe scripts\\test_push_invite_order_reward.py # 随机金额发 1 条
.venv\\Scripts\\python.exe scripts\\test_push_invite_order_reward.py --cents 200 # 固定 2.00 元
.venv\\Scripts\\python.exe scripts\\test_push_invite_order_reward.py --invitee-phone 12000000001 # 真实昵称
.venv\\Scripts\\python.exe scripts\\test_push_invite_order_reward.py --phone 22222222222 --token gkha
`--token` 是可选不区分大小写的 token 子串过滤器会推送给该账号下所有 token 包含该
子串的设备没有匹配设备时脚本直接停止不会生成站内消息
结果判读(看输出日志):
push sent = 厂商接口受理成功,手机应弹好友下单奖励到账通知
@@ -66,11 +62,6 @@ def _notif_count(db, uid: int) -> int:
def main() -> None:
parser = argparse.ArgumentParser(description="#12 好友下单到账 推送联调(每次 1 条,金额默认随机)")
parser.add_argument("--phone", default="11111111111", help="邀请人(收通知方)手机号,默认 11111111111")
parser.add_argument(
"--token",
default="",
help="可选的厂商 token 子串(不区分大小写);仅推送到该账号下所有匹配设备",
)
parser.add_argument("--cents", type=int, default=None,
help="奖励金额,单位分(默认随机 1~9999;线上真实值 200)")
parser.add_argument("--invitee-phone", default="",
@@ -96,37 +87,13 @@ def main() -> None:
invitee_id = random.randint(900000, 999999) # 假 id,昵称兜底「好友」,永不去重
targets = device_repo.list_push_targets(db, user_id=user.id)
token_filter = args.token.strip()
normalized_filter = token_filter.casefold()
matched_targets = [
target
for target in targets
if not token_filter
or normalized_filter in (target.push_token or "").casefold()
]
if token_filter and not matched_targets:
print(
f"❌ 用户 {args.phone}{len(targets)} 台有效推送设备中,"
f"没有 token 包含 {token_filter!r};未生成站内消息,也未发送推送。"
)
return
print(
f"邀请人 {args.phone}(id={user.id});有效推送设备 {len(targets)} 台;"
f"本次匹配 {len(matched_targets)} 台:"
f"{[target.push_vendor for target in matched_targets] or '无 ← 手机收不到!先在 App 上报 push token'}"
)
if token_filter:
print(f"token 子串过滤(不区分大小写): {token_filter!r}")
vendors = [t.push_vendor for t in targets]
print(f"邀请人 {args.phone}(id={user.id});推送设备 {len(targets)} 台:{vendors or '无 ← 手机收不到!先在 App 上报 push token'}")
before = _notif_count(db, user.id)
print(f"→ 本次奖励 【{cents / 100:.2f} 元】(invitee_user_id={invitee_id}),手机上按金额认领这条通知")
notification_events.notify_invite_order_reward(
db,
inviter_user_id=user.id,
invitee_user_id=invitee_id,
cash_cents=cents,
push_token_contains=token_filter or None,
db, inviter_user_id=user.id, invitee_user_id=invitee_id, cash_cents=cents
)
created = _notif_count(db, user.id) - before
-69
View File
@@ -140,75 +140,6 @@ def test_harvest_done_failed_derives_fail_reason(client) -> None:
assert rec.information == "比价过程出错,请稍后重试" # 原文案仍留存
def _done_params_platforms_no_isbest() -> dict:
"""id 3304 型:done 帧 platforms 全平台「相似替换/仅供参考」(has_dish_diff),pricebot
一个 is_best 都没标,但有有价目标(美团 57.8 < 源淘宝闪购 60.8)"""
return {
"platforms": [
{"role": "source", "platform_id": "eleme", "platform_name": "淘宝闪购",
"package": "me.ele", "price": 60.8, "is_best": False, "has_dish_diff": False,
"store_name": "窑鸡王", "items": [{"name": "招牌窑鸡 整只-香辣", "qty": 1}]},
{"role": "target", "platform_id": "meituan_waimai", "platform_name": "美团外卖",
"package": "com.sankuai.meituan.takeoutnew", "price": 57.8, "is_best": False,
"has_dish_diff": True},
{"role": "target", "platform_id": "jd_waimai_standalone", "platform_name": "京东外卖",
"package": "com.jd.waimai", "price": 74.9, "is_best": False, "has_dish_diff": True},
],
"information": "美团更便宜(含相似商品替换)",
}
def test_harvest_done_platforms_no_isbest_falls_back_to_cheapest_target(client) -> None:
"""回归(id 3304):platforms 全无 is_best(全平台相似替换)但有有价目标 → best 应兜底取
有价目标里最低那家,不能让 best_*/saved 整条落 NULL(否则首页价 0.00 / 记录页无最低红框 /
省额丢失)"""
tid = _tid()
with SessionLocal() as db:
crud.harvest_running(db, trace_id=tid, user_id=None)
rec, _ = crud.harvest_done(db, trace_id=tid, user_id=None,
done_params=_done_params_platforms_no_isbest())
assert rec.status == "success"
assert rec.best_platform_id == "meituan_waimai" # 有价目标里最低
assert rec.best_price_cents == 5780 # 57.8 元
assert rec.source_price_cents == 6080 # 源淘宝闪购 60.8
assert rec.saved_amount_cents == 300 # 60.8 - 57.8
assert rec.is_source_best is False # 兜底选的是目标,非源
def _done_params_no_isbest_source_cheapest() -> dict:
"""无 is_best 且源本身最便宜:源淘宝闪购 50.0 < 全部 dish-diff 目标(美团 57.8 / 京东 74.9)。
此时不能强选更贵的目标当 best(否则 saved 变负污染累计发现可省),应回落源"""
return {
"platforms": [
{"role": "source", "platform_id": "eleme", "platform_name": "淘宝闪购",
"package": "me.ele", "price": 50.0, "is_best": False, "has_dish_diff": False,
"store_name": "窑鸡王", "items": [{"name": "招牌窑鸡 整只-香辣", "qty": 1}]},
{"role": "target", "platform_id": "meituan_waimai", "platform_name": "美团外卖",
"package": "com.sankuai.meituan.takeoutnew", "price": 57.8, "is_best": False,
"has_dish_diff": True},
{"role": "target", "platform_id": "jd_waimai_standalone", "platform_name": "京东外卖",
"package": "com.jd.waimai", "price": 74.9, "is_best": False, "has_dish_diff": True},
],
"information": "源平台已是最低(其余为相似替换)",
}
def test_harvest_done_no_isbest_source_cheapest_falls_back_to_source(client) -> None:
"""回归:无 is_best 且源最便宜 → best 回落源(与 _derive「全目标缺菜回落源、不虚报省」同语义),
saved=0is_source_best=True,绝不因强选更贵目标而让 saved 变负倒扣累计发现可省"""
tid = _tid()
with SessionLocal() as db:
crud.harvest_running(db, trace_id=tid, user_id=None)
rec, _ = crud.harvest_done(db, trace_id=tid, user_id=None,
done_params=_done_params_no_isbest_source_cheapest())
assert rec.status == "success"
assert rec.best_platform_id == "eleme" # 回落到源(源最便宜)
assert rec.best_price_cents == 5000 # 源 50.0
assert rec.source_price_cents == 5000
assert rec.saved_amount_cents == 0 # 没省到,绝不为负
assert rec.is_source_best is True # 源就是最便宜
def test_harvest_abort_cancels_running(client) -> None:
tid = _tid()
with SessionLocal() as db: