Compare commits

...

6 Commits

Author SHA1 Message Date
exinglang a0196b8d64 fix(push): 完善厂商推送排障日志 2026-07-29 19:52:06 +08:00
zuochenyong 53c3b7f60f feat(guide-video): 支持领券和比价独立视频奖励配置 (#196)
Co-authored-by: exinglang <exinglang@qq.com>
Reviewed-on: #196
Co-authored-by: zuochenyong <zuochenyong@wonderable.ai>
Co-committed-by: zuochenyong <zuochenyong@wonderable.ai>
2026-07-29 16:12:05 +08:00
marco 4bd4e66678 重构比价结果页取数据逻辑 (#195)
Reviewed-on: #195
2026-07-29 01:59:31 +08:00
linkeyu 90c6fe599a 修复:统一用户Draw信息流eCPM统计口径 (#190)
## 问题

业务收益详情的平均 Draw eCPM 仅平均成功发奖记录,会排除未发奖的真实展示,导致数值系统性偏高,且与广告收益页口径不一致。

## 修复

- `feed_avg_ecpm` 改为从 `ad_ecpm_record` 的全部 `draw/feed` 实际展示计算
- 成功发奖、未发奖展示均纳入,每次展示等权
- 日期、正式/测试环境、业务代码位、领券/比价场景支持与广告收益页对齐
- 奖励份数仍基于成功发奖表,不混用展示数据源
- 复用广告收益报表的业务代码位集合

## 线上数据复算

2026-07-25、正式业务、用户 #33:

- 旧口径(只看成功发奖):`29.9117 元/千次`
- 新口径(333 次真实展示):`19.9926 元/千次`
- 新值与广告收益报表一致

## 验证

- 新增成功/未发奖、场景、环境、业务代码位回归用例
- `tests/test_admin_read.py` + `tests/test_admin_ad_revenue_scope.py`:29 项全通过
- Ruff 改动文件检查通过

## 上线顺序

本 PR 需先于管理后台配套 PR 上线。

---------

Co-authored-by: guke <guke@wonderable.ai>
Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #190
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-07-28 17:58:23 +08:00
linkeyu e529112a90 修复中途退出比价的 LLM 成本回填 (#191)
## 问题

比价记录进入中途退出后未触发 LLM 成本回填,周期补偿也未扫描 cancelled,导致实际已有 LLM 调用的记录长期显示成本、LLM、TOKEN 为空。

## 修改

- finalize 落库后立即追加 LLM 成本回填
- 周期补偿范围加入 cancelled
- 保持无有效调用和全调用失败记录不伪造成本
- 增加即时回填和周期补偿回归测试

## 验证

- ruff 检查通过
- 相关测试 28 项通过
- 全仓 626 项通过;主干既有失败已在未修改的 origin/main 复现

---------

Co-authored-by: unknown <798648091@qq.com>
Reviewed-on: #191
Co-authored-by: linkeyu <linkeyu@wonderable.ai>
Co-committed-by: linkeyu <linkeyu@wonderable.ai>
2026-07-28 17:57:52 +08:00
guke 50da718e35 比价记录失败卡展示具体原因(新增 fail_reason) (#189)
失败记录不再一律「网络开小差」:新增记录级 fail_reason 派生列——information
具体则直出,笼统则从 platform_results 救出业务原因(找不到店/菜、未起送、打烊、
单点不配送等),纯系统失败为 None → 端侧品牌兜底。store_closed/no_delivery 被
pricebot 漏成 status=failed 的按 reason 补判,打烊脏店名统一简短模板。接入
harvest_done 与灰度期 upsert_record 两条写路径。

- models: comparison_record.fail_reason 列
- repositories: _derive_fail_display + 补判/清洗 helper,两条写路径接入
- schemas: ComparisonRecordOut 暴露 fail_reason
- alembic: 加列 + 回填老 specific 失败记录
- tests: _derive_fail_display 单测(8 例)+ harvest 失败落库集成测试

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: guke <guke@autohome.com.cn>
Reviewed-on: #189
2026-07-28 14:04:36 +08:00
25 changed files with 1132 additions and 111 deletions
@@ -0,0 +1,26 @@
"""merge comparison platforms + fail_reason heads
Revision ID: 6d2309208549
Revises: comparison_platforms_col, comparison_record_fail_reason
Create Date: 2026-07-29 01:48:41.868083
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
# revision identifiers, used by Alembic.
revision: str = '6d2309208549'
down_revision: Union[str, Sequence[str], None] = ('comparison_platforms_col', 'comparison_record_fail_reason')
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
pass
def downgrade() -> None:
pass
@@ -0,0 +1,47 @@
"""add platforms unified array column to comparison_record
展示模型统一数组(pricebot done.params.platforms 原样存): 每平台一行、自带
status/is_best/display, 记录页据此直接渲染, 不再靠 comparison_results + 客户端合并 + 前端派生。
纯新增列, 老记录为空 → 前端回退老 comparison_results。
Revision ID: comparison_platforms_col
Revises: user_manual_risk_fields
"""
from __future__ import annotations
from collections.abc import Sequence
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
from alembic import op
revision: str = "comparison_platforms_col"
down_revision: str | None = "user_manual_risk_fields"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
def upgrade() -> None:
# 幂等: 线上为了提前给历史数据补 platforms(2026-07-29), 已手动
# `ALTER TABLE comparison_record ADD COLUMN IF NOT EXISTS platforms jsonb
# NOT NULL DEFAULT '[]'::jsonb`(与本 migration 定义一致)。列已存在时跳过,
# 否则上线 alembic upgrade head 会撞 DuplicateColumn 直接部署失败。
bind = op.get_bind()
cols = {c["name"] for c in sa.inspect(bind).get_columns("comparison_record")}
if "platforms" in cols:
return
with op.batch_alter_table("comparison_record") as batch_op:
batch_op.add_column(
sa.Column(
"platforms", _JSON, nullable=False,
server_default=sa.text("'[]'"),
)
)
def downgrade() -> None:
with op.batch_alter_table("comparison_record") as batch_op:
batch_op.drop_column("platforms")
@@ -0,0 +1,49 @@
"""comparison_record.fail_reason (失败卡展示原因)
Revision ID: comparison_record_fail_reason
Revises: user_manual_risk_fields
Create Date: 2026-07-28 12:00:00.000000
失败记录的展示原因:information 具体则=它;笼统则由写路径从 platform_results 捞出的
业务原因;纯系统失败为 None(端侧品牌兜底)。见 repositories.comparison._derive_fail_display。
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
# revision identifiers, used by Alembic.
revision: str = 'comparison_record_fail_reason'
down_revision: Union[str, Sequence[str], None] = 'user_manual_risk_fields'
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
with op.batch_alter_table('comparison_record', schema=None) as batch_op:
batch_op.add_column(sa.Column('fail_reason', sa.String(length=256), nullable=True))
# 回填老失败记录:information 具体的直接搬过来(笼统/系统失败留 None → 端侧品牌兜底)。
# 新记录由写路径 _derive_fail_display 落库(含 platform_results 救援/补判),不走这条。
# platform_results 只在 raw_payload 里,SQL 里不易解析,故老记录不做救援/补判(可接受:
# 老 mixed/打烊记录回退品牌兜底);具体 information 的老记录本次即可显示真实原因。
op.execute(
"""
UPDATE comparison_record
SET fail_reason = information
WHERE status = 'failed'
AND information IS NOT NULL
AND information <> ''
AND information NOT IN (
'比价过程出错,请稍后重试',
'比价出错',
'比价未完成',
'done 参数缺少可验证的目标平台结果'
)
"""
)
def downgrade() -> None:
with op.batch_alter_table('comparison_record', schema=None) as batch_op:
batch_op.drop_column('fail_reason')
@@ -0,0 +1,31 @@
"""guide video play count is independent for coupon and comparison
Revision ID: guide_video_scene_unique
Revises: 6d2309208549
"""
from alembic import op
revision = "guide_video_scene_unique"
down_revision = "6d2309208549"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.drop_index("uq_guide_video_play_user_seq", table_name="guide_video_play")
op.create_index(
"uq_guide_video_play_user_scene_seq",
"guide_video_play",
["user_id", "scene", "seq"],
unique=True,
)
def downgrade() -> None:
op.drop_index("uq_guide_video_play_user_scene_seq", table_name="guide_video_play")
op.create_index(
"uq_guide_video_play_user_seq",
"guide_video_play",
["user_id", "seq"],
unique=True,
)
+5 -5
View File
@@ -43,7 +43,7 @@ _KNOWN_PROD_BUSINESS_CODE_IDS = frozenset({"104098712", "104099389"})
_TEST_BUSINESS_CODE_IDS = frozenset({"104127529", "104127626", "104137445"}) _TEST_BUSINESS_CODE_IDS = frozenset({"104127529", "104127626", "104137445"})
def _business_code_ids(db: Session, app_env: str | None) -> set[str]: def business_code_ids(db: Session, app_env: str | None) -> set[str]:
"""返回指定应用环境下可用于业务收益对账的 GroMore 聚合代码位。""" """返回指定应用环境下可用于业务收益对账的 GroMore 聚合代码位。"""
prod_config = app_config.get_ad_config(db) prod_config = app_config.get_ad_config(db)
prod_ids = set(_KNOWN_PROD_BUSINESS_CODE_IDS) | { prod_ids = set(_KNOWN_PROD_BUSINESS_CODE_IDS) | {
@@ -320,10 +320,10 @@ def ad_revenue_report(
# 业务口径仅保留正式配置/测试业务链路实际使用的代码位。穿山甲“全量”还包含广告测试 # 业务口径仅保留正式配置/测试业务链路实际使用的代码位。穿山甲“全量”还包含广告测试
# demo、插屏等没有客户端收益上报的曝光,两边直接比较会天然产生假差额。 # demo、插屏等没有客户端收益上报的曝光,两边直接比较会天然产生假差额。
business_code_ids: set[str] | None = None business_ids: set[str] | None = None
if revenue_scope == "business": if revenue_scope == "business":
business_code_ids = _business_code_ids(db, app_env) business_ids = business_code_ids(db, app_env)
events = [e for e in events if e.get("our_code_id") in business_code_ids] events = [e for e in events if e.get("our_code_id") in business_ids]
# 排序:time=按时间倒序(新→旧);ecpm=按 eCPM 数值倒序(eCPM 原值是字符串「分」,转数值排; # 排序:time=按时间倒序(新→旧);ecpm=按 eCPM 数值倒序(eCPM 原值是字符串「分」,转数值排;
# 纯发奖行用其发奖采用的 eCPM,缺失/非法计 0 排末尾)。 # 纯发奖行用其发奖采用的 eCPM,缺失/非法计 0 排末尾)。
@@ -381,7 +381,7 @@ def ad_revenue_report(
date_from=date_from, date_from=date_from,
date_to=date_to, date_to=date_to,
app_env=app_env, app_env=app_env,
our_code_ids=business_code_ids, our_code_ids=business_ids,
) )
if pangle_aggs: if pangle_aggs:
by_date = {a["date"]: a for a in pangle_aggs} by_date = {a["date"]: a for a in pangle_aggs}
+57 -14
View File
@@ -14,6 +14,7 @@ from sqlalchemy.orm import Session
from app.core import rewards from app.core import rewards
from app.core.config import settings from app.core.config import settings
from app.models.ad_ecpm import AdEcpmRecord
from app.models.ad_feed_reward import AdFeedRewardRecord from app.models.ad_feed_reward import AdFeedRewardRecord
from app.models.ad_reward import AdRewardRecord from app.models.ad_reward import AdRewardRecord
from app.models.admin import AdminAuditLog from app.models.admin import AdminAuditLog
@@ -1257,11 +1258,15 @@ def user_reward_stats(
date_from: datetime | None = None, date_from: datetime | None = None,
date_to: datetime | None = None, date_to: datetime | None = None,
withdraw_source: str | None = None, withdraw_source: str | None = None,
app_env: str | None = None,
revenue_scope: str = "all",
feed_scene: str | None = None,
) -> dict: ) -> dict:
"""提现详情「用户统计区」10 项。窗口作用于除「现金余额」外的所有项(余额是当前快照)。 """提现详情「用户统计区」10 项。窗口作用于除「现金余额」外的所有项(余额是当前快照)。
口径:激励视频/信息流只统计 granted;数量——视频按条数、信息流按份数(unit_count 累加); 口径:激励视频/信息流奖励数量只统计 granted;数量——视频按条数、信息流按份数(unit_count 累加)
平均 eCPM 用原始分值(分/千次)按记录取算术平均;各「提现」= 该来源累计金币折现。 平均 Draw eCPM 与广告收益报表一致:基于 ad_ecpm_record 的全部 draw/feed 展示记录计算,
不以是否发奖为筛选条件。各「提现」= 该来源累计金币折现。
传统任务 = 窗口内正向金币中,排除广告(reward_video/feed_ad_reward)与人工调整后的折现。 传统任务 = 窗口内正向金币中,排除广告(reward_video/feed_ad_reward)与人工调整后的折现。
""" """
withdraw_source_conds = ( withdraw_source_conds = (
@@ -1296,30 +1301,68 @@ def user_reward_stats(
# 只投影本统计实际使用的列。避免滚动发布或旧本地库尚未补齐无关新列时, # 只投影本统计实际使用的列。避免滚动发布或旧本地库尚未补齐无关新列时,
# SQLAlchemy 因 select(ORM) 自动展开整表字段而让提现详情整体 500。 # SQLAlchemy 因 select(ORM) 自动展开整表字段而让提现详情整体 500。
business_ids: set[str] | None = None
if revenue_scope == "business":
# 与广告收益报表共用正式/测试业务代码位集合,避免两个页面随配置切换后再次漂移。
from app.admin.repositories.ad_revenue import business_code_ids
business_ids = business_code_ids(db, app_env)
rv_conds = [
AdRewardRecord.user_id == user_id,
AdRewardRecord.reward_scene == "reward_video",
AdRewardRecord.status == "granted",
*_window_conds(AdRewardRecord.created_at, date_from, date_to),
]
if app_env is not None:
rv_conds.append(AdRewardRecord.app_env == app_env)
if business_ids is not None:
rv_conds.append(AdRewardRecord.our_code_id.in_(business_ids))
rv = db.execute( rv = db.execute(
select(AdRewardRecord.ecpm_raw, AdRewardRecord.coin).where( select(AdRewardRecord.ecpm_raw, AdRewardRecord.coin).where(
AdRewardRecord.user_id == user_id, *rv_conds,
AdRewardRecord.reward_scene == "reward_video",
AdRewardRecord.status == "granted",
*_window_conds(AdRewardRecord.created_at, date_from, date_to),
) )
).all() ).all()
rv_ecpms = [rewards.parse_ecpm_fen(r.ecpm_raw) for r in rv if r.ecpm_raw] rv_ecpms = [rewards.parse_ecpm_fen(r.ecpm_raw) for r in rv if r.ecpm_raw]
rv_coins = sum(r.coin for r in rv) rv_coins = sum(r.coin for r in rv)
feed = db.execute( feed_reward_conds = [
AdFeedRewardRecord.user_id == user_id,
AdFeedRewardRecord.status == "granted",
*_window_conds(AdFeedRewardRecord.created_at, date_from, date_to),
]
if app_env is not None:
feed_reward_conds.append(AdFeedRewardRecord.app_env == app_env)
if feed_scene is not None:
feed_reward_conds.append(AdFeedRewardRecord.feed_scene == feed_scene)
if business_ids is not None:
feed_reward_conds.append(AdFeedRewardRecord.our_code_id.in_(business_ids))
feed_rewards = db.execute(
select( select(
AdFeedRewardRecord.unit_count, AdFeedRewardRecord.unit_count,
AdFeedRewardRecord.ecpm_raw,
AdFeedRewardRecord.coin, AdFeedRewardRecord.coin,
).where( ).where(
AdFeedRewardRecord.user_id == user_id, *feed_reward_conds,
AdFeedRewardRecord.status == "granted",
*_window_conds(AdFeedRewardRecord.created_at, date_from, date_to),
) )
).all() ).all()
feed_ecpms = [rewards.parse_ecpm_fen(f.ecpm_raw) for f in feed if f.ecpm_raw] feed_coins = sum(f.coin for f in feed_rewards)
feed_coins = sum(f.coin for f in feed)
feed_impression_conds = [
AdEcpmRecord.user_id == user_id,
AdEcpmRecord.ad_type.in_(("draw", "feed")),
*_window_conds(AdEcpmRecord.created_at, date_from, date_to),
]
if app_env is not None:
feed_impression_conds.append(AdEcpmRecord.app_env == app_env)
if feed_scene is not None:
feed_impression_conds.append(AdEcpmRecord.feed_scene == feed_scene)
if business_ids is not None:
feed_impression_conds.append(AdEcpmRecord.our_code_id.in_(business_ids))
feed_impressions = db.execute(
select(AdEcpmRecord.ecpm_raw).where(*feed_impression_conds)
).all()
# 与 ad_revenue.category_stats 相同:每次展示权重相同,非法原值按 parse_ecpm_fen 记 0。
feed_ecpms = [rewards.parse_ecpm_fen(row.ecpm_raw) for row in feed_impressions]
trad_coins = db.execute( trad_coins = db.execute(
select(func.coalesce(func.sum(CoinTransaction.amount), 0)).where( select(func.coalesce(func.sum(CoinTransaction.amount), 0)).where(
@@ -1338,7 +1381,7 @@ def user_reward_stats(
"reward_video_count": len(rv), "reward_video_count": len(rv),
"reward_video_avg_ecpm": round(sum(rv_ecpms) / len(rv_ecpms), 2) if rv_ecpms else 0.0, "reward_video_avg_ecpm": round(sum(rv_ecpms) / len(rv_ecpms), 2) if rv_ecpms else 0.0,
"reward_video_cash_cents": _coins_to_cents(rv_coins), "reward_video_cash_cents": _coins_to_cents(rv_coins),
"feed_count": int(sum(f.unit_count for f in feed)), "feed_count": int(sum(f.unit_count for f in feed_rewards)),
"feed_avg_ecpm": round(sum(feed_ecpms) / len(feed_ecpms), 2) if feed_ecpms else 0.0, "feed_avg_ecpm": round(sum(feed_ecpms) / len(feed_ecpms), 2) if feed_ecpms else 0.0,
"feed_cash_cents": _coins_to_cents(feed_coins), "feed_cash_cents": _coins_to_cents(feed_coins),
} }
+29 -14
View File
@@ -9,7 +9,7 @@ client_max_body_size,见 shaguabijia-admin-web/deploy/nginx/admin.shaguabijia.co
""" """
from __future__ import annotations from __future__ import annotations
from typing import Annotated from typing import Annotated, Literal
from fastapi import APIRouter, Depends, File, HTTPException, Request, UploadFile from fastapi import APIRouter, Depends, File, HTTPException, Request, UploadFile
@@ -27,14 +27,21 @@ router = APIRouter(
) )
def _out(db: AdminDb) -> GuideVideoConfigOut: GuideScene = Literal["coupon", "comparison"]
def _out(db: AdminDb, scene: GuideScene) -> GuideVideoConfigOut:
"""配置 + 播放统计合成响应(四个写接口都以最新状态返回,前端一次同步到位)。""" """配置 + 播放统计合成响应(四个写接口都以最新状态返回,前端一次同步到位)。"""
return GuideVideoConfigOut(**guide_video.get_config(db), **guide_video.play_stats(db)) return GuideVideoConfigOut(
scene=scene,
**guide_video.get_config(db, scene),
**guide_video.play_stats(db, scene),
)
@router.get("", response_model=GuideVideoConfigOut, summary="新手引导视频配置(领券浮层)") @router.get("", response_model=GuideVideoConfigOut, summary="新手引导视频配置(领券浮层)")
def get_config(db: AdminDb) -> GuideVideoConfigOut: def get_config(db: AdminDb, scene: GuideScene = "coupon") -> GuideVideoConfigOut:
return _out(db) return _out(db, scene)
@router.patch("", response_model=GuideVideoConfigOut, summary="改开关/次数/金币(带审计)") @router.patch("", response_model=GuideVideoConfigOut, summary="改开关/次数/金币(带审计)")
@@ -43,21 +50,23 @@ def update_config(
request: Request, request: Request,
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
scene: GuideScene = "coupon",
) -> GuideVideoConfigOut: ) -> GuideVideoConfigOut:
before, after = guide_video.update_config( before, after = guide_video.update_config(
db, db,
enabled=body.enabled, enabled=body.enabled,
max_plays=body.max_plays, max_plays=body.max_plays,
reward_coin=body.reward_coin, reward_coin=body.reward_coin,
scene=scene,
admin_id=admin.id, admin_id=admin.id,
commit=False, commit=False,
) )
write_audit( write_audit(
db, admin, action="guide_video.update", target_type="guide_video", target_id=None, db, admin, action="guide_video.update", target_type="guide_video", target_id=None,
detail={"before": before, "after": after}, ip=get_client_ip(request), commit=False, detail={"scene": scene, "before": before, "after": after}, ip=get_client_ip(request), commit=False,
) )
db.commit() db.commit()
return _out(db) return _out(db, scene)
@router.post("/video", response_model=GuideVideoConfigOut, summary="上传新手引导视频(MP4,带审计)") @router.post("/video", response_model=GuideVideoConfigOut, summary="上传新手引导视频(MP4,带审计)")
@@ -65,23 +74,26 @@ async def upload_video(
request: Request, request: Request,
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
file: UploadFile = File(...), file: Annotated[UploadFile, File()],
scene: GuideScene = "coupon",
) -> GuideVideoConfigOut: ) -> GuideVideoConfigOut:
data = await file.read() data = await file.read()
try: try:
url = media.save_guide_video(data) url = media.save_guide_video(data)
except media.MediaError as e: except media.MediaError as e:
raise HTTPException(status_code=400, detail=str(e)) from e raise HTTPException(status_code=400, detail=str(e)) from e
before, after = guide_video.set_video(db, url, admin_id=admin.id, commit=False) before, after = guide_video.set_video(
db, url, scene=scene, admin_id=admin.id, commit=False
)
write_audit( write_audit(
db, admin, action="guide_video.set_video", target_type="guide_video", target_id=None, db, admin, action="guide_video.set_video", target_type="guide_video", target_id=None,
detail={"before": before.get("video_url"), "after": url, "bytes": len(data)}, detail={"scene": scene, "before": before.get("video_url"), "after": url, "bytes": len(data)},
ip=get_client_ip(request), commit=False, ip=get_client_ip(request), commit=False,
) )
db.commit() db.commit()
# 提交成功后再删旧片,避免新片没落库就把旧片丢了 # 提交成功后再删旧片,避免新片没落库就把旧片丢了
media.delete_guide_video(before.get("video_url")) media.delete_guide_video(before.get("video_url"))
return _out(db) return _out(db, scene)
@router.delete("/video", response_model=GuideVideoConfigOut, summary="移除新手引导视频(带审计)") @router.delete("/video", response_model=GuideVideoConfigOut, summary="移除新手引导视频(带审计)")
@@ -89,13 +101,16 @@ def delete_video(
request: Request, request: Request,
admin: Annotated[AdminUser, Depends(require_role("operator"))], admin: Annotated[AdminUser, Depends(require_role("operator"))],
db: AdminDb, db: AdminDb,
scene: GuideScene = "coupon",
) -> GuideVideoConfigOut: ) -> GuideVideoConfigOut:
"""移除后 /guide-video/start 一律返回 should_play=false,领券浮层回到「只放广告」。""" """移除后 /guide-video/start 一律返回 should_play=false,领券浮层回到「只放广告」。"""
before, after = guide_video.set_video(db, None, admin_id=admin.id, commit=False) before, after = guide_video.set_video(
db, None, scene=scene, admin_id=admin.id, commit=False
)
write_audit( write_audit(
db, admin, action="guide_video.delete_video", target_type="guide_video", target_id=None, db, admin, action="guide_video.delete_video", target_type="guide_video", target_id=None,
detail={"before": before.get("video_url")}, ip=get_client_ip(request), commit=False, detail={"scene": scene, "before": before.get("video_url")}, ip=get_client_ip(request), commit=False,
) )
db.commit() db.commit()
media.delete_guide_video(before.get("video_url")) media.delete_guide_video(before.get("video_url"))
return _out(db) return _out(db, scene)
+8
View File
@@ -88,6 +88,11 @@ def get_user_reward_stats(
withdraw_source: Annotated[ withdraw_source: Annotated[
str | None, Query(pattern="^(coin_cash|invite_cash)$") str | None, Query(pattern="^(coin_cash|invite_cash)$")
] = None, ] = None,
app_env: Annotated[str | None, Query(pattern="^(prod|test)$")] = None,
revenue_scope: Annotated[str, Query(pattern="^(business|all)$")] = "all",
feed_scene: Annotated[
str | None, Query(pattern="^(comparison|coupon|welfare)$")
] = None,
) -> UserRewardStats: ) -> UserRewardStats:
"""提现详情抽屉「用户统计区」。date_from/date_to 都不传 = 注册至今(全量)。""" """提现详情抽屉「用户统计区」。date_from/date_to 都不传 = 注册至今(全量)。"""
if not user_repo.user_exists(db, user_id): if not user_repo.user_exists(db, user_id):
@@ -99,6 +104,9 @@ def get_user_reward_stats(
date_from=date_from, date_from=date_from,
date_to=date_to, date_to=date_to,
withdraw_source=withdraw_source, withdraw_source=withdraw_source,
app_env=app_env,
revenue_scope=revenue_scope,
feed_scene=feed_scene,
) )
) )
+1
View File
@@ -7,6 +7,7 @@ from app.repositories.guide_video import MAX_PLAYS_LIMIT, REWARD_COIN_LIMIT
class GuideVideoConfigOut(BaseModel): class GuideVideoConfigOut(BaseModel):
scene: str
enabled: bool enabled: bool
video_url: str | None = None # 相对地址 /media/guide_video/xxx.mp4;未配片 = None video_url: str | None = None # 相对地址 /media/guide_video/xxx.mp4;未配片 = None
max_plays: int max_plays: int
+1 -1
View File
@@ -60,7 +60,7 @@ class UserRewardStats(BaseModel):
reward_video_avg_ecpm: float # 平均激励视频 eCPM(分/千次) reward_video_avg_ecpm: float # 平均激励视频 eCPM(分/千次)
reward_video_cash_cents: int # 激励视频提现(金币折现) reward_video_cash_cents: int # 激励视频提现(金币折现)
feed_count: int # 累计信息流广告数(granted 份数,unit_count 累加) feed_count: int # 累计信息流广告数(granted 份数,unit_count 累加)
feed_avg_ecpm: float # 平均信息流广告 eCPM(分/千次) feed_avg_ecpm: float # 全部 Draw/feed 实际展示的平均 eCPM(分/千次,含未发奖展示)
feed_cash_cents: int # 信息流广告提现(金币折现) feed_cash_cents: int # 信息流广告提现(金币折现)
+11 -3
View File
@@ -108,7 +108,7 @@ def _harvest_done_blocking(
def _harvest_abort_blocking( def _harvest_abort_blocking(
trace_id: str, status_hint: str, reason: str | None, trace_url: str | None, trace_id: str, status_hint: str, reason: str | None, trace_url: str | None,
) -> None: ) -> int | None:
with SessionLocal() as db: with SessionLocal() as db:
rec = crud_compare.harvest_abort( rec = crud_compare.harvest_abort(
db, trace_id=trace_id, status=status_hint, reason=reason, trace_url=trace_url, db, trace_id=trace_id, status=status_hint, reason=reason, trace_url=trace_url,
@@ -118,6 +118,7 @@ def _harvest_abort_blocking(
extra={"phase": "harvest_abort", extra={"phase": "harvest_abort",
"status": (rec.status if rec else None), "reason": reason}, "status": (rec.status if rec else None), "reason": reason},
) )
return rec.id if rec is not None else None
async def _forward( async def _forward(
@@ -291,7 +292,10 @@ async def trace_epilogue(
@router.post("/trace/finalize", summary="比价 trace 收尾上云 (透传 + 夭折落库)") @router.post("/trace/finalize", summary="比价 trace 收尾上云 (透传 + 夭折落库)")
async def trace_finalize( async def trace_finalize(
request: Request, user: OptionalUser, db: DbSession request: Request,
background_tasks: BackgroundTasks,
user: OptionalUser,
db: DbSession,
) -> dict[str, Any]: ) -> dict[str, Any]:
_ensure_compare_allowed(user, db) _ensure_compare_allowed(user, db)
# 用户终止 / Phase1 未识别没到 done 帧: pricebot 打包半截上云返回 {trace_url}; # 用户终止 / Phase1 未识别没到 done 帧: pricebot 打包半截上云返回 {trace_url};
@@ -302,12 +306,16 @@ async def trace_finalize(
request, "/api/trace/finalize", user, harvest_first_frame=False, request, "/api/trace/finalize", user, harvest_first_frame=False,
) )
try: try:
await run_in_threadpool( record_id = await run_in_threadpool(
_harvest_abort_blocking, trace_id, _harvest_abort_blocking, trace_id,
(meta.get("status") or "cancelled"), (meta.get("status") or "cancelled"),
(meta.get("reason") or meta.get("information")), (meta.get("reason") or meta.get("information")),
(resp.get("trace_url") if isinstance(resp, dict) else None), (resp.get("trace_url") if isinstance(resp, dict) else None),
) )
if record_id is not None:
background_tasks.add_task(
backfill_comparison_llm_cost, record_id, trace_id
)
except Exception as e: # noqa: BLE001 except Exception as e: # noqa: BLE001
logger.warning("harvest_abort failed trace=%s: %s", trace_id, e) logger.warning("harvest_abort failed trace=%s: %s", trace_id, e)
return resp return resp
+242 -35
View File
@@ -22,7 +22,7 @@ import uuid
from collections.abc import Callable from collections.abc import Callable
from dataclasses import dataclass from dataclasses import dataclass
from typing import Any from typing import Any
from urllib.parse import quote from urllib.parse import quote, urlsplit
import httpx import httpx
@@ -69,6 +69,18 @@ class _CachedToken:
_token_cache: dict[str, _CachedToken] = {} _token_cache: dict[str, _CachedToken] = {}
_LOG_SUMMARY_MAX_CHARS = 1500
_LOG_STRING_MAX_CHARS = 200
_SENSITIVE_LOG_KEYS = {
"accesstoken",
"appkey",
"authtoken",
"authorization",
"clientsecret",
"mastersecret",
"sign",
}
def normalize_vendor(push_vendor: str | None) -> str | None: def normalize_vendor(push_vendor: str | None) -> str | None:
if not push_vendor: if not push_vendor:
@@ -115,8 +127,16 @@ def send_notification(
if mock: if mock:
logger.info( logger.info(
"[mock push] vendor=%s token=%s... title=%s body=%s extras=%s", "vendor push mock vendor=%s request=%s",
vendor, token[:12], title, body, extras, vendor,
_log_summary(
{
"push_token": token,
"title": title,
"body": body,
"extras": extras,
}
),
) )
return { return {
"mock": True, "mock": True,
@@ -133,7 +153,28 @@ def send_notification(
"xiaomi": _send_xiaomi, "xiaomi": _send_xiaomi,
"oppo": _send_oppo, "oppo": _send_oppo,
} }
return dispatch[vendor](token, title, body, extras) started = time.perf_counter()
try:
result = dispatch[vendor](token, title, body, extras)
except Exception:
logger.exception(
"vendor push dispatch failed vendor=%s token=%s notification_type=%s "
"elapsed_ms=%.1f",
vendor,
token,
extras.get("type", ""),
_elapsed_ms(started),
)
raise
logger.info(
"vendor push dispatch succeeded vendor=%s token=%s notification_type=%s "
"elapsed_ms=%.1f",
vendor,
token,
extras.get("type", ""),
_elapsed_ms(started),
)
return result
def send_accessibility_disabled( def send_accessibility_disabled(
@@ -159,13 +200,120 @@ def _require(value: str, name: str) -> str:
return value return value
def _elapsed_ms(started: float) -> float:
return (time.perf_counter() - started) * 1000
def _redacted_value(value: Any) -> str:
raw = json.dumps(value, ensure_ascii=False, sort_keys=True, default=str)
digest = hashlib.sha256(raw.encode()).hexdigest()[:10]
return f"<redacted len={len(raw)} sha256={digest}>"
def _sanitize_for_log(value: Any, *, key: str = "", depth: int = 0) -> Any:
"""生成有排障价值、但不暴露服务端鉴权凭据的紧凑日志摘要。"""
normalized_key = "".join(char for char in key.lower() if char.isalnum())
if normalized_key in _SENSITIVE_LOG_KEYS or normalized_key.endswith("secret"):
return _redacted_value(value)
if depth >= 10:
return f"<{type(value).__name__}>"
if isinstance(value, dict):
items = list(value.items())
sanitized = {
str(item_key): _sanitize_for_log(item_value, key=str(item_key), depth=depth + 1)
for item_key, item_value in items[:30]
}
if len(items) > 30:
sanitized["<omitted_keys>"] = len(items) - 30
return sanitized
if isinstance(value, (list, tuple)):
items = list(value)
sanitized = [_sanitize_for_log(item, key=key, depth=depth + 1) for item in items[:20]]
if len(items) > 20:
sanitized.append(f"<omitted_items={len(items) - 20}>")
return sanitized
if isinstance(value, str):
if value.lstrip().startswith(("{", "[")):
try:
decoded = json.loads(value)
except ValueError:
pass
else:
return _sanitize_for_log(decoded, key=key, depth=depth + 1)
if len(value) > _LOG_STRING_MAX_CHARS:
return f"{value[:_LOG_STRING_MAX_CHARS]}…<len={len(value)}>"
return value
def _log_summary(value: Any) -> str:
rendered = json.dumps(
_sanitize_for_log(value),
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
default=str,
)
if len(rendered) > _LOG_SUMMARY_MAX_CHARS:
return f"{rendered[:_LOG_SUMMARY_MAX_CHARS]}…<len={len(rendered)}>"
return rendered
def _raw_log_summary(value: Any) -> str:
"""厂商响应摘要:保留原始字段和值,仅限制单条日志长度。"""
rendered = json.dumps(
value,
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
default=str,
)
if len(rendered) > _LOG_SUMMARY_MAX_CHARS:
return f"{rendered[:_LOG_SUMMARY_MAX_CHARS]}…<len={len(rendered)}>"
return rendered
def _endpoint_for_log(url: str) -> str:
parsed = urlsplit(url)
return f"{parsed.scheme}://{parsed.netloc}{parsed.path}"
def _request_summary(kwargs: dict[str, Any]) -> str:
summary: dict[str, Any] = {}
if "json" in kwargs:
summary["body_type"] = "json"
summary["request_payload"] = kwargs["json"]
elif "data" in kwargs:
summary["body_type"] = "form"
summary["request_payload"] = kwargs["data"]
if kwargs.get("params"):
summary["params"] = kwargs["params"]
if kwargs.get("headers"):
summary["headers"] = kwargs["headers"]
return _log_summary(summary)
def _response_summary(resp: Any, parsed: Any | None = None) -> str:
if parsed is None:
try:
parsed = resp.json()
except ValueError:
pass
if parsed is not None:
return _raw_log_summary(parsed)
return _raw_log_summary(getattr(resp, "text", ""))
def _request_json( def _request_json(
method: str, method: str,
url: str, url: str,
*, *,
vendor: str,
operation: str,
expected_status: tuple[int, ...] = (200,), expected_status: tuple[int, ...] = (200,),
**kwargs: Any, **kwargs: Any,
) -> dict[str, Any]: ) -> dict[str, Any]:
started = time.perf_counter()
request_summary = _request_summary(kwargs)
try: try:
resp = httpx.request( resp = httpx.request(
method, method,
@@ -174,41 +322,82 @@ def _request_json(
**kwargs, **kwargs,
) )
except httpx.HTTPError as e: except httpx.HTTPError as e:
logger.error(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=network_error request=%s response=%s elapsed_ms=%.1f error=%s",
vendor,
operation,
method,
_endpoint_for_log(url),
request_summary,
"<no_response>",
_elapsed_ms(started),
type(e).__name__,
)
raise VendorPushError(f"push http error: {e}") from e raise VendorPushError(f"push http error: {e}") from e
if resp.status_code not in expected_status: if resp.status_code not in expected_status:
logger.error("vendor push http failed url=%s http=%s body=%s", url, resp.status_code, resp.text[:500]) logger.error(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=http_error http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
resp.status_code,
request_summary,
_response_summary(resp),
_elapsed_ms(started),
)
raise VendorPushError(f"push http {resp.status_code}") raise VendorPushError(f"push http {resp.status_code}")
try: try:
return resp.json() data = resp.json()
except ValueError as e: except ValueError as e:
logger.error(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=invalid_json http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
resp.status_code,
request_summary,
_response_summary(resp),
_elapsed_ms(started),
)
raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e
logger.info(
"vendor push http completed vendor=%s operation=%s method=%s endpoint=%s "
"outcome=success http_status=%s request=%s response=%s elapsed_ms=%.1f",
vendor,
operation,
method,
_endpoint_for_log(url),
resp.status_code,
request_summary,
_response_summary(resp, data),
_elapsed_ms(started),
)
return data
def _request_form( def _request_form(
method: str, method: str,
url: str, url: str,
*, *,
vendor: str,
operation: str,
expected_status: tuple[int, ...] = (200,), expected_status: tuple[int, ...] = (200,),
**kwargs: Any, **kwargs: Any,
) -> dict[str, Any]: ) -> dict[str, Any]:
try: return _request_json(
resp = httpx.request( method,
method, url,
url, vendor=vendor,
timeout=settings.PUSH_REQUEST_TIMEOUT_SEC, operation=operation,
**kwargs, expected_status=expected_status,
) **kwargs,
except httpx.HTTPError as e: )
raise VendorPushError(f"push http error: {e}") from e
if resp.status_code not in expected_status:
logger.error("vendor push http failed url=%s http=%s body=%s", url, resp.status_code, resp.text[:500])
raise VendorPushError(f"push http {resp.status_code}")
try:
return resp.json()
except ValueError as e:
raise VendorPushError(f"push invalid json: {resp.text[:200]}") from e
def _cache_get(key: str) -> str | None: def _cache_get(key: str) -> str | None:
@@ -234,6 +423,8 @@ def _honor_access_token() -> str:
data = _request_form( data = _request_form(
"POST", "POST",
settings.HONOR_PUSH_TOKEN_ENDPOINT, settings.HONOR_PUSH_TOKEN_ENDPOINT,
vendor="honor",
operation="authenticate",
data={ data={
"grant_type": "client_credentials", "grant_type": "client_credentials",
"client_id": client_id, "client_id": client_id,
@@ -243,7 +434,7 @@ def _honor_access_token() -> str:
) )
token = data.get("access_token") token = data.get("access_token")
if not token: if not token:
raise VendorPushError(f"honor auth failed: {data}") raise VendorPushError(f"honor auth failed: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), data.get("expires_in")) return _cache_put(cache_key, str(token), data.get("expires_in"))
@@ -270,6 +461,8 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di
data = _request_json( data = _request_json(
"POST", "POST",
settings.HONOR_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id), settings.HONOR_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id),
vendor="honor",
operation="send_notification",
json=payload, json=payload,
headers={ headers={
"Content-Type": "application/json; charset=UTF-8", "Content-Type": "application/json; charset=UTF-8",
@@ -279,7 +472,7 @@ def _send_honor(token: str, title: str, body: str, extras: dict[str, str]) -> di
) )
code = data.get("code") code = data.get("code")
if code is not None and int(code) != 200: if code is not None and int(code) != 200:
raise VendorPushError(f"honor push failed: {data}") raise VendorPushError(f"honor push failed: {_raw_log_summary(data)}")
return data return data
@@ -294,6 +487,8 @@ def _huawei_access_token() -> str:
data = _request_form( data = _request_form(
"POST", "POST",
settings.HUAWEI_PUSH_TOKEN_ENDPOINT, settings.HUAWEI_PUSH_TOKEN_ENDPOINT,
vendor="huawei",
operation="authenticate",
data={ data={
"grant_type": "client_credentials", "grant_type": "client_credentials",
"client_id": app_id, "client_id": app_id,
@@ -303,7 +498,7 @@ def _huawei_access_token() -> str:
) )
token = data.get("access_token") token = data.get("access_token")
if not token: if not token:
raise VendorPushError(f"huawei auth failed: {data}") raise VendorPushError(f"huawei auth failed: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), data.get("expires_in")) return _cache_put(cache_key, str(token), data.get("expires_in"))
@@ -333,6 +528,8 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d
data = _request_json( data = _request_json(
"POST", "POST",
settings.HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id), settings.HUAWEI_PUSH_SEND_ENDPOINT_TEMPLATE.format(app_id=app_id),
vendor="huawei",
operation="send_notification",
json=payload, json=payload,
headers={ headers={
"Content-Type": "application/json; charset=UTF-8", "Content-Type": "application/json; charset=UTF-8",
@@ -340,7 +537,7 @@ def _send_huawei(token: str, title: str, body: str, extras: dict[str, str]) -> d
}, },
) )
if str(data.get("code", "")) != "80000000": if str(data.get("code", "")) != "80000000":
raise VendorPushError(f"huawei push failed: {data}") raise VendorPushError(f"huawei push failed: {_raw_log_summary(data)}")
return data return data
@@ -357,6 +554,8 @@ def _vivo_auth_token() -> str:
data = _request_json( data = _request_json(
"POST", "POST",
settings.VIVO_PUSH_AUTH_ENDPOINT, settings.VIVO_PUSH_AUTH_ENDPOINT,
vendor="vivo",
operation="authenticate",
json={ json={
"appId": app_id, "appId": app_id,
"appKey": app_key, "appKey": app_key,
@@ -366,10 +565,10 @@ def _vivo_auth_token() -> str:
headers={"Content-Type": "application/json"}, headers={"Content-Type": "application/json"},
) )
if int(data.get("result", -1)) != 0: if int(data.get("result", -1)) != 0:
raise VendorPushError(f"vivo auth failed: {data}") raise VendorPushError(f"vivo auth failed: {_raw_log_summary(data)}")
token = data.get("authToken") token = data.get("authToken")
if not token: if not token:
raise VendorPushError(f"vivo auth missing authToken: {data}") raise VendorPushError(f"vivo auth missing authToken: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), 24 * 3600) return _cache_put(cache_key, str(token), 24 * 3600)
@@ -401,6 +600,8 @@ def _send_vivo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
data = _request_json( data = _request_json(
"POST", "POST",
settings.VIVO_PUSH_SEND_ENDPOINT, settings.VIVO_PUSH_SEND_ENDPOINT,
vendor="vivo",
operation="send_notification",
json=payload, json=payload,
headers={ headers={
"Content-Type": "application/json", "Content-Type": "application/json",
@@ -408,7 +609,7 @@ def _send_vivo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
}, },
) )
if int(data.get("result", -1)) != 0: if int(data.get("result", -1)) != 0:
raise VendorPushError(f"vivo push failed: {data}") raise VendorPushError(f"vivo push failed: {_raw_log_summary(data)}")
return data return data
@@ -446,14 +647,16 @@ def _send_xiaomi(token: str, title: str, body: str, extras: dict[str, str]) -> d
data = _request_form( data = _request_form(
"POST", "POST",
settings.XIAOMI_PUSH_SEND_ENDPOINT, settings.XIAOMI_PUSH_SEND_ENDPOINT,
vendor="xiaomi",
operation="send_notification",
data=form, data=form,
headers={"Authorization": f"key={app_secret}"}, headers={"Authorization": f"key={app_secret}"},
) )
code = data.get("code") code = data.get("code")
if code not in (0, "0", None): if code not in (0, "0", None):
raise VendorPushError(f"xiaomi push failed: {data}") raise VendorPushError(f"xiaomi push failed: {_raw_log_summary(data)}")
if str(data.get("result", "ok")).lower() not in ("ok", "success"): if str(data.get("result", "ok")).lower() not in ("ok", "success"):
raise VendorPushError(f"xiaomi push failed: {data}") raise VendorPushError(f"xiaomi push failed: {_raw_log_summary(data)}")
return data return data
@@ -522,6 +725,8 @@ def _oppo_auth_token() -> str:
data = _request_form( data = _request_form(
"POST", "POST",
settings.OPPO_PUSH_AUTH_ENDPOINT, settings.OPPO_PUSH_AUTH_ENDPOINT,
vendor="oppo",
operation="authenticate",
data={ data={
"app_key": app_key, "app_key": app_key,
"timestamp": timestamp, "timestamp": timestamp,
@@ -530,10 +735,10 @@ def _oppo_auth_token() -> str:
headers={"Content-Type": "application/x-www-form-urlencoded"}, headers={"Content-Type": "application/x-www-form-urlencoded"},
) )
if int(data.get("code", -1)) != 0: if int(data.get("code", -1)) != 0:
raise VendorPushError(f"oppo auth failed: {data}") raise VendorPushError(f"oppo auth failed: {_raw_log_summary(data)}")
token = (data.get("data") or {}).get("auth_token") or data.get("auth_token") token = (data.get("data") or {}).get("auth_token") or data.get("auth_token")
if not token: if not token:
raise VendorPushError(f"oppo auth missing auth_token: {data}") raise VendorPushError(f"oppo auth missing auth_token: {_raw_log_summary(data)}")
return _cache_put(cache_key, str(token), 24 * 3600) return _cache_put(cache_key, str(token), 24 * 3600)
@@ -573,6 +778,8 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
data = _request_form( data = _request_form(
"POST", "POST",
settings.OPPO_PUSH_SEND_ENDPOINT, settings.OPPO_PUSH_SEND_ENDPOINT,
vendor="oppo",
operation="send_notification",
data={ data={
"auth_token": auth_token, "auth_token": auth_token,
"message": json.dumps(message, ensure_ascii=False), "message": json.dumps(message, ensure_ascii=False),
@@ -580,5 +787,5 @@ def _send_oppo(token: str, title: str, body: str, extras: dict[str, str]) -> dic
headers={"Content-Type": "application/x-www-form-urlencoded"}, headers={"Content-Type": "application/x-www-form-urlencoded"},
) )
if int(data.get("code", -1)) != 0: if int(data.get("code", -1)) != 0:
raise VendorPushError(f"oppo push failed: {data}") raise VendorPushError(f"oppo push failed: {_raw_log_summary(data)}")
return data return data
+9
View File
@@ -102,12 +102,21 @@ class ComparisonRecord(Base):
# done 帧 information 文案。成功:"在美团找到同店,到手价 ¥X…"; # done 帧 information 文案。成功:"在美团找到同店,到手价 ¥X…";
# 失败:具体原因(如"美团、京东外卖均未找到该商品")。前端在比价失败时当原因展示。 # 失败:具体原因(如"美团、京东外卖均未找到该商品")。前端在比价失败时当原因展示。
information: Mapped[str | None] = mapped_column(String(256), nullable=True) information: Mapped[str | None] = mapped_column(String(256), nullable=True)
# 失败卡「原因」行的展示文案(仅 status=failed 时非空):information 具体则=它;笼统则从
# platform_results 捞出的业务原因(打烊/未起送/找不到店或菜/单点不配送);纯系统失败为 None
# → 端侧显示品牌兜底「网络开小差…」。写路径(harvest_done / upsert_record)落库时派生。
# 见 repositories.comparison._derive_fail_display。
fail_reason: Mapped[str | None] = mapped_column(String(256), nullable=True)
# ===== 明细(JSON,越详细越好)===== # ===== 明细(JSON,越详细越好)=====
# 下单菜品 [{name, qty, specs?}] # 下单菜品 [{name, qty, specs?}]
items: Mapped[list] = mapped_column(_JSON, nullable=False, default=list) items: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
# 逐平台对比 [{platform_id, platform_name, package, price, is_source, rank, coupon_saved, coupon_name, applied_coupons}](price/coupon_saved 单位:元,原样存;coupon_name=优惠来源名;applied_coupons=[{name,amount}] 多券明细) # 逐平台对比 [{platform_id, platform_name, package, price, is_source, rank, coupon_saved, coupon_name, applied_coupons}](price/coupon_saved 单位:元,原样存;coupon_name=优惠来源名;applied_coupons=[{name,amount}] 多券明细)
comparison_results: Mapped[list] = mapped_column(_JSON, nullable=False, default=list) comparison_results: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
# 展示模型统一数组(pricebot done.params.platforms 原样存): 每平台一行、自带
# status/is_best/display/display_order, 记录页据此直接渲染, 不再靠 comparison_results
# + 客户端合并 + 前端派生。老记录/旧客户端为空 → 前端回退老 comparison_results 渲染。
platforms: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
# 目标平台未找到、跳过的菜名 # 目标平台未找到、跳过的菜名
skipped_dish_names: Mapped[list] = mapped_column(_JSON, nullable=False, default=list) skipped_dish_names: Mapped[list] = mapped_column(_JSON, nullable=False, default=list)
# 客户端上报的原始 payload(calibration + done.params 全量),未来取数兜底 # 客户端上报的原始 payload(calibration + done.params 全量),未来取数兜底
+1 -1
View File
@@ -38,7 +38,7 @@ class GuideVideoPlay(Base):
# start_play 捕获 IntegrityError 降级成"这次不放视频"。 # start_play 捕获 IntegrityError 降级成"这次不放视频"。
# 用 unique Index 而非 UniqueConstraint:与迁移里的 create_index 对齐(SQLite 加约束 # 用 unique Index 而非 UniqueConstraint:与迁移里的 create_index 对齐(SQLite 加约束
# 要整表重建),autogenerate 才不会每次报一条假 diff。 # 要整表重建),autogenerate 才不会每次报一条假 diff。
Index("uq_guide_video_play_user_seq", "user_id", "seq", unique=True), 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) id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
+189 -12
View File
@@ -52,6 +52,81 @@ def _product_names_from_items(items: list | None) -> str | None:
return joined[:500] or None return joined[:500] or None
# ---- 失败记录的展示文案(记录页失败卡「原因」行)------------------------------
# information 具体就直出;笼统(_GENERIC_INFO)则从 platform_results 捞一条用户可读的业务
# 原因;捞不到 → None(端侧显示品牌兜底「网络开小差…」)。pricebot 把 store_closed /
# no_delivery 漏成了 status=failed,这里按 reason 关键字补判;打烊类 reason 常带一坨脏店名
# (店名+月售+起送+配送…),统一成简短模板。自动化黑话(搜索失败/读价失败/购物车残留/裸
# FAILED…)不给用户看 → 归入品牌兜底。
# pricebot 组不出具体原因时的笼统 information(线上统计的大头),一律走品牌兜底。
_GENERIC_INFO = {
"比价过程出错,请稍后重试",
"比价出错",
"比价未完成",
"done 参数缺少可验证的目标平台结果",
}
# 干净业务结局 status(直接可信),按展示优先级(越靠前越先选)。
_BIZ_STATUS_PRIORITY = (
"below_minimum",
"no_delivery",
"store_closed",
"items_not_found",
"store_not_found",
)
def _store_closed_text(reason: str | None) -> str:
"""打烊/暂停营业/休息类 reason 常带脏店名元数据 → 只留结论,套简短模板。"""
r = reason or ""
if "暂停营业" in r:
state = "暂停营业"
elif "休息" in r:
state = "休息中"
else:
state = "已打烊"
return f"门店{state},无法比价"
def _target_display_reason(platform_results: dict | None) -> str | None:
"""从逐平台结果里挑一条"可展示给用户"的失败原因;挑不到返回 None。
① status 命中干净业务结局集 → 直接采信(打烊套模板,其余用 reason);
② 补判 pricebot 漏成 status=failed 的两类:打烊(套模板)、单点不配送(reason 本身干净);
自动化黑话(搜索失败/读价失败/购物车残留/裸 FAILED…)一律不展示 → None。"""
pr = platform_results or {}
targets = [
v for v in pr.values() if isinstance(v, dict) and not v.get("is_source")
]
for want in _BIZ_STATUS_PRIORITY: # ① 干净 status 优先
for v in targets:
if v.get("status") == want:
if want == "store_closed":
return _store_closed_text(v.get("reason"))
if v.get("reason"):
return v["reason"]
for v in targets: # ② 漏成 failed 的业务结局补判
if v.get("status") != "failed":
continue
reason = (v.get("reason") or "").strip()
if any(k in reason for k in ("打烊", "暂停营业", "休息")):
return _store_closed_text(reason)
if "单点不配送" in reason:
return reason
return None
def _derive_fail_display(
information: str | None, platform_results: dict | None
) -> str | None:
"""失败记录展示文案:information 具体则直出;笼统则从 platform_results 捞/补判;
都拿不到 → None(端侧品牌兜底)。仅在 status=failed 时调用。"""
info = (information or "").strip()
text = info if (info and info not in _GENERIC_INFO) else _target_display_reason(
platform_results
)
return text[:256] if text else None
def _derive(payload: ComparisonRecordIn) -> dict: def _derive(payload: ComparisonRecordIn) -> dict:
"""从上报 payload 派生结构化列(best/saved/is_source_best/status)。""" """从上报 payload 派生结构化列(best/saved/is_source_best/status)。"""
results = payload.comparison_results results = payload.comparison_results
@@ -89,8 +164,9 @@ def _derive(payload: ComparisonRecordIn) -> dict:
is_source_best = best.is_source if best is not None else None is_source_best = best.is_source if best is not None else None
# status:客户端显式给了就用;否则有"非源且有价"的结果=success,否则 failed # status:优先 pricebot record_status(区分 below_minimum/store_closed) → 客户端显式 status
status = payload.status # → 兜底"非源且有价"=success/否则 failed。record_status 让"未满起送"不再塌缩成 failed。
status = payload.record_status or payload.status
if status is None: if status is None:
has_valid_target = any( has_valid_target = any(
(not r.is_source) and r.price is not None for r in results (not r.is_source) and r.price is not None for r in results
@@ -105,6 +181,11 @@ def _derive(payload: ComparisonRecordIn) -> dict:
"saved_amount_cents": saved_amount_cents, "saved_amount_cents": saved_amount_cents,
"is_source_best": is_source_best, "is_source_best": is_source_best,
"status": status, "status": status,
"fail_reason": (
_derive_fail_display(payload.information, _pr)
if status == "failed"
else None
),
} }
@@ -117,16 +198,35 @@ def upsert_record(
灰度期老客户端 POST /compare/record 走这条,与后端 harvest 按 trace_id reconcile; 灰度期老客户端 POST /compare/record 走这条,与后端 harvest 按 trace_id reconcile;
新客户端不再 POST(改由 compare.py 透传壳 harvest 落库)。 新客户端不再 POST(改由 compare.py 透传壳 harvest 落库)。
""" """
derived = _derive(payload) # 单源派生: 与 harvest_done 一致, payload 带 platforms 时从它派生(唯一真相源
# _derive_from_platforms), 老客户端不带 platforms 时回退 _derive(从 comparison_results)。
if payload.platforms:
derived = _derive_from_platforms(payload.platforms, payload.record_status)
# 对齐 _derive 返回键(#189 fail_reason): 两路径 fields 键集一致, 覆盖已有行时不残留旧值
derived["fail_reason"] = (
_derive_fail_display(payload.information, payload.platform_results or {})
if derived["status"] == "failed"
else None
)
# 单源派生取自 platforms 源行(常无源平台元数据/店名)→ 空则用 payload 兜底不丢字段。
# 下面 fields 不再显式写这四个键, 统一由 derived 提供(否则 dict(store_name=..., **derived)
# 与 _derive_from_platforms 同名键撞键 TypeError)。
for _k in ("store_name", "source_platform_id", "source_platform_name", "source_package"):
if not derived.get(_k):
derived[_k] = getattr(payload, _k)
else:
derived = _derive(payload)
# _derive 只从 comparison_results 派生, 不含源平台四件套 / store_name → 从 payload 补,
# 与上面 platforms 分支键集对齐(fields 统一靠 **derived 提供这些列)。
for _k in ("store_name", "source_platform_id", "source_platform_name", "source_package"):
derived[_k] = getattr(payload, _k)
items = [it.model_dump(exclude_none=True) for it in payload.items] items = [it.model_dump(exclude_none=True) for it in payload.items]
fields = dict( fields = dict(
device_id=payload.device_id, device_id=payload.device_id,
business_type=payload.business_type, business_type=payload.business_type,
store_name=payload.store_name,
product_names=_product_names_from_items(items), product_names=_product_names_from_items(items),
source_platform_id=payload.source_platform_id, # store_name / source_platform_id / source_platform_name / source_package 统一由
source_platform_name=payload.source_platform_name, # derived 提供(见上方两分支补齐), 不在此显式写 —— 否则与 _derive_from_platforms 撞键。
source_package=payload.source_package,
information=payload.information, information=payload.information,
best_deeplink=payload.best_deeplink, best_deeplink=payload.best_deeplink,
trace_url=payload.trace_url, trace_url=payload.trace_url,
@@ -134,6 +234,7 @@ def upsert_record(
skipped_dish_count=payload.skipped_dish_count, skipped_dish_count=payload.skipped_dish_count,
items=items, items=items,
comparison_results=[r.model_dump() for r in payload.comparison_results], comparison_results=[r.model_dump() for r in payload.comparison_results],
platforms=list(payload.platforms or []),
skipped_dish_names=list(payload.skipped_dish_names), skipped_dish_names=list(payload.skipped_dish_names),
# 客户端环境 / 性能(debug,客户端上报;旧客户端为 None) # 客户端环境 / 性能(debug,客户端上报;旧客户端为 None)
device_model=payload.device_model, device_model=payload.device_model,
@@ -206,7 +307,8 @@ def upsert_record(
def _derive_from_results( def _derive_from_results(
results: list[dict], platform_results: dict | None = None results: list[dict], platform_results: dict | None = None,
record_status: str | None = None,
) -> dict: ) -> dict:
"""从 done 帧 comparison_results(pricebot 原始 dict 列表)派生结构化列。 """从 done 帧 comparison_results(pricebot 原始 dict 列表)派生结构化列。
等价 _derive,但吃原始字段(is_source/price/rank/platform_id/store_name...)而非 pydantic 对象。 等价 _derive,但吃原始字段(is_source/price/rank/platform_id/store_name...)而非 pydantic 对象。
@@ -254,7 +356,53 @@ def _derive_from_results(
"saved_amount_cents": saved_amount_cents, "saved_amount_cents": saved_amount_cents,
"is_source_best": best.get("is_source") if best else None, "is_source_best": best.get("is_source") if best else None,
"store_name": (src_row or {}).get("store_name") or None, "store_name": (src_row or {}).get("store_name") or None,
"status": "success" if has_valid_target else "failed", # 记录级结局: 优先用 pricebot 下发的 record_status(区分 below_minimum/store_closed,
# 不再把"未满起送"塌缩成 failed → 记录页不再误报"网络开小差"); 旧 pricebot 未下发时
# 回退老的 success/failed 二态派生, 向后兼容。
"status": record_status or ("success" if has_valid_target else "failed"),
}
def _derive_from_platforms(
platforms: list, record_status: str | None = None,
) -> dict:
"""从 done 帧 platforms(每平台一行、渲染就绪)派生结构化列——**单一真相源**。
best_* 直接取 platforms 里 is_best 的那一行、source_* 取 role=source 行,与前端读的
platforms 天然一致(不再像 _derive_from_results 那样从 comparison_results 二次评最优,
消除"标量列 vs platforms"双源不一致)。platforms 非空时优先走这里;老 pricebot 无
platforms 时调用方回退 _derive_from_results(向后兼容)。"""
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)
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
if source_price_cents is not None and best_price_cents is not None:
saved_amount_cents = source_price_cents - best_price_cents
has_valid_target = any(
p.get("role") != "source" and p.get("price") is not None for p in rows
)
# store_name: 优先源行; recompare 场景源平台自己当目标、源行被目标覆盖(pricebot
# _build_platform_rows 有意去重, platforms 无 role=source 行)→ 回退 best 行 → 首个有店名
# 的行(显示现场实际比到的店), 免得记录页店名空掉兜底显示成"比价"。正常比价有源行不走回退。
store_name = (
(src or {}).get("store_name")
or (best or {}).get("store_name")
or next((p.get("store_name") for p in rows if p.get("store_name")), None)
)
return {
"source_platform_id": (src or {}).get("platform_id"),
"source_platform_name": (src or {}).get("platform_name"),
"source_package": (src or {}).get("package"),
"source_price_cents": source_price_cents,
"best_platform_id": (best or {}).get("platform_id"),
"best_platform_name": (best or {}).get("platform_name"),
"best_price_cents": best_price_cents,
"saved_amount_cents": saved_amount_cents,
"is_source_best": (best.get("role") == "source") if best else None,
"store_name": store_name or None,
"status": record_status or ("success" if has_valid_target else "failed"),
} }
@@ -413,12 +561,40 @@ def harvest_done(
返回 (记录, 是否本次**新**落成 success)——供调用方据此幂等发一次邀请奖。 返回 (记录, 是否本次**新**落成 success)——供调用方据此幂等发一次邀请奖。
行不存在(理论上帧0已建;防御)则新建。""" 行不存在(理论上帧0已建;防御)则新建。"""
results = done_params.get("comparison_results") or [] results = done_params.get("comparison_results") or []
derived = _derive_from_results(results, done_params.get("platform_results")) # 展示模型统一数组(pricebot 新增, 每平台一行自带 status/is_best): 原样存, 记录页据此直渲染。
# 菜品:pricebot 已把源单菜品塞进 comparison_results[源行].items # record_status: 记录级结局(success/below_minimum/store_closed/failed), 覆盖老二态派生。
items = next((r.get("items") or [] for r in results if r.get("is_source")), []) platforms = done_params.get("platforms") or []
record_status = done_params.get("record_status")
# 单源派生: platforms(含 pricebot 权威 is_best)是唯一真相源, best_*/source_*/saved/status
# 全从它取 → 与前端读的 platforms 天然一致; 菜品也取 platforms 源行。老 pricebot 无
# platforms 时回退从 comparison_results 派生(向后兼容)。
if platforms:
derived = _derive_from_platforms(platforms, record_status)
# 菜品优先源行; recompare 无源行 → 回退 best 行 → 首个有菜品的行(同 store_name 回退)
_item_row = (
next((p for p in platforms if isinstance(p, dict) and p.get("role") == "source"), None)
or next((p for p in platforms if isinstance(p, dict) and p.get("is_best")), None)
or next((p for p in platforms if isinstance(p, dict) and p.get("items")), None)
)
items = (_item_row or {}).get("items") or []
else:
derived = _derive_from_results(
results, done_params.get("platform_results"), record_status
)
# pricebot 已把源单菜品塞进 comparison_results[源行].items
items = next((r.get("items") or [] for r in results if r.get("is_source")), [])
# 失败展示原因(#189): platforms / results 两个派生分支的 status 都可能 failed, 统一在此算
fail_reason = (
_derive_fail_display(
done_params.get("information"), done_params.get("platform_results")
)
if derived["status"] == "failed"
else None
)
fields = dict( fields = dict(
business_type=business_type or "food", business_type=business_type or "food",
information=done_params.get("information") or None, information=done_params.get("information") or None,
fail_reason=fail_reason,
# best_deeplink 来自客户端剪贴板采集,harvest 拿不到 → 留空(灰度期 fromComparison 会补; # best_deeplink 来自客户端剪贴板采集,harvest 拿不到 → 留空(灰度期 fromComparison 会补;
# 纯 harvest 行「再次比价」退化为按 package 拉起 App。要精确深链需客户端另传,后续)。 # 纯 harvest 行「再次比价」退化为按 package 拉起 App。要精确深链需客户端另传,后续)。
trace_url=trace_url or done_params.get("trace_url"), trace_url=trace_url or done_params.get("trace_url"),
@@ -426,6 +602,7 @@ def harvest_done(
skipped_dish_count=done_params.get("skipped_dish_count"), skipped_dish_count=done_params.get("skipped_dish_count"),
skipped_dish_names=list(done_params.get("skipped_dish_names") or []), skipped_dish_names=list(done_params.get("skipped_dish_names") or []),
comparison_results=results, comparison_results=results,
platforms=platforms,
items=items, items=items,
product_names=_product_names_from_items(items), product_names=_product_names_from_items(items),
raw_payload=done_params, raw_payload=done_params,
+45 -22
View File
@@ -30,7 +30,11 @@ from app.models.app_config import AppConfig
from app.models.guide_video import GuideVideoPlay from app.models.guide_video import GuideVideoPlay
from app.repositories import wallet as crud_wallet from app.repositories import wallet as crud_wallet
_KEY = "coupon_guide_video" SCENES = ("coupon", "comparison")
_KEY_BY_SCENE = {
"coupon": "coupon_guide_video",
"comparison": "comparison_guide_video",
}
#: 金币流水 biz_type。客户端收益明细按它显示「新手引导视频奖励」。 #: 金币流水 biz_type。客户端收益明细按它显示「新手引导视频奖励」。
BIZ_TYPE = "guide_video" BIZ_TYPE = "guide_video"
@@ -41,7 +45,7 @@ _DEFAULTS: dict[str, Any] = {
"enabled": True, "enabled": True,
"video_url": None, # None/空 = 未配片 → 不下发,浮层照旧放广告 "video_url": None, # None/空 = 未配片 → 不下发,浮层照旧放广告
"max_plays": 3, # 每个账号前 N 次浮层放引导视频 "max_plays": 3, # 每个账号前 N 次浮层放引导视频
"reward_coin": 120, # 每次固定金币 "reward_coin": 100, # 每次固定金币
} }
_FIELDS = tuple(_DEFAULTS.keys()) _FIELDS = tuple(_DEFAULTS.keys())
@@ -65,19 +69,28 @@ def _merge(raw: Any) -> dict[str, Any]:
return out return out
def get_config(db: Session) -> dict[str, Any]: def _config_key(scene: str) -> str:
if scene not in _KEY_BY_SCENE:
raise ValueError(f"unsupported guide video scene: {scene}")
return _KEY_BY_SCENE[scene]
def get_config(db: Session, scene: str = "coupon") -> dict[str, Any]:
"""完整配置 + updated_at(admin 读 / 业务读共用)。""" """完整配置 + updated_at(admin 读 / 业务读共用)。"""
row = db.get(AppConfig, _KEY) row = db.get(AppConfig, _config_key(scene))
cfg = _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 cfg["updated_at"] = row.updated_at.isoformat() if row is not None and row.updated_at else None
return cfg return cfg
def _write(db: Session, value: dict[str, Any], *, admin_id: int, commit: bool) -> dict[str, Any]: def _write(
db: Session, value: dict[str, Any], *, scene: str, admin_id: int, commit: bool
) -> dict[str, Any]:
"""整体覆写该行(value 须为完整字段 dict),返回合并后的完整配置(含 updated_at)。""" """整体覆写该行(value 须为完整字段 dict),返回合并后的完整配置(含 updated_at)。"""
row = db.get(AppConfig, _KEY) key = _config_key(scene)
row = db.get(AppConfig, key)
if row is None: if row is None:
row = AppConfig(key=_KEY, value=value, updated_by_admin_id=admin_id) row = AppConfig(key=key, value=value, updated_by_admin_id=admin_id)
db.add(row) db.add(row)
else: else:
row.value = value # 整体重新赋值,SQLAlchemy 才侦测得到变更 row.value = value # 整体重新赋值,SQLAlchemy 才侦测得到变更
@@ -98,11 +111,12 @@ def update_config(
enabled: bool | None = None, enabled: bool | None = None,
max_plays: int | None = None, max_plays: int | None = None,
reward_coin: int | None = None, reward_coin: int | None = None,
scene: str = "coupon",
admin_id: int, admin_id: int,
commit: bool = True, commit: bool = True,
) -> tuple[dict[str, Any], dict[str, Any]]: ) -> tuple[dict[str, Any], dict[str, Any]]:
"""改开关 / 次数 / 金币(只改传了的字段;视频走 set_video)。返回 (before, after) 供审计。""" """改开关 / 次数 / 金币(只改传了的字段;视频走 set_video)。返回 (before, after) 供审计。"""
row = db.get(AppConfig, _KEY) row = db.get(AppConfig, _config_key(scene))
before = _merge(row.value if row is not None else None) before = _merge(row.value if row is not None else None)
new_value = {k: before[k] for k in _FIELDS} new_value = {k: before[k] for k in _FIELDS}
if enabled is not None: if enabled is not None:
@@ -111,45 +125,52 @@ def update_config(
new_value["max_plays"] = max(0, min(int(max_plays), MAX_PLAYS_LIMIT)) new_value["max_plays"] = max(0, min(int(max_plays), MAX_PLAYS_LIMIT))
if reward_coin is not None: if reward_coin is not None:
new_value["reward_coin"] = max(0, min(int(reward_coin), REWARD_COIN_LIMIT)) new_value["reward_coin"] = max(0, min(int(reward_coin), REWARD_COIN_LIMIT))
after = _write(db, new_value, admin_id=admin_id, commit=commit) after = _write(db, new_value, scene=scene, admin_id=admin_id, commit=commit)
return before, after return before, after
def set_video( def set_video(
db: Session, video_url: str | None, *, 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]]: ) -> tuple[dict[str, Any], dict[str, Any]]:
"""设置/清空引导视频地址。返回 (before, after);before['video_url'] 供调用方删旧文件。""" """设置/清空引导视频地址。返回 (before, after);before['video_url'] 供调用方删旧文件。"""
row = db.get(AppConfig, _KEY) row = db.get(AppConfig, _config_key(scene))
before = _merge(row.value if row is not None else None) before = _merge(row.value if row is not None else None)
new_value = {k: before[k] for k in _FIELDS} new_value = {k: before[k] for k in _FIELDS}
new_value["video_url"] = video_url new_value["video_url"] = video_url
after = _write(db, new_value, admin_id=admin_id, commit=commit) after = _write(db, new_value, scene=scene, admin_id=admin_id, commit=commit)
return before, after return before, after
# ===== 播放计次 ===== # ===== 播放计次 =====
def used_plays(db: Session, user_id: int) -> int: def used_plays(db: Session, user_id: int, scene: str = "coupon") -> int:
"""该账号已用掉的引导视频次数(开播即算,含未发币的)。""" """该账号已用掉的引导视频次数(开播即算,含未发币的)。"""
return int( return int(
db.execute( db.execute(
select(func.count()).select_from(GuideVideoPlay).where( select(func.count()).select_from(GuideVideoPlay).where(
GuideVideoPlay.user_id == user_id GuideVideoPlay.user_id == user_id,
GuideVideoPlay.scene == scene,
) )
).scalar_one() ).scalar_one()
) )
def play_stats(db: Session) -> dict[str, int]: def play_stats(db: Session, scene: str = "coupon") -> dict[str, int]:
"""全站播放统计(admin 页展示):总播放次数 / 其中已发币次数。""" """全站播放统计(admin 页展示):总播放次数 / 其中已发币次数。"""
total = int( total = int(
db.execute(select(func.count()).select_from(GuideVideoPlay)).scalar_one() db.execute(
select(func.count()).select_from(GuideVideoPlay).where(
GuideVideoPlay.scene == scene
)
).scalar_one()
) )
granted = int( granted = int(
db.execute( db.execute(
select(func.count()).select_from(GuideVideoPlay).where( select(func.count()).select_from(GuideVideoPlay).where(
GuideVideoPlay.status == "granted" GuideVideoPlay.status == "granted",
GuideVideoPlay.scene == scene,
) )
).scalar_one() ).scalar_one()
) )
@@ -168,11 +189,12 @@ def start_play(
reward_coin 播完/中途关闭都发的固定金币 reward_coin 播完/中途关闭都发的固定金币
seq / remaining 第几次 / 发完这次还剩几次(仅展示与排查用) seq / remaining 第几次 / 发完这次还剩几次(仅展示与排查用)
""" """
cfg = get_config(db) _config_key(scene)
cfg = get_config(db, scene)
video_url = (cfg.get("video_url") or "").strip() video_url = (cfg.get("video_url") or "").strip()
max_plays = int(cfg.get("max_plays") or 0) max_plays = int(cfg.get("max_plays") or 0)
reward_coin = int(cfg.get("reward_coin") or 0) reward_coin = int(cfg.get("reward_coin") or 0)
used = used_plays(db, user_id) used = used_plays(db, user_id, scene)
def _miss(used_now: int) -> dict[str, Any]: def _miss(used_now: int) -> dict[str, Any]:
return { return {
@@ -194,7 +216,7 @@ def start_play(
scene=scene, scene=scene,
seq=seq, seq=seq,
video_url=video_url, video_url=video_url,
coin=0, coin=reward_coin,
status="playing", status="playing",
completed=0, completed=0,
started_at=datetime.now(rewards.CN_TZ).replace(tzinfo=None), started_at=datetime.now(rewards.CN_TZ).replace(tzinfo=None),
@@ -210,7 +232,7 @@ def start_play(
db.flush() db.flush()
except IntegrityError: except IntegrityError:
db.rollback() db.rollback()
return _miss(used_plays(db, user_id)) return _miss(used_plays(db, user_id, scene))
return { return {
"should_play": True, "should_play": True,
"video_url": video_url, "video_url": video_url,
@@ -241,7 +263,8 @@ def grant_play(
""" """
token = (play_token or "").strip() token = (play_token or "").strip()
# 金币额度以**服务端配置**为准,不信客户端(客户端只上报"播完/关闭")。 # 金币额度以**服务端配置**为准,不信客户端(客户端只上报"播完/关闭")。
coin = int(get_config(db).get("reward_coin") or 0) play = _find_play(db, user_id, token)
coin = int(play.coin if play is not None else 0)
# 幂等核心:把 status 放进 WHERE 做条件更新(compare-and-set),而不是"先读再判再写"。 # 幂等核心:把 status 放进 WHERE 做条件更新(compare-and-set),而不是"先读再判再写"。
# 「播完」与「✕ 关闭」抢跑、或客户端超时重试时,两个请求会都读到 status='playing', # 「播完」与「✕ 关闭」抢跑、或客户端超时重试时,两个请求会都读到 status='playing',
+12
View File
@@ -107,6 +107,13 @@ class ComparisonRecordIn(BaseModel):
# 明细 # 明细
items: list[ComparisonItemIn] = Field(default_factory=list) items: list[ComparisonItemIn] = Field(default_factory=list)
comparison_results: list[ComparisonResultIn] = Field(default_factory=list) comparison_results: list[ComparisonResultIn] = Field(default_factory=list)
# 展示模型统一数组(pricebot done.params.platforms 原样透传): 每平台一行、自带
# status/is_best/display/display_order,记录页据此直渲染。宽松 list[dict] 存(结构由
# pricebot 定,server 只原样落库),前端读它、老记录空时回退 comparison_results。
platforms: list[dict] = Field(default_factory=list)
# 记录级结局(pricebot 下发): success/below_minimum/store_closed/failed。让"未满起送"不再
# 被塌缩成 failed。_derive 优先用它、其次客户端 status、再兜底二态派生。
record_status: str | None = None
# 逐平台结局摘要(含失败平台的细分原因 status: store_not_found/items_not_found/below_minimum/ # 逐平台结局摘要(含失败平台的细分原因 status: store_not_found/items_not_found/below_minimum/
# unsupported/...)。来自 done.params.platform_results,客户端透传;落 raw_payload(不单列), # unsupported/...)。来自 done.params.platform_results,客户端透传;落 raw_payload(不单列),
# admin「卡在哪一步」从这里读。dict{platform_id: {...}} 宽松存(结构由 pricebot 定——是 # admin「卡在哪一步」从这里读。dict{platform_id: {...}} 宽松存(结构由 pricebot 定——是
@@ -173,8 +180,13 @@ class ComparisonRecordOut(BaseModel):
skipped_dish_count: int | None = None skipped_dish_count: int | None = None
status: str status: str
information: str | None = None information: str | None = None
# 失败卡「原因」文案:具体失败给具体原因,纯系统失败为 None(端侧品牌兜底)。见模型 fail_reason。
fail_reason: str | None = None
items: list = [] items: list = []
comparison_results: list = [] comparison_results: list = []
# 展示模型统一数组(每平台一行、自带 status/is_best/display/display_order): 记录页据此
# 直渲染, 不再靠 comparison_results + 前端派生。老记录为空 → 前端回退 comparison_results。
platforms: list = []
skipped_dish_names: list = [] skipped_dish_names: list = []
total_ms: int | None = None total_ms: int | None = None
# 「已下单」(店级):该店名在该用户真实下单(source='compare')里出现过即 True。 # 「已下单」(店级):该店名在该用户真实下单(source='compare')里出现过即 True。
+3 -1
View File
@@ -1,13 +1,15 @@
"""新手引导视频(领券等候浮层前 N 次替代广告)的客户端请求/响应契约。""" """新手引导视频(领券等候浮层前 N 次替代广告)的客户端请求/响应契约。"""
from __future__ import annotations from __future__ import annotations
from typing import Literal
from pydantic import BaseModel, Field from pydantic import BaseModel, Field
class GuideVideoStartIn(BaseModel): class GuideVideoStartIn(BaseModel):
"""开播询问。scene 目前只有 coupon(领券浮层);预留给日后比价等场景。""" """开播询问。scene 目前只有 coupon(领券浮层);预留给日后比价等场景。"""
scene: str = Field(default="coupon", max_length=16) scene: Literal["coupon", "comparison"] = "coupon"
class GuideVideoStartOut(BaseModel): class GuideVideoStartOut(BaseModel):
+1 -1
View File
@@ -135,7 +135,7 @@ def repair_missing_comparison_llm_costs(
select(ComparisonRecord.id, ComparisonRecord.trace_id) select(ComparisonRecord.id, ComparisonRecord.trace_id)
.where( .where(
*date_conditions, *date_conditions,
ComparisonRecord.status.in_(("success", "failed")), ComparisonRecord.status.in_(("success", "failed", "cancelled")),
ComparisonRecord.llm_cost_yuan.is_(None), ComparisonRecord.llm_cost_yuan.is_(None),
) )
.order_by(ComparisonRecord.created_at.desc(), ComparisonRecord.id.desc()) .order_by(ComparisonRecord.created_at.desc(), ComparisonRecord.id.desc())
+109
View File
@@ -219,6 +219,115 @@ def test_user_reward_stats_can_scope_withdrawals_by_account(
assert invite.json()["cash_balance_cents"] == 456 assert invite.json()["cash_balance_cents"] == 456
def test_user_reward_stats_draw_ecpm_uses_all_filtered_impressions(
admin_client: TestClient, admin_token: str
) -> None:
"""Draw 平均 eCPM 应与收益报表一致,不能只平均成功发奖记录。"""
from app.models.ad_ecpm import AdEcpmRecord
from app.models.ad_feed_reward import AdFeedRewardRecord
uid = _seed_user_with_data("13800000024")
created_at = datetime(2038, 1, 15, 4, tzinfo=UTC)
db = SessionLocal()
try:
db.add_all(
[
AdFeedRewardRecord(
client_event_id="reward-stats-granted-high",
ad_session_id="reward-stats-granted-high",
user_id=uid,
reward_date="2038-01-15",
duration_seconds=10,
unit_count=1,
ecpm_raw="9000",
ad_type="draw",
feed_scene="coupon",
app_env="prod",
our_code_id="104098712",
coin=9,
status="granted",
created_at=created_at,
),
AdEcpmRecord(
user_id=uid,
ad_type="draw",
feed_scene="coupon",
ad_session_id="reward-stats-impression-low",
app_env="prod",
our_code_id="104098712",
ecpm_raw="1000",
report_date="2038-01-15",
created_at=created_at,
),
AdEcpmRecord(
user_id=uid,
ad_type="feed",
feed_scene="coupon",
ad_session_id="reward-stats-impression-mid",
app_env="prod",
our_code_id="104098712",
ecpm_raw="3000",
report_date="2038-01-15",
created_at=created_at,
),
# 同用户但不同场景/环境/非业务代码位,均不应进入本次详情筛选。
AdEcpmRecord(
user_id=uid,
ad_type="draw",
feed_scene="comparison",
ad_session_id="reward-stats-other-scene",
app_env="prod",
our_code_id="104098712",
ecpm_raw="7000",
report_date="2038-01-15",
created_at=created_at,
),
AdEcpmRecord(
user_id=uid,
ad_type="draw",
feed_scene="coupon",
ad_session_id="reward-stats-test-env",
app_env="test",
our_code_id="104127529",
ecpm_raw="8000",
report_date="2038-01-15",
created_at=created_at,
),
AdEcpmRecord(
user_id=uid,
ad_type="draw",
feed_scene="coupon",
ad_session_id="reward-stats-non-business",
app_env="prod",
our_code_id="demo-slot",
ecpm_raw="9000",
report_date="2038-01-15",
created_at=created_at,
),
]
)
db.commit()
finally:
db.close()
response = admin_client.get(
f"/admin/api/users/{uid}/reward-stats",
params={
"date_from": "2038-01-15T00:00:00Z",
"date_to": "2038-01-15T23:59:59Z",
"app_env": "prod",
"revenue_scope": "business",
"feed_scene": "coupon",
},
headers=_auth(admin_token),
)
assert response.status_code == 200, response.text
data = response.json()
assert data["feed_count"] == 1
# 全部真实展示 (1000 + 3000) / 2;不能返回成功发奖记录的 9000。
assert data["feed_avg_ecpm"] == 2000.0
def test_user_coin_record_sort_accepts_mixed_timezone_datetimes() -> None: def test_user_coin_record_sort_accepts_mixed_timezone_datetimes() -> None:
"""线上 PostgreSQL 返回 awareSQLite/历史转换可能返回 naive,二者必须可混排。""" """线上 PostgreSQL 返回 awareSQLite/历史转换可能返回 naive,二者必须可混排。"""
naive = datetime(2038, 1, 1, 8, 0) naive = datetime(2038, 1, 1, 8, 0)
+34 -1
View File
@@ -102,6 +102,7 @@ def test_harvest_done_derives_and_newly_success_once(client) -> None:
assert rec.is_source_best is False assert rec.is_source_best is False
assert rec.store_name == "测试店" assert rec.store_name == "测试店"
assert rec.information == "美团更便宜" assert rec.information == "美团更便宜"
assert rec.fail_reason is None # 成功记录不派生失败原因
assert rec.items == [{"name": "肥牛饭", "qty": 1}] assert rec.items == [{"name": "肥牛饭", "qty": 1}]
assert rec.trace_url.endswith("/done/") assert rec.trace_url.endswith("/done/")
# 再来一次(重试 done)→ 已 success,newly_success=False(发奖不重复触发) # 再来一次(重试 done)→ 已 success,newly_success=False(发奖不重复触发)
@@ -110,6 +111,35 @@ def test_harvest_done_derives_and_newly_success_once(client) -> None:
assert newly2 is False assert newly2 is False
def test_harvest_done_failed_derives_fail_reason(client) -> None:
"""failed 记录:记录级 information 笼统,但 fail_reason 从 platform_results 救出具体原因
(id 3030 型:美团系统失败 + 京东 items_not_found → 展示京东那条)。"""
tid = _tid()
done_failed = {
"comparison_results": [
{"platform_id": "taobao_flash", "platform_name": "淘宝闪购",
"package": "com.taobao.taobao", "price": 23.04, "is_source": True, "rank": 1,
"items": [{"name": "肥牛饭", "qty": 1}]},
],
"platform_results": {
"taobao_flash": {"is_source": True, "status": "source", "price": 23.04},
"meituan_waimai": {"is_source": False, "status": "failed",
"reason": "搜索店铺失败, 无法跳转到搜索页"},
"jd_waimai_standalone": {"is_source": False, "status": "items_not_found",
"reason": "京东外卖此店内未找到这些菜品"},
},
"information": "比价过程出错,请稍后重试",
}
with SessionLocal() as db:
crud.harvest_running(db, trace_id=tid, user_id=None)
rec, newly = crud.harvest_done(db, trace_id=tid, user_id=None,
done_params=done_failed)
assert newly is False # 没落成 success
assert rec.status == "failed"
assert rec.fail_reason == "京东外卖此店内未找到这些菜品"
assert rec.information == "比价过程出错,请稍后重试" # 原文案仍留存
def test_harvest_abort_cancels_running(client) -> None: def test_harvest_abort_cancels_running(client) -> None:
tid = _tid() tid = _tid()
with SessionLocal() as db: with SessionLocal() as db:
@@ -233,7 +263,9 @@ def test_trace_finalize_harvests_abort(client) -> None:
with SessionLocal() as db: # 先有 running 行(帧0建的) with SessionLocal() as db: # 先有 running 行(帧0建的)
crud.harvest_running(db, trace_id=tid, user_id=None) crud.harvest_running(db, trace_id=tid, user_id=None)
p, _cap = _mock_pricebot({"trace_url": "https://price.shaguabijia.com/traces/fin/"}) p, _cap = _mock_pricebot({"trace_url": "https://price.shaguabijia.com/traces/fin/"})
with p: with p, patch(
"app.api.v1.compare.backfill_comparison_llm_cost"
) as backfill:
r = client.post("/api/v1/trace/finalize", r = client.post("/api/v1/trace/finalize",
json={"trace_id": tid, "status": "cancelled", "reason": "用户终止"}) json={"trace_id": tid, "status": "cancelled", "reason": "用户终止"})
assert r.status_code == 200 assert r.status_code == 200
@@ -241,6 +273,7 @@ def test_trace_finalize_harvests_abort(client) -> None:
rec = _get(db, tid) rec = _get(db, tid)
assert rec is not None and rec.status == "cancelled" assert rec is not None and rec.status == "cancelled"
assert rec.trace_url.endswith("/fin/") assert rec.trace_url.endswith("/fin/")
backfill.assert_called_once_with(rec.id, tid)
def test_price_step_binds_user_when_authed(client) -> None: def test_price_step_binds_user_when_authed(client) -> None:
+4
View File
@@ -72,6 +72,7 @@ def test_backfill_retries_then_persists_cost(monkeypatch):
def test_repair_batch_only_targets_terminal_missing_rows(monkeypatch): def test_repair_batch_only_targets_terminal_missing_rows(monkeypatch):
missing_id = _record("llm-repair-missing") missing_id = _record("llm-repair-missing")
cancelled_id = _record("llm-repair-cancelled", status="cancelled")
running_id = _record("llm-repair-running", status="running") running_id = _record("llm-repair-running", status="running")
calls = [ calls = [
{ {
@@ -96,12 +97,15 @@ def test_repair_batch_only_targets_terminal_missing_rows(monkeypatch):
) )
assert result["repaired"] >= 1 assert result["repaired"] >= 1
assert "llm-repair-missing" in seen assert "llm-repair-missing" in seen
assert "llm-repair-cancelled" in seen
assert "llm-repair-running" not in seen assert "llm-repair-running" not in seen
with SessionLocal() as db: with SessionLocal() as db:
assert db.get(ComparisonRecord, missing_id).llm_cost_yuan is not None assert db.get(ComparisonRecord, missing_id).llm_cost_yuan is not None
assert db.get(ComparisonRecord, cancelled_id).llm_cost_yuan is not None
assert db.get(ComparisonRecord, running_id).llm_cost_yuan is None assert db.get(ComparisonRecord, running_id).llm_cost_yuan is None
finally: finally:
_delete(missing_id) _delete(missing_id)
_delete(cancelled_id)
_delete(running_id) _delete(running_id)
+122
View File
@@ -0,0 +1,122 @@
"""失败卡展示原因派生(repositories.comparison._derive_fail_display)单元测试。
用例取自线上真实 failed 记录(platform_results 形态),覆盖:
- information 具体 → 直出
- information 笼统 + platform_results 有干净业务结局 → 救援出该原因(id 3030/2964 型)
- information 笼统 + 仅系统失败(搜索失败等黑话)→ None(端侧品牌兜底,id 3027 型)
- store_closed / no_delivery 被 pricebot 漏成 status=failed → 按 reason 关键字补判
- 打烊类脏店名 blob → 统一简短模板
- platform_results 为空 / 非对象 → None
"""
from __future__ import annotations
from app.repositories import comparison as crud
def test_specific_information_passthrough() -> None:
# information 本身具体(未达起送/找不到菜等)→ 直出,不看 platform_results
assert (
crud._derive_fail_display("淘宝闪购未达起送门槛,可加菜凑单后下单", {})
== "淘宝闪购未达起送门槛,可加菜凑单后下单"
)
assert crud._derive_fail_display("未识别到商品", {}) == "未识别到商品"
def test_generic_info_rescued_from_items_not_found() -> None:
# id 3030 型:美团系统失败 + 京东 items_not_found,记录级 information 笼统 → 救出京东那条
pr = {
"eleme": {"is_source": True, "status": "source", "price": 23.04},
"meituan_waimai": {
"is_source": False, "status": "failed",
"reason": "搜索店铺失败, 无法跳转到搜索页",
},
"jd_waimai_standalone": {
"is_source": False, "status": "items_not_found",
"reason": "京东外卖此店内未找到这些菜品",
},
}
assert (
crud._derive_fail_display("比价过程出错,请稍后重试", pr)
== "京东外卖此店内未找到这些菜品"
)
def test_generic_info_rescued_from_store_not_found() -> None:
# id 2964 型:美团系统失败 + 京东 store_not_found → 救出京东相似店铺文案
pr = {
"taobao_flash": {"is_source": True, "status": "source", "price": 127.98},
"meituan": {
"is_source": False, "status": "failed",
"reason": "比价过程出错,请稍后重试",
},
"jd_waimai": {
"is_source": False, "status": "store_not_found",
"reason": "未在京东找到「黔珍味·贵州牛肉蘸水健康菜 (望京店)」相似店铺",
},
}
assert (
crud._derive_fail_display("比价过程出错,请稍后重试", pr)
== "未在京东找到「黔珍味·贵州牛肉蘸水健康菜 (望京店)」相似店铺"
)
def test_generic_info_pure_system_failure_returns_none() -> None:
# id 3027 型:唯一目标平台是自动化黑话失败 → 不给用户看 → None(端侧品牌兜底)
pr = {
"eleme": {"is_source": True, "status": "source", "price": 18.83},
"meituan_waimai": {
"is_source": False, "status": "failed",
"reason": "搜索店铺失败, 无法跳转到搜索页",
},
}
assert crud._derive_fail_display("比价过程出错,请稍后重试", pr) is None
def test_store_closed_leaked_to_failed_is_rescued_and_cleaned() -> None:
# 打烊被漏成 status=failed;reason 常带脏店名 blob → 统一简短模板
pr = {
"jd_waimai": {"is_source": True, "status": "source", "price": 25},
"taobao_flash": {
"is_source": False, "status": "failed",
"reason": "淘宝闪购「沙胆彪炭炉牛杂煲(...),蜂鸟准时达,月售300+,起送¥20」本店已休息,无法比价",
},
}
assert crud._derive_fail_display("比价过程出错,请稍后重试", pr) == "门店休息中,无法比价"
pr2 = {
"taobao_flash": {"is_source": True, "status": "source", "price": 31.83},
"meituan": {
"is_source": False, "status": "failed",
"reason": "美团「奈雪的茶(北京王府井奥莱·香江」门店已打烊,无法比价",
},
}
assert crud._derive_fail_display("比价出错", pr2) == "门店已打烊,无法比价"
def test_no_delivery_leaked_to_failed_is_rescued() -> None:
# 单点不配送被漏成 status=failed;reason 本身干净 → 直接用
reason = "京东外卖该商家所选商品单点不配送,无法进入结算比价"
pr = {
"taobao_flash": {"is_source": True, "status": "source", "price": 20.1},
"jd_waimai_standalone": {
"is_source": False, "status": "failed", "reason": reason,
},
}
assert crud._derive_fail_display("比价过程出错,请稍后重试", pr) == reason
def test_empty_or_missing_platform_results_returns_none() -> None:
# 「比价出错」+ 空 {} / None / 非对象:引擎早夭,无可展示原因 → None
assert crud._derive_fail_display("比价出错", {}) is None
assert crud._derive_fail_display("比价过程出错,请稍后重试", None) is None
assert crud._derive_fail_display("比价出错", []) is None # 老 array 形态,防御
def test_clean_status_wins_over_priority_order() -> None:
# 多个业务结局同现时按 _BIZ_STATUS_PRIORITY 选(below_minimum 优先于 store_not_found)
pr = {
"src": {"is_source": True, "status": "source", "price": 30},
"a": {"is_source": False, "status": "store_not_found", "reason": "未找到店铺A"},
"b": {"is_source": False, "status": "below_minimum", "reason": "B未达起送门槛"},
}
assert crud._derive_fail_display("比价过程出错,请稍后重试", pr) == "B未达起送门槛"
+1 -1
View File
@@ -201,7 +201,7 @@ def test_start_loses_seq_race_degrades_to_ad(client, guide_configured, monkeypat
assert first["should_play"] is True and first["seq"] == 1 assert first["should_play"] is True and first["seq"] == 1
# 本次请求读到的是过期计数 → 仍会算出 seq=1 # 本次请求读到的是过期计数 → 仍会算出 seq=1
monkeypatch.setattr(crud_guide, "used_plays", lambda db, user_id: 0) monkeypatch.setattr(crud_guide, "used_plays", lambda db, user_id, scene="coupon": 0)
with SessionLocal() as db: with SessionLocal() as db:
result = crud_guide.start_play(db, uid) result = crud_guide.start_play(db, uid)
+95
View File
@@ -7,7 +7,9 @@ mock/真发/变量覆盖/站内联动/设备反查/参数校验。厂商 HTTP
from __future__ import annotations from __future__ import annotations
import json import json
import logging
import httpx
import pytest import pytest
from fastapi.testclient import TestClient from fastapi.testclient import TestClient
@@ -110,6 +112,99 @@ def test_huawei_non_success_code_raises(monkeypatch) -> None:
vendor_push.send_notification("huawei", "bad-token", title="t", body="b") vendor_push.send_notification("huawei", "bad-token", title="t", body="b")
def test_vendor_http_logs_device_token_raw_response_and_elapsed(monkeypatch, caplog) -> None:
vendor_push._token_cache.clear()
def _fake_request(method, url, **kwargs): # noqa: ANN001
if url == vendor_push.settings.HUAWEI_PUSH_TOKEN_ENDPOINT:
return _Resp({"access_token": "raw-access-token", "expires_in": 3600})
return _Resp({"code": "80000000", "msg": "Success", "requestId": "vendor-request-1"})
monkeypatch.setattr(vendor_push.settings, "HUAWEI_PUSH_APP_ID", "10086001")
monkeypatch.setattr(vendor_push.settings, "HUAWEI_PUSH_APP_SECRET", "raw-app-secret")
monkeypatch.setattr(vendor_push.httpx, "request", _fake_request)
with caplog.at_level(logging.INFO, logger=vendor_push.logger.name):
vendor_push.send_notification(
"huawei",
"raw-device-token",
title="需要进入日志的标题",
body="需要进入日志的正文",
extras={"type": "withdraw_success", "notificationId": "90001"},
)
messages = "\n".join(record.getMessage() for record in caplog.records)
assert "operation=authenticate" in messages
assert "operation=send_notification" in messages
assert "request=" in messages
assert "response=" in messages
assert "elapsed_ms=" in messages
assert "vendor-request-1" in messages
assert "withdraw_success" in messages
assert "raw-access-token" in messages
assert "raw-device-token" in messages
assert "需要进入日志的标题" in messages
assert "需要进入日志的正文" in messages
assert "<redacted" in messages
assert "raw-app-secret" not in messages
send_http_log = next(
record.getMessage()
for record in caplog.records
if "operation=send_notification" in record.getMessage()
and "vendor push http completed" in record.getMessage()
)
assert "raw-device-token" in send_http_log
assert "raw-access-token" not in send_http_log # 请求头中的厂商鉴权 token 仍脱敏
def test_vendor_http_network_error_logs_request_context(monkeypatch, caplog) -> None:
def _fake_request(*args, **kwargs): # noqa: ANN002, ANN003
raise httpx.ConnectError("connection refused")
monkeypatch.setattr(vendor_push.settings, "XIAOMI_PUSH_APP_SECRET", "raw-xiaomi-secret")
monkeypatch.setattr(vendor_push.httpx, "request", _fake_request)
with (
caplog.at_level(logging.ERROR, logger=vendor_push.logger.name),
pytest.raises(vendor_push.VendorPushError, match="push http error"),
):
vendor_push.send_notification(
"xiaomi",
"raw-xiaomi-token",
title="title",
body="body",
extras={"type": "push_test"},
)
messages = "\n".join(record.getMessage() for record in caplog.records)
assert "outcome=network_error" in messages
assert "response=<no_response>" in messages
assert "operation=send_notification" in messages
assert "elapsed_ms=" in messages
assert "raw-xiaomi-secret" not in messages
assert "raw-xiaomi-token" in messages
@pytest.mark.parametrize(
("request_kwargs", "device_token"),
[
({"json": {"token": ["honor-device-token"]}}, "honor-device-token"),
({"json": {"message": {"token": ["huawei-device-token"]}}}, "huawei-device-token"),
({"json": {"regId": "vivo-device-token"}}, "vivo-device-token"),
({"data": {"registration_id": "xiaomi-device-token"}}, "xiaomi-device-token"),
(
{"data": {"message": '{"target_value":"oppo-device-token"}'}},
"oppo-device-token",
),
],
)
def test_all_vendor_device_token_fields_remain_visible_in_request_summary(
request_kwargs, device_token
) -> None:
summary = vendor_push._request_summary(request_kwargs)
assert device_token in summary
def test_vendor_aliases_normalize() -> None: def test_vendor_aliases_normalize() -> None:
assert vendor_push.normalize_vendor("华为") == "huawei" assert vendor_push.normalize_vendor("华为") == "huawei"
assert vendor_push.normalize_vendor("HMS") == "huawei" assert vendor_push.normalize_vendor("HMS") == "huawei"