Compare commits

..

8 Commits

Author SHA1 Message Date
unknown f4508863cd 修复:收敛白名单迁移链
同步主分支推送绑定隔离迁移,将白名单迁移串接到 push_binding_isolation 之后,避免 Alembic 双 head 阻断部署。
2026-07-31 16:48:32 +08:00
unknown a40c35d41b 合并:同步推送绑定隔离主分支更新 2026-07-31 16:45:15 +08:00
unknown fab49da883 修复:设备白名单支持跨分类配置
设备 ID 按普通字符串处理,移除不同业务来源之间的混选限制。补充批量新增、追加和整主体替换的跨分类回归测试。
2026-07-31 14:38:31 +08:00
unknown 085d610c99 合并:同步短信容灾与广告发奖主分支更新
# Conflicts:
#	app/api/v1/auth.py
#	app/integrations/sms/__init__.py
2026-07-31 11:59:38 +08:00
unknown c67407668e 合并:同步最新主分支 2026-07-30 20:21:29 +08:00
unknown fdbf7a893b 修复:完善白名单配置兼容与审计一致性 2026-07-30 20:21:16 +08:00
unknown 125f0e1a37 Merge remote-tracking branch 'origin/main' into codex/limit-policy-whitelist
# Conflicts:
#	app/admin/routers/guide_video.py
#	app/models/guide_video.py
#	app/repositories/guide_video.py
#	app/schemas/guide_video.py
2026-07-30 15:01:20 +08:00
unknown 3f2ec19ec2 功能:统一限制策略与白名单管理
新增统一限制规则、临时不限和风控免告警白名单,接入比价、短信登录、广告、引导视频与账号冷却等业务链路。补充设备选择、权限、审计、迁移及回归测试。
2026-07-29 19:07:51 +08:00
38 changed files with 100 additions and 1866 deletions
@@ -1,76 +0,0 @@
"""normalize granular comparison outcomes into terminal record statuses
Revision ID: comparison_below_min_success
Revises: limit_policy_global_bundle
"""
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_below_min_success"
down_revision: str | Sequence[str] | None = "limit_policy_global_bundle"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
_FAILED_OUTCOMES = (
"store_closed",
"store_not_found",
"items_not_found",
"no_delivery",
"unsupported",
)
def _comparison_record() -> sa.TableClause:
return sa.table(
"comparison_record",
sa.column("status", sa.String(16)),
sa.column("fail_reason", sa.String(256)),
sa.column("raw_payload", _JSON),
)
def upgrade() -> None:
comparison_record = _comparison_record()
op.execute(
comparison_record.update()
.where(comparison_record.c.status == "below_minimum")
.values(status="success", fail_reason=None)
)
op.execute(
comparison_record.update()
.where(comparison_record.c.status.in_(_FAILED_OUTCOMES))
.values(status="failed")
)
def downgrade() -> None:
comparison_record = _comparison_record()
# The write paths deliberately preserve the granular outcome. Restore only
# rows that this change normalized, without touching ordinary successes.
raw_outcome = sa.func.coalesce(
comparison_record.c.raw_payload["record_status"].as_string(),
comparison_record.c.raw_payload["status"].as_string(),
)
op.execute(
comparison_record.update()
.where(
comparison_record.c.status == "success",
raw_outcome == "below_minimum",
)
.values(status="below_minimum")
)
op.execute(
comparison_record.update()
.where(
comparison_record.c.status == "failed",
raw_outcome.in_(_FAILED_OUTCOMES),
)
.values(status=raw_outcome)
)
-211
View File
@@ -1,211 +0,0 @@
"""correct DeepSeek V4 Flash token price and frozen historical costs
Revision ID: deepseek_v4_flash_price
Revises: comparison_below_min_success
"""
from __future__ import annotations
import json
from collections.abc import Sequence
from typing import Any
import sqlalchemy as sa
from sqlalchemy.dialects import postgresql
from alembic import op
revision: str = "deepseek_v4_flash_price"
down_revision: str | Sequence[str] | None = "comparison_below_min_success"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
_JSON = sa.JSON().with_variant(postgresql.JSONB(), "postgresql")
_CONFIG_KEY = "llm_token_price"
_MODEL = "deepseek-v4-flash"
_OLD_INPUT_PRICE = 3.0
_OLD_OUTPUT_PRICE = 15.0
_NEW_INPUT_PRICE = 1.0
_NEW_OUTPUT_PRICE = 2.0
_CORRECTION_MARKER = "deepseek_v4_flash_price"
_CONFIG_MARKER_KEY = "migration_deepseek_v4_flash_price"
def _decode_object(value: Any) -> dict[str, Any] | None:
if isinstance(value, dict):
return value
if isinstance(value, str):
try:
parsed = json.loads(value)
except (TypeError, ValueError):
return None
return parsed if isinstance(parsed, dict) else None
return None
def _model_tokens(calls: Any) -> tuple[int, int]:
if not isinstance(calls, list):
return 0, 0
input_tokens = 0
output_tokens = 0
for call in calls:
if not isinstance(call, dict) or call.get("error") or call.get("model") != _MODEL:
continue
usage = _decode_object(call.get("usage"))
if usage is None:
continue
input_tokens += int(usage.get("prompt_tokens") or 0)
output_tokens += int(usage.get("completion_tokens") or 0)
return input_tokens, output_tokens
def _app_config_table() -> sa.TableClause:
return sa.table(
"app_config",
sa.column("key", sa.String(64)),
sa.column("value", _JSON),
sa.column("updated_at", sa.DateTime(timezone=True)),
)
def _comparison_table() -> sa.TableClause:
return sa.table(
"comparison_record",
sa.column("id", sa.Integer),
sa.column("llm_calls", _JSON),
sa.column("llm_cost_yuan", sa.Float),
sa.column("llm_price_snapshot", _JSON),
)
def _update_config(conn, *, upgrade: bool) -> None:
table = _app_config_table()
row = conn.execute(
sa.select(table.c.value).where(table.c.key == _CONFIG_KEY)
).mappings().first()
if row is None or not isinstance(row["value"], dict):
return
config = dict(row["value"])
per_model = dict(config.get("per_model") or {})
current = per_model.get(_MODEL)
if upgrade:
# Preserve an operator's explicit model price. The production defect is specifically
# the missing key falling through to the generic 3/15 price.
if current is not None:
return
per_model[_MODEL] = {
"input_per_1m": _NEW_INPUT_PRICE,
"output_per_1m": _NEW_OUTPUT_PRICE,
}
conn.execute(table.insert().values(key=_CONFIG_MARKER_KEY, value=True))
else:
marker_exists = conn.execute(
sa.select(table.c.key).where(table.c.key == _CONFIG_MARKER_KEY)
).scalar_one_or_none()
if marker_exists is None:
return
if current != {
"input_per_1m": _NEW_INPUT_PRICE,
"output_per_1m": _NEW_OUTPUT_PRICE,
}:
conn.execute(table.delete().where(table.c.key == _CONFIG_MARKER_KEY))
return
per_model.pop(_MODEL, None)
conn.execute(table.delete().where(table.c.key == _CONFIG_MARKER_KEY))
config["per_model"] = per_model
conn.execute(
table.update()
.where(table.c.key == _CONFIG_KEY)
# 这是对历史误配置的追溯修正,不是从部署时刻开始的新价格。保留原 updated_at,
# 否则缺失成本回填会把部署前的记录全部排除。
.values(value=config)
)
def _correct_frozen_costs(conn, *, upgrade: bool) -> None:
table = _comparison_table()
rows = conn.execute(
sa.select(
table.c.id,
table.c.llm_calls,
table.c.llm_cost_yuan,
table.c.llm_price_snapshot,
).where(
table.c.llm_cost_yuan.is_not(None),
table.c.llm_price_snapshot.is_not(None),
)
).mappings()
for row in rows:
snapshot = row["llm_price_snapshot"]
if not isinstance(snapshot, dict):
continue
prices = snapshot.get("prices")
if not isinstance(prices, dict):
continue
model_price = prices.get(_MODEL)
if not isinstance(model_price, dict):
continue
if upgrade:
if not (
model_price.get("_source") == "default"
and model_price.get("input_per_1m") == _OLD_INPUT_PRICE
and model_price.get("output_per_1m") == _OLD_OUTPUT_PRICE
):
continue
elif snapshot.get("pricing_correction") != _CORRECTION_MARKER:
continue
input_tokens, output_tokens = _model_tokens(row["llm_calls"])
if input_tokens == 0 and output_tokens == 0:
continue
if upgrade:
delta = (
input_tokens / 1_000_000 * (_OLD_INPUT_PRICE - _NEW_INPUT_PRICE)
+ output_tokens / 1_000_000 * (_OLD_OUTPUT_PRICE - _NEW_OUTPUT_PRICE)
)
corrected_price = {
"input_per_1m": _NEW_INPUT_PRICE,
"output_per_1m": _NEW_OUTPUT_PRICE,
"_source": "per_model",
}
snapshot["pricing_correction"] = _CORRECTION_MARKER
new_cost = max(0.0, float(row["llm_cost_yuan"]) - delta)
else:
delta = (
input_tokens / 1_000_000 * (_OLD_INPUT_PRICE - _NEW_INPUT_PRICE)
+ output_tokens / 1_000_000 * (_OLD_OUTPUT_PRICE - _NEW_OUTPUT_PRICE)
)
corrected_price = {
"input_per_1m": _OLD_INPUT_PRICE,
"output_per_1m": _OLD_OUTPUT_PRICE,
"_source": "default",
}
snapshot.pop("pricing_correction", None)
new_cost = float(row["llm_cost_yuan"]) + delta
updated_prices = dict(prices)
updated_prices[_MODEL] = corrected_price
updated_snapshot = dict(snapshot)
updated_snapshot["prices"] = updated_prices
conn.execute(
table.update()
.where(table.c.id == row["id"])
.values(
llm_cost_yuan=round(new_cost, 6),
llm_price_snapshot=updated_snapshot,
)
)
def upgrade() -> None:
conn = op.get_bind()
_update_config(conn, upgrade=True)
_correct_frozen_costs(conn, upgrade=True)
def downgrade() -> None:
conn = op.get_bind()
_correct_frozen_costs(conn, upgrade=False)
_update_config(conn, upgrade=False)
+10 -29
View File
@@ -166,18 +166,6 @@ def _session_to_row(
point_stats: dict | None = None,
) -> dict:
"""CouponSession ORM → 明细行 dict(主表「领券数据」与「用户全部领券」抽屉共用)。"""
# 中途退出可能发生在第一张券产生终态之前,此时没有逐券事件。
# 明确返回 0/0,让前端区分「退出前无单券结果」与其它状态的埋点缺失。
if point_stats is not None:
point_success_count = point_stats["succeeded"]
point_total_count = point_stats["tried"]
elif r.status == "abandoned":
point_success_count = 0
point_total_count = 0
else:
point_success_count = None
point_total_count = None
point_event_count = point_stats["events"] if point_stats is not None else 0
return {
"id": r.id,
"trace_id": r.trace_id,
@@ -194,9 +182,8 @@ def _session_to_row(
"app_env": r.app_env,
"started_at": r.started_at,
"claimed_count": r.claimed_count,
"point_success_count": point_success_count,
"point_total_count": point_total_count,
"point_event_count": point_event_count,
"point_success_count": point_stats["succeeded"] if point_stats else None,
"point_total_count": point_stats["tried"] if point_stats else None,
"trace_url": r.trace_url,
"ad_revenue_yuan": ad_revenue_yuan,
}
@@ -207,24 +194,21 @@ def _point_scores_by_trace(db: Session, trace_ids: list[str]) -> dict[str, dict[
if not trace_ids:
return {}
succeeded = func.sum(case((CouponClaimEvent.status.in_(_SLOT_OK), 1), else_=0))
tried = func.sum(case((CouponClaimEvent.status.in_(_SLOT_TRIED), 1), else_=0))
rows = db.execute(
select(
CouponClaimEvent.trace_id,
succeeded.label("succeeded"),
tried.label("tried"),
func.count().label("events"),
func.count().label("tried"),
)
.where(
CouponClaimEvent.trace_id.in_(trace_ids),
CouponClaimEvent.status.in_(_SLOT_TRIED),
)
.where(CouponClaimEvent.trace_id.in_(trace_ids))
.group_by(CouponClaimEvent.trace_id)
).all()
return {
trace_id: {
"succeeded": int(success_count or 0),
"tried": int(tried_count or 0),
"events": int(event_count or 0),
}
for trace_id, success_count, tried_count, event_count in rows
trace_id: {"succeeded": int(success_count or 0), "tried": int(tried or 0)}
for trace_id, success_count, tried in rows
if trace_id is not None
}
@@ -416,15 +400,12 @@ def coupon_user_records(db: Session, *, user_id: int, limit: int = 100) -> dict:
total = db.execute(
select(func.count()).select_from(CouponSession).where(CouponSession.user_id == user_id)
).scalar_one()
trace_ids = [r.trace_id for r in rows]
rev_map = crud_ecpm.revenue_yuan_by_trace(db, trace_ids)
point_stats_map = _point_scores_by_trace(db, trace_ids)
rev_map = crud_ecpm.revenue_yuan_by_trace(db, [r.trace_id for r in rows])
return {
"items": [
_session_to_row(
r,
ad_revenue_yuan=rev_map.get(r.trace_id, 0.0),
point_stats=point_stats_map.get(r.trace_id),
)
for r in rows
],
+6 -34
View File
@@ -39,34 +39,6 @@ from app.repositories import activity, ad_ecpm
# 折算成可提现现金时,非广告金币来源的排除集(广告单独统计、人工调整不算"赚取")
_NON_TASK_BIZ_TYPES = ("reward_video", "feed_ad_reward", "admin_grant", "admin_deduct")
# comparison_record historically persisted a few granular business outcomes as
# top-level statuses. Admin filters and metrics expose lifecycle buckets while
# retaining the raw values until the data migration has run everywhere.
_COMPARISON_STATUS_ALIASES = {
"success": ("success", "below_minimum"),
"failed": (
"failed",
"store_closed",
"store_not_found",
"items_not_found",
"no_delivery",
"unsupported",
),
"cancelled": ("cancelled",),
"running": ("running",),
}
_COMPARISON_SUCCESS_STATUSES = _COMPARISON_STATUS_ALIASES["success"]
_COMPARISON_FAILED_STATUSES = _COMPARISON_STATUS_ALIASES["failed"]
_COMPARISON_COMPLETED_STATUSES = (
*_COMPARISON_SUCCESS_STATUSES,
*_COMPARISON_FAILED_STATUSES,
)
def _comparison_status_condition(status: str):
values = _COMPARISON_STATUS_ALIASES.get(status, (status,))
return ComparisonRecord.status.in_(values)
# 信息流点位场景 → 金币记录「赚取途径」展示名;NULL/未知 = 历史未分类。
_FEED_SCENE_LABEL = {
"comparison": "比价信息流",
@@ -362,7 +334,7 @@ def _comparison_conditions(
)
)
if status:
conditions.append(_comparison_status_condition(status))
conditions.append(ComparisonRecord.status == status)
if business_type:
conditions.append(ComparisonRecord.business_type == business_type)
if store:
@@ -447,7 +419,7 @@ def _comparison_duration_aggregate_stmt(conditions: list, status: str, quantiles
),
).where(
*conditions,
_comparison_status_condition(status),
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
@@ -472,7 +444,7 @@ def _comparison_duration_aggregates(
select(ComparisonRecord.total_ms)
.where(
*conditions,
_comparison_status_condition(status),
ComparisonRecord.status == status,
ComparisonRecord.total_ms.is_not(None),
)
.order_by(ComparisonRecord.total_ms)
@@ -502,11 +474,11 @@ def comparison_records_summary(
row = db.execute(
select(
func.count(ComparisonRecord.id),
func.sum(case((ComparisonRecord.status.in_(_COMPARISON_COMPLETED_STATUSES), 1), else_=0)),
func.sum(case((ComparisonRecord.status.in_(_COMPARISON_SUCCESS_STATUSES), 1), else_=0)),
func.sum(case((ComparisonRecord.status.in_(("success", "failed")), 1), else_=0)),
func.sum(case((ComparisonRecord.status == "success", 1), else_=0)),
func.avg(ComparisonRecord.llm_cost_yuan),
func.sum(case((
ComparisonRecord.status.in_(_COMPARISON_SUCCESS_STATUSES)
(ComparisonRecord.status == "success")
& (ComparisonRecord.saved_amount_cents > 0), 1
), else_=0)),
func.sum(case((ComparisonRecord.status == "cancelled", 1), else_=0)),
+1 -9
View File
@@ -566,10 +566,6 @@ def dashboard_overview(
)
).all()
coupon_started = len(period_coupon_sessions)
coupon_abandoned = sum(s.status == "abandoned" for s in period_coupon_sessions)
# 用户主动中途退出不代表领券流程失败,不进入整场成功率样本。
# started / failed 仍留在分母:前者是尚未形成终态的流失,后者是实际执行失败。
coupon_success_denominator = coupon_started - coupon_abandoned
coupon_completed_elapsed = sorted(
s.elapsed_ms
for s in period_coupon_sessions
@@ -735,13 +731,9 @@ def dashboard_overview(
},
"coupon": {
"started": coupon_started,
"abandoned": coupon_abandoned,
"success_denominator": coupon_success_denominator,
"all_success": coupon_all_success,
"success_rate": (
round(coupon_all_success / coupon_success_denominator, 4)
if coupon_success_denominator
else None
round(coupon_all_success / coupon_started, 4) if coupon_started else None
),
"point_success": coupon_point_success,
"points_per_session": coupon_points_per_session,
+2 -2
View File
@@ -35,7 +35,7 @@ def list_comparison_records(
db: AdminDb,
user_id: Annotated[int | None, Query()] = None,
phone: Annotated[str | None, Query(description="手机号前缀")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled|running)$")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled)$")] = None,
business_type: Annotated[str | None, Query()] = None,
store: Annotated[str | None, Query(description="店名子串模糊匹配")] = None,
product: Annotated[str | None, Query(description="商品名子串模糊匹配")] = None,
@@ -66,7 +66,7 @@ def comparison_records_summary(
db: AdminDb,
user_id: Annotated[int | None, Query()] = None,
phone: Annotated[str | None, Query(description="手机号前缀")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled|running)$")] = None,
status: Annotated[str | None, Query(pattern="^(success|failed|cancelled)$")] = None,
business_type: Annotated[str | None, Query()] = None,
store: Annotated[str | None, Query(description="店名子串模糊匹配")] = None,
product: Annotated[str | None, Query(description="商品名子串模糊匹配")] = None,
+1 -2
View File
@@ -20,7 +20,7 @@ class AdminComparisonListItem(BaseModel):
trace_id: str
# admin 是 debug 工具,无条件下发 trace_url(不看 user.debug_trace_enabled)
trace_url: str | None = None
status: str # success / failed / cancelled / running;旧细分值由前端兼容映射
status: str
information: str | None = None
store_name: str | None = None
product_names: str | None = None # 下单商品名派生串(顿号分隔;「商品」列展示 + 商品搜索)
@@ -83,7 +83,6 @@ class AdminComparisonDetail(AdminComparisonListItem):
skipped_dish_count: int | None = None
device_id: str | None = None
items: list = []
platforms: list = [] # pricebot 渲染就绪的逐平台卡片模型(status=ok/业务失败细分)
comparison_results: list = [] # 逐平台对比(价格/rank/coupon/打烊...)
skipped_dish_names: list = []
# 全量环境
+2 -8
View File
@@ -79,16 +79,10 @@ class CouponDataRow(BaseModel):
started_at: datetime = Field(..., description="发起时刻(明细「时间」列)")
claimed_count: int | None = None
point_success_count: int | None = Field(
None,
description="本次成功单券数(success+already_claimed);中途退出且无逐券结果为0,其它无事件为空",
None, description="本次成功单券数(success+already_claimed);无逐券事件为空"
)
point_total_count: int | None = Field(
None,
description="本次尝试单券数(success+already_claimed+failed,不含 skipped);中途退出且无逐券结果为0,其它无事件为空",
)
point_event_count: int = Field(
0,
description="本次全部逐券事件数(含 skipped);用于区分无有效计分事件与完全无事件",
None, description="本次尝试单券数(success+already_claimed+failed,不含 skipped);无逐券事件为空"
)
trace_url: str | None = Field(None, description="pricebot 公网 trace 链接(仅 completed 有);admin 渲染可点链接,无则显示可复制 trace_id")
ad_revenue_yuan: float = Field(
-3
View File
@@ -70,9 +70,6 @@ class DashboardPeriodCoupon(BaseModel):
成功口径 success+already_claimed(与「我的」页累计领券一致)。"""
started: int = 0
# 用户主动中途退出,不计入整场成功率分母。
abandoned: int = 0
success_denominator: int = 0
# 全部领成功的次数:completed 且当日该设备全部点位成功
all_success: int = 0
success_rate: float | None = None
+2 -2
View File
@@ -21,6 +21,7 @@ from __future__ import annotations
import json
import logging
import time
import uuid
from typing import Any
import httpx
@@ -31,7 +32,6 @@ from app.api.deps import DbSession, OptionalUser
from app.core.config import settings
from app.core.logging import trace_id_ctx
from app.core.pricebot_client import get_pricebot_client
from app.core.trace_ids import new_trace_id
from app.core.pricebot_router import pick_pricebot
from app.db.session import SessionLocal
from app.repositories import comparison as crud_compare
@@ -142,7 +142,7 @@ async def _forward(
trace_id = meta.get("trace_id")
minted = False
if not trace_id:
trace_id = new_trace_id()
trace_id = str(uuid.uuid4())
meta["trace_id"] = trace_id
raw = json.dumps(meta).encode() # 仅首帧重新序列化(注入 trace_id);后续帧走原始 bytes
minted = True
+1 -7
View File
@@ -17,7 +17,6 @@ from fastapi import APIRouter, BackgroundTasks, HTTPException, Query, status
from app.api.deps import CurrentUser, DbSession
from app.core import limit_policy
from app.core.trace_ids import new_trace_id
from app.repositories import comparison as crud_compare
from app.repositories import risk as risk_repo
from app.schemas.compare_record import (
@@ -54,10 +53,6 @@ def reserve_compare_start(
scope=risk_repo.SCOPE_ECONOMIC_ACCOUNT,
):
raise HTTPException(status_code=403, detail="账号存在异常,该功能暂不可用")
# trace_id 统一由服务端签发(客户端不带时):预占额度本就是任务的第一个请求,
# 签发与建 running 行合一,此后 Phase1/Phase2/记录/前端日志全链用同一个 id。
# 客户端带了则沿用——老客户端兼容 + 同 trace 重试幂等(reserve_daily_start 按 trace_id 去重)。
trace_id = payload.trace_id or new_trace_id()
try:
policy = limit_policy.resolve(
db,
@@ -68,7 +63,7 @@ def reserve_compare_start(
rec, used = crud_compare.reserve_daily_start(
db,
user_id=user.id,
trace_id=trace_id,
trace_id=payload.trace_id,
business_type=payload.business_type,
device_id=payload.device_id,
limit=policy.limit,
@@ -99,7 +94,6 @@ def reserve_compare_start(
limit=policy.limit,
used=used,
remaining=max(policy.limit - used, 0) if policy.limit is not None else None,
trace_id=trace_id,
)
+5 -31
View File
@@ -21,7 +21,6 @@ from fastapi.concurrency import run_in_threadpool
from app.api.deps import CurrentUser, DbSession
from app.core.config import settings
from app.core.pricebot_client import get_pricebot_client
from app.core.trace_ids import new_trace_id
from app.core.pricebot_router import pick_pricebot
from app.db.session import SessionLocal
from app.repositories import coupon_state as coupon_repo
@@ -31,7 +30,6 @@ from app.schemas.coupon_state import (
CouponPromptShouldShowOut,
CouponPromptShownIn,
CouponSessionIn,
CouponSessionOut,
CouponStatsOut,
)
@@ -177,12 +175,6 @@ async def coupon_step(
)
resp_json = resp.json()
# 每帧响应顶层回传本次任务 trace_id(对齐 compare _forward 的 setdefault):客户端任一帧
# 都能从响应拿到全链 id。**只回显请求里带的、不 mint**——step 是循环接口,每帧签新 id
# 会把一次任务打散;领券 trace_id 的唯一签发点在 /coupon/session (status=started)。
# pricebot 响应顶层本无 trace_id(只有 trace_url),setdefault 不会覆盖任何上游值。
if isinstance(resp_json, dict) and trace_id:
resp_json.setdefault("trace_id", trace_id)
# 领券结果沉淀:每日资产 + 逐次事件;中间帧和 done 全量帧均幂等写库。
# 当前只记录、不参与"要不要领"判断(MVP 先不去重)。写库失败不影响返回。
@@ -212,33 +204,15 @@ async def coupon_step(
return resp_json
@router.post(
"/session",
response_model=CouponSessionOut,
summary="领券任务流水上报(admin 领券数据看板数据源;started 兼签发本轮 trace_id)",
)
def coupon_session(payload: CouponSessionIn, db: DbSession) -> CouponSessionOut:
@router.post("/session", summary="领券任务流水上报(admin 领券数据看板数据源)")
def coupon_session(payload: CouponSessionIn, db: DbSession) -> dict[str, bool]:
"""客户端两段上报一次领券流水(发起 started / 收尾 completed-failed-abandoned),按 trace_id upsert
coupon_session不鉴权(同领券循环 MVP, device_id/trace_id); admin领券数据看板算
发起/完成数耗时分位机型维度写库失败不应连累客户端(本就 fire-and-forget),吞掉返回 ok
trace_id 统一由后端签发:started 不带 trace_id 签发 uuid 并随响应返回,客户端全程用它
(领券 step 循环 / 收尾上报 / 前端运行日志)签发不依赖写库成功写库失败照样返回 trace_id,
后续收尾上报 upsert 会补建行 started trace_id 不签发(收尾没有 id 只能是异常调用,
签发新 id 只会造出一行查不到发起信息的孤儿),不写库trace_id=null 返回
"""
trace_id = payload.trace_id or (
new_trace_id() if payload.status == "started" else None
)
if trace_id is None:
logger.warning(
"coupon session missing trace_id for status=%s (skip write)", payload.status
)
return CouponSessionOut(ok=True, trace_id=None)
发起/完成数耗时分位机型维度写库失败不应连累客户端(本就 fire-and-forget),吞掉返回 ok"""
try:
coupon_repo.upsert_coupon_session(
db,
trace_id=trace_id,
trace_id=payload.trace_id,
device_id=payload.device_id,
status=payload.status,
started_at_ms=payload.started_at_ms,
@@ -255,7 +229,7 @@ def coupon_session(payload: CouponSessionIn, db: DbSession) -> CouponSessionOut:
)
except Exception as e: # noqa: BLE001
logger.warning("coupon session write failed: %s", e)
return CouponSessionOut(ok=True, trace_id=trace_id)
return {"ok": True}
@router.post("/prompt/shown", summary="领券引导窗弹出即上报(按 App 记 shown)")
+2 -4
View File
@@ -54,13 +54,11 @@ def _app_status(db_status: str) -> str:
def _record_out(fb) -> FeedbackRecordOut:
images = fb.images or []
return FeedbackRecordOut(
id=fb.id,
content=fb.content,
scene=getattr(fb, "scene", None),
images=images,
image_thumbnails=[media.feedback_thumbnail_url(url) for url in images],
images=fb.images or [],
status=_app_status(fb.status),
reject_reason=getattr(fb, "reject_reason", None),
reward_coins=getattr(fb, "reward_coins", None),
@@ -86,7 +84,7 @@ async def submit_feedback(
device_model: str = Form(default=""),
rom_name: str = Form(default=""),
android_version: str = Form(default=""),
images: list[UploadFile] = File(default=[]), # noqa: B008 - FastAPI dependency declaration
images: list[UploadFile] = File(default=[]),
) -> FeedbackOut:
content = content.strip()
contact = contact.strip()
+1 -5
View File
@@ -227,11 +227,7 @@ CONFIG_DEFS: dict[str, dict[str, Any]] = {
# 编辑框;set_value 不校验类型,嵌套 JSON 照存。
"llm_token_price": {
"default": {
"per_model": {
"qwen3.5-flash": {"input_per_1m": 0.8, "output_per_1m": 2.0},
# DashScope 华北 2 公网调用原价;必须显式配置,不能落到 3/15 的未知模型兜底价。
"deepseek-v4-flash": {"input_per_1m": 1.0, "output_per_1m": 2.0},
},
"per_model": {"qwen3.5-flash": {"input_per_1m": 0.8, "output_per_1m": 2.0}},
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
"currency": "CNY", "unit": "per_1m_tokens",
},
+2 -90
View File
@@ -11,7 +11,6 @@
from __future__ import annotations
import json
import logging
import os
import secrets
import subprocess
@@ -19,17 +18,8 @@ from decimal import ROUND_HALF_UP, Decimal, InvalidOperation
from pathlib import Path
from typing import Any
from PIL import Image, ImageOps
from app.core.config import settings
logger = logging.getLogger("shagua.media")
_FEEDBACK_DIR = "feedback"
_FEEDBACK_THUMB_DIR = "feedback_thumbs"
_FEEDBACK_THUMB_MAX_PX = 256
_FEEDBACK_THUMB_QUALITY = 78
class MediaError(Exception):
"""上传文件不合法(类型/大小)。调用方转 400。"""
@@ -78,86 +68,8 @@ def save_avatar(user_id: int, data: bytes) -> str:
def save_feedback_image(user_id: int, data: bytes) -> str:
"""保存反馈截图并预生成历史页缩略图,返回原图相对 URL"""
url = _save_image(_FEEDBACK_DIR, user_id, data)
# 缩略图失败不影响反馈受理;读取缩略图 URL 时会按需重试并回退原图。
ensure_feedback_thumbnail(url)
return url
def feedback_thumbnail_url(image_url: str) -> str:
"""把反馈原图 URL 映射成确定的缩略图 URL,不在 records 接口内做图片解码。
上传文件名由服务端生成且不会覆盖旧数据在客户端真正请求可见图片时按需补图
"""
paths = _feedback_thumbnail_paths(image_url)
return paths[2] if paths is not None else image_url
def _feedback_thumbnail_paths(image_url: str) -> tuple[Path, Path, str] | None:
prefix = f"{settings.MEDIA_URL_PREFIX}/{_FEEDBACK_DIR}/"
if not image_url.startswith(prefix):
return None
filename = image_url.removeprefix(prefix)
# 只接受当前目录下的单个文件名,避免数据库脏数据造成路径穿越。
if not filename or Path(filename).name != filename:
return None
source = _media_dir(_FEEDBACK_DIR) / filename
thumb_name = f"{Path(filename).stem}.jpg"
thumb = _media_dir(_FEEDBACK_THUMB_DIR) / thumb_name
thumb_url = f"{settings.MEDIA_URL_PREFIX}/{_FEEDBACK_THUMB_DIR}/{thumb_name}"
return source, thumb, thumb_url
def ensure_feedback_thumbnail(image_url: str) -> Path | None:
"""确保缩略图存在并返回文件;生成失败时回退原图,供动态缩略图路由使用。"""
paths = _feedback_thumbnail_paths(image_url)
if paths is None:
return None
source, thumb, _ = paths
if thumb.is_file():
return thumb
if not source.is_file():
return None
temp = thumb.with_name(f".{thumb.name}.{secrets.token_hex(4)}.tmp")
try:
with Image.open(source) as opened:
image = ImageOps.exif_transpose(opened)
image.thumbnail(
(_FEEDBACK_THUMB_MAX_PX, _FEEDBACK_THUMB_MAX_PX),
Image.Resampling.LANCZOS,
)
if image.mode != "RGB":
image = image.convert("RGB")
image.save(
temp,
format="JPEG",
quality=_FEEDBACK_THUMB_QUALITY,
optimize=True,
)
os.replace(temp, thumb)
return thumb
except (Image.DecompressionBombError, OSError, ValueError):
logger.warning("生成反馈缩略图失败: %s", source, exc_info=True)
return source
finally:
temp.unlink(missing_ok=True)
def feedback_thumbnail_file(filename: str) -> Path | None:
"""由缩略图文件名找到原反馈图并按需生成,非法/不存在返回 None。"""
if not filename or Path(filename).name != filename or Path(filename).suffix.lower() != ".jpg":
return None
stem = Path(filename).stem
for ext in (".jpg", ".png", ".webp"):
original = _media_dir(_FEEDBACK_DIR) / f"{stem}{ext}"
if original.is_file():
original_url = f"{settings.MEDIA_URL_PREFIX}/{_FEEDBACK_DIR}/{original.name}"
return ensure_feedback_thumbnail(original_url)
return None
"""保存反馈截图,返回相对 URL(`/media/feedback/<file>`)"""
return _save_image("feedback", user_id, data)
def save_report_image(user_id: int, data: bytes) -> str:
+4 -8
View File
@@ -234,17 +234,13 @@ def calculate_ad_reward_coin(ecpm: str | int | float | None, count_after_this: i
会铸天量金币;钳在这唯一入口,feed reward_video 回退客户端 eCPM 的路径都护住,且阈值高于
所有真实值,不影响正规发奖
下限(2026-08,产品口径看了就保底 1):**任何输入**都至少 1 金币,与前端展示公式
FeedRewardFormula.singleUnitCoin 完全对齐(那边注释:"无论 eCPM 是否为空、非法或非正数,
单条广告最低都发 1 金币,不能出现 +0")。此前 eCPM 缺失/为 0 返 0,造成两端不一致:
小球显示 +1后端信息流记 too_short 零发;激励视频侧 "0" 字符串还是 truthy绕过
ecpm_missing `if not ecpm_raw` 判定,落成 granted 0 币且白占当日额度/LT 计数
防刷影响:伪造 eCPM0 每天至多多骗 每日上限×1 金币(500 金币=0.05 ),量级可控;
天价伪造仍由 AD_ECPM_MAX_FEN 钳顶把守
下限(2026-07):有真实正 eCPM 时单份至少 1 金币 eCPM 单份收益四舍五入成 0 时兜底为 1,
避免用户看了广告却因数值太小被记 too_short 零发eCPM 缺失/ 0/非法(没有真实广告价值)仍返 0,
不凭空铸币不破坏 ecpm_missing 语义;防刷上限仍由 AD_ECPM_MAX_FEN 钳顶把守
"""
ecpm_yuan = min(parse_ecpm_yuan(ecpm), AD_ECPM_MAX_FEN / 100.0)
if ecpm_yuan <= 0:
return 1 # 保底:缺失/为 0/非法也发 1(镜像前端 validEcpmFen 判非法 → 直接返 1)
return 0
yuan = (ecpm_yuan / 1000.0) * ad_ecpm_factor(ecpm_yuan) * ad_lt_factor(count_after_this)
return max(1, round(yuan * COIN_PER_YUAN))
-29
View File
@@ -1,29 +0,0 @@
"""trace_id 签发(全后端唯一签发口径, 2026-07 起替代裸 uuid4)。
格式: "YYYYMMDD_HHMMSS_" + 12 位小写 hex 随机, 28 字符,
20260731_162254_a1b2c3d4e5f6
Why 带时间前缀: pricebot 落盘目录名/trace_url 尾段**直接用 trace_id 本身**
( pricebot app/utils/trace_ids.py), id/目录/URL 三者合一此前 uuid trace_id
{首帧时刻}_{uuid[:16]} 目录名是两套标识, URL 只有 pricebot 能拼按前缀反查
还有同秒歧义时间用北京时间(CN_TZ)不依赖各机器 TZ 配置, 与业务时区一致
唯一性: 秒级前缀 + 48bit 随机(hex12), 同一秒内碰撞概率可忽略(比价/领券发起 QPS
远低于产生生日碰撞的量级); pricebot 侧同秒多 trace 靠随机段区分(目录精确匹配
llm jsonl 按尾 12 分文件, 不做前缀模糊匹配)
兼容: 三个签发点(compare/startcoupon/session startedcompare.py _forward mint)
统一走这里; 客户端自带 trace_id(老客户端/重试幂等)仍原样沿用pricebot 对老
uuid 格式保持既有目录/短标识行为, 两代 id 并行不冲突
"""
from __future__ import annotations
import uuid
from datetime import datetime
from app.core.rewards import CN_TZ
def new_trace_id() -> str:
"""签发一个自描述 trace_id: 北京时间前缀 + 12 位 hex 随机。"""
return f"{datetime.now(CN_TZ):%Y%m%d_%H%M%S}_{uuid.uuid4().hex[:12]}"
+2 -32
View File
@@ -10,7 +10,7 @@ from collections.abc import AsyncIterator
from contextlib import asynccontextmanager
from pathlib import Path
from fastapi import FastAPI, HTTPException
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse
from fastapi.staticfiles import StaticFiles
@@ -44,7 +44,6 @@ from app.api.v1.tasks import router as tasks_router
from app.api.v1.user import router as user_router
from app.api.v1.wallet import router as wallet_router
from app.api.v1.wxpay import router as wxpay_router
from app.core import media
from app.core.config import settings
from app.core.cps_reconcile_worker import (
start_cps_reconcile_worker,
@@ -83,19 +82,6 @@ setup_logging(debug=settings.APP_DEBUG)
logger = logging.getLogger("shagua.main")
class FeedbackMediaStaticFiles(StaticFiles):
"""反馈原图/缩略图文件名不可变,可长期缓存,避免列表反复回源。"""
async def get_response(self, path: str, scope):
response = await super().get_response(path, scope)
media_path = path.replace("\\", "/").lstrip("/")
if response.status_code == 200 and media_path.startswith(
("feedback/", "feedback_thumbs/")
):
response.headers["Cache-Control"] = "public, max-age=31536000, immutable"
return response
@asynccontextmanager
async def lifespan(_: FastAPI) -> AsyncIterator[None]:
# 提示而非强制建表:生产用 alembic upgrade head,本地 dev 也建议先跑一次 migration。
@@ -226,24 +212,8 @@ def download_apk() -> FileResponse:
)
@app.get(
f"{settings.MEDIA_URL_PREFIX}/feedback_thumbs/{{filename}}",
tags=["feedback"],
include_in_schema=False,
)
def feedback_thumbnail(filename: str) -> FileResponse:
"""旧反馈图按首次可见请求补缩略图;新图上传时已预生成。"""
path = media.feedback_thumbnail_file(filename)
if path is None:
raise HTTPException(status_code=404, detail="图片不存在")
return FileResponse(
path,
headers={"Cache-Control": "public, max-age=31536000, immutable"},
)
app.mount(
settings.MEDIA_URL_PREFIX,
FeedbackMediaStaticFiles(directory=str(_media_root)),
StaticFiles(directory=str(_media_root)),
name="media",
)
+1 -1
View File
@@ -97,7 +97,7 @@ class ComparisonRecord(Base):
total_dish_count: Mapped[int | None] = mapped_column(Integer, nullable=True)
skipped_dish_count: Mapped[int | None] = mapped_column(Integer, nullable=True)
# success(流程正常完成,含 below_minimum)/ failed(技术异常或未形成可比报价,含店铺打烊等)
# success(拿到有效对比)/ failed(出错或没采到目标价)
status: Mapped[str] = mapped_column(String(16), nullable=False, default="success")
# done 帧 information 文案。成功:"在美团找到同店,到手价 ¥X…";
# 失败:具体原因(如"美团、京东外卖均未找到该商品")。前端在比价失败时当原因展示。
+10 -39
View File
@@ -76,30 +76,6 @@ _BIZ_STATUS_PRIORITY = (
)
def _normalize_record_status(status: str | None) -> str | None:
"""Map granular business outcomes onto the record lifecycle status.
``record_status`` describes the business outcome, while
``comparison_record.status`` is also the completed-comparison flag used by
milestones, stats and idempotent rewards. ``below_minimum`` is a completed
success because the target cart produced a trustworthy conclusion. Other
known target-side outcomes did not produce a comparable quote and belong
to the failed record bucket. The granular outcome remains in ``raw_payload``
and ``platform_results`` for result rendering.
"""
if status == "below_minimum":
return "success"
if status in {
"store_closed",
"store_not_found",
"items_not_found",
"no_delivery",
"unsupported",
}:
return "failed"
return status
def _store_closed_text(reason: str | None) -> str:
"""打烊/暂停营业/休息类 reason 常带脏店名元数据 → 只留结论,套简短模板。"""
r = reason or ""
@@ -188,10 +164,9 @@ def _derive(payload: ComparisonRecordIn) -> dict:
is_source_best = best.is_source if best is not None else None
# status:优先 pricebot record_status → 客户端显式 status → 兜底派生。
# below_minimum 是已形成可信结论的正常完成态,记录级归 success;细分结局仍完整保留在
# raw_payload/platform_results,供结果卡展示"未满起送"。
status = _normalize_record_status(payload.record_status or payload.status)
# status:优先 pricebot record_status(区分 below_minimum/store_closed) → 客户端显式 status
# → 兜底"非源且有价"=success/否则 failed。record_status 让"未满起送"不再塌缩成 failed。
status = payload.record_status or payload.status
if status is None:
has_valid_target = any(
(not r.is_source) and r.price is not None for r in results
@@ -226,9 +201,7 @@ def upsert_record(
# 单源派生: 与 harvest_done 一致, payload 带 platforms 时从它派生(唯一真相源
# _derive_from_platforms), 老客户端不带 platforms 时回退 _derive(从 comparison_results)。
if payload.platforms:
derived = _derive_from_platforms(
payload.platforms, payload.record_status or payload.status
)
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 {})
@@ -383,10 +356,10 @@ def _derive_from_results(
"saved_amount_cents": saved_amount_cents,
"is_source_best": best.get("is_source") if best else None,
"store_name": (src_row or {}).get("store_name") or None,
# below_minimum 已完成到购物车并形成可信结论,记录级计 success;细分结局仍在
# raw_payload/platform_results。旧 pricebot 未下发 record_status 时回退二态派生。
"status": _normalize_record_status(record_status)
or ("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"),
}
@@ -439,8 +412,7 @@ def _derive_from_platforms(
"saved_amount_cents": saved_amount_cents,
"is_source_best": (best.get("role") == "source") if best else None,
"store_name": store_name or None,
"status": _normalize_record_status(record_status)
or ("success" if has_valid_target else "failed"),
"status": record_status or ("success" if has_valid_target else "failed"),
}
@@ -612,8 +584,7 @@ def harvest_done(
行不存在(理论上帧0已建;防御)则新建"""
results = done_params.get("comparison_results") or []
# 展示模型统一数组(pricebot 新增, 每平台一行自带 status/is_best): 原样存, 记录页据此直渲染。
# record_status: 记录级业务结局(success/below_minimum/store_closed/failed)。其中
# below_minimum 是正常完成态,持久化 status 归 success,原值仍随 done_params 落 raw_payload。
# record_status: 记录级结局(success/below_minimum/store_closed/failed), 覆盖老二态派生
platforms = done_params.get("platforms") or []
record_status = done_params.get("record_status")
# 单源派生: platforms(含 pricebot 权威 is_best)是唯一真相源, best_*/source_*/saved/status
+7 -21
View File
@@ -25,7 +25,6 @@ from datetime import datetime, timedelta
from sqlalchemy import delete, func, select, update
from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.rewards import CN_TZ
from app.models.comparison import ComparisonRecord
from app.models.ops_marquee_seed import OpsMarqueeSeed
@@ -55,7 +54,7 @@ _SEED_MAX_CENTS = 100000
# 展示层随机(抽样/金额/时间/名字合成)仍每次重算,缓存只省查询;新记录最多晚 30s 进轮播,可接受。
_REAL_ROWS_TTL_SECONDS = 30
_REAL_ROWS_FETCH_CAP = 600 # 一次多取些,够 limit≤30 去重后取数;命中缓存后复用
_real_rows_cache: dict = {"at": None, "rows": None, "test_phones": None}
_real_rows_cache: dict = {"at": None, "rows": None}
# ===== 用户标识脱敏(对齐 PRD) + 种子无真实昵称时的假名合成 =====
# 脱敏规则(按字符数,中英文皆适用):有昵称→n≥5「首+***+末」、n=4「首+**+末」、n≤3「首+**」;
@@ -199,16 +198,10 @@ def _recent_real_rows(db: Session) -> list[tuple[int, int, str | None]]:
返回纯元组(脱离 session),可安全跨请求复用极端并发下偶尔多查一次(无锁幂等),纯门面无副作用
"""
now = datetime.now(CN_TZ)
test_phones = tuple(sorted(settings.test_account_phones))
cached, at = _real_rows_cache["rows"], _real_rows_cache["at"]
if (
cached is not None
and at is not None
and _real_rows_cache["test_phones"] == test_phones
and (now - at).total_seconds() < _REAL_ROWS_TTL_SECONDS
):
if cached is not None and at is not None and (now - at).total_seconds() < _REAL_ROWS_TTL_SECONDS:
return cached
stmt = (
rows = db.execute(
select(ComparisonRecord.user_id, ComparisonRecord.saved_amount_cents, User.nickname)
.join(User, User.id == ComparisonRecord.user_id)
.where(
@@ -218,12 +211,9 @@ def _recent_real_rows(db: Session) -> list[tuple[int, int, str | None]]:
)
.order_by(ComparisonRecord.created_at.desc())
.limit(_REAL_ROWS_FETCH_CAP)
)
if test_phones:
stmt = stmt.where(User.phone.not_in(test_phones))
rows = db.execute(stmt).all()
).all()
out = [(int(uid), int(sc), nick) for uid, sc, nick in rows]
_real_rows_cache.update(rows=out, at=now, test_phones=test_phones)
_real_rows_cache["rows"], _real_rows_cache["at"] = out, now
return out
@@ -350,7 +340,7 @@ def list_real_records(
# pool: [(cluster_key, item)];cluster_key 供去连簇——真实=user_id、种子=各自唯一负数(互不聚簇)
pool: list[tuple[int, dict]] = []
if mode != "seed":
stmt = (
rows = db.execute(
select(
ComparisonRecord.user_id,
ComparisonRecord.saved_amount_cents,
@@ -365,11 +355,7 @@ def list_real_records(
)
.order_by(ComparisonRecord.created_at.desc())
.limit(_REAL_BROWSE_CAP)
)
test_phones = tuple(sorted(settings.test_account_phones))
if test_phones:
stmt = stmt.where(User.phone.not_in(test_phones))
rows = db.execute(stmt).all()
).all()
for uid, sc, nick, ca in rows:
pool.append((
int(uid),
+4 -12
View File
@@ -111,9 +111,8 @@ class ComparisonRecordIn(BaseModel):
# 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。
# below_minimum 表示流程正常完成,持久化主状态归 success;store_closed/items_not_found 等
# 已知无报价结局归 failed。原值仍随 raw_payload 落库,admin/记录页从 platform_results 展示细分结论。
# 记录级结局(pricebot 下发): success/below_minimum/store_closed/failed。让"未满起送"不再
# 被塌缩成 failed。_derive 优先用它、其次客户端 status、再兜底二态派生。
record_status: str | None = None
# 逐平台结局摘要(含失败平台的细分原因 status: store_not_found/items_not_found/below_minimum/
# unsupported/...)。来自 done.params.platform_results,客户端透传;落 raw_payload(不单列),
@@ -215,13 +214,9 @@ class ComparisonRecordCreatedOut(BaseModel):
class CompareStartReserveIn(BaseModel):
"""Reserve one authenticated comparison start before the agent begins.
"""Reserve one authenticated comparison start before the agent begins."""
trace_id 可选不带 = 请服务端签发统一 trace_id 由后端下发前端/SLS 日志/
pricebot 全链用同一个 id = 沿用客户端值老客户端兼容 + 网络重试幂等
"""
trace_id: str | None = Field(default=None, min_length=1, max_length=64)
trace_id: str = Field(..., min_length=1, max_length=64)
business_type: str = Field(default="food", min_length=1, max_length=16)
device_id: str | None = Field(default=None, max_length=64)
@@ -230,9 +225,6 @@ class CompareStartReserveOut(BaseModel):
limit: int | None
used: int
remaining: int | None
# 本次比价全链 trace_id(服务端签发的,或回显客户端带来的)。客户端必须以它为准,
# 贯穿 Phase1/Phase2 step、比价记录、trace 收尾与前端运行日志上报。
trace_id: str
class CompareStatsOut(BaseModel):
+1 -18
View File
@@ -58,13 +58,9 @@ class CouponSessionIn(BaseModel):
- 发起(status=started):带勾选平台 + 机型/ROM/app_env + started_at_ms(发起墙钟毫秒)
- 收尾(completed/failed/abandoned): elapsed_ms(全程耗时)+ platform_elapsed(各平台耗时)+ claimed_count
不鉴权(同领券循环 MVP, device_id/trace_id),user_id 登录态带上做留痕(可空)
trace_id 可选:started 不带 = 请服务端签发本轮领券 trace_id(统一 trace_id 由后端下发,
响应 CouponSessionOut.trace_id 返回,客户端全程用它); = 沿用客户端值(老客户端兼容)
started trace_id 不签发(防孤儿行),返回 trace_id=null 且不写库
"""
trace_id: str | None = None
trace_id: str
device_id: str
status: str # started / completed / failed / abandoned
started_at_ms: int # 发起墙钟毫秒(客户端 System.currentTimeMillis)
@@ -78,16 +74,3 @@ class CouponSessionIn(BaseModel):
platform_elapsed: dict[str, int] | None = None
claimed_count: int | None = None
trace_url: str | None = None
class CouponSessionOut(BaseModel):
"""POST /api/v1/coupon/session 响应。
trace_id = 本轮领券全链 id(服务端签发的,或回显客户端带来的);客户端以它为准贯穿
/coupon/step 循环收尾上报与前端运行日志 不能沿用旧的 dict[str, bool] 返回注解
FastAPI 会按注解校验响应,字符串 trace_id bool 校验必炸,故显式建模
started 且缺 trace_id 时为 null(不签发防孤儿行)
"""
ok: bool = True
trace_id: str | None = None
-2
View File
@@ -33,8 +33,6 @@ class FeedbackRecordOut(BaseModel):
# 比价反馈的问题场景(找错商品/优惠不对…);普通反馈为 None
scene: str | None = None
images: list[str] = Field(default_factory=list)
# 与 images 下标一一对应;生成失败时该项回退原图 URL,兼容历史数据。
image_thumbnails: list[str] = Field(default_factory=list)
status: str
reject_reason: str | None = None
reward_coins: int | None = None
@@ -1,304 +0,0 @@
# 线上 PostgreSQL 数据库备份 / 恢复方案 设计
- 日期:2026-08-03
- 状态:设计已评审,待写实现计划
- 相关代码(复用的现成范式):
- [deploy/daily-exchange.service](../../../deploy/daily-exchange.service) / [deploy/daily-exchange.timer](../../../deploy/daily-exchange.timer)systemd oneshot + timer 定时任务范式:`TZ=Asia/Shanghai``Persistent=true` 补跑、`ProtectSystem=strict` 加固、文件锁)
- [app/core/heartbeat_monitor_worker.py](../../../app/core/heartbeat_monitor_worker.py)(单实例文件锁 + 优雅退出范式)
- [app/integrations/notifier.py](../../../app/integrations/notifier.py)(可插拔通知器:`LogNotifier` 占位、协议不变后续替换)
- [app/admin/deps.py](../../../app/admin/deps.py)`require_role("super_admin")` 权限守卫、审计 IP
- [app/admin/routers/withdraw.py](../../../app/admin/routers/withdraw.py)admin router 风格范式)
- [scripts/init_postgres.py](../../../scripts/init_postgres.py)PG 运维脚本范式)
## 1. 背景与目标
线上业务库 `shaguabijia`PostgreSQL 16,当前数据量 1–10GB)目前没有系统化的备份/恢复方案。本方案提供覆盖「定时 + 手动 + 平台恢复 + 平台不可用兜底」四条路径的逻辑备份体系。
四条明确需求:
1. **每日定时备份** —— 无人值守,不依赖平台进程。
2. **平台手动备份** —— admin 后台一键触发。
3. **平台指定备份文件快速恢复** —— 在 admin 后台选一份备份,恢复到**旁库**供核对。
4. **平台不可用时手动恢复** —— 登录服务器,交互式选择备份恢复,与平台走**同一套脚本**。
数据规模 110GB → **逻辑备份(`pg_dump`)足够**,不引入物理备份 / PITR 的复杂度(见 §2 预留路径)。
## 2. 非目标(本期不做 / 预留升级路径)
- **物理备份 + PITR**`pg_basebackup` + WAL 归档):数据量涨到「凌晨全量 dump 也影响业务」时再上。§3 的架构不阻碍后续叠加。
- **异地对象存储(S3 / 阿里云 OSS)**:本期只落**本地磁盘**。备份脚本预留 `upload_to_remote()` 可插拔 hook`BACKUP_S3_ENABLED=false` 时直接返回),等 S3 服务确认后填充,主流程不改。
- **从只读副本 dump**:搭流复制 standby、改从副本备份以对主库零影响,属后续升级;那个副本还能顺带承载 PITR。
- **平台「一键覆盖生产」按钮**:刻意不做。扶正生产(旁库→主库)永远是人工 SSH 执行的高危脚本(§6)。
- **平台下载备份文件**:整库备份含手机号 / 微信 / 提现等敏感数据,经浏览器下载放大泄露面且文件大。取文件走 SSH。
- **真实推送告警渠道**:项目当前无真实推送能力(`notifier.py``LogNotifier` 占位、心跳 worker 也只打印)。本期告警走「日志 + openobserve + 平台新鲜度徽标」,推送做占位 hook(§8)。
## 3. 核心设计原则与整体架构
### 3.1 两条铁律
1. **执行逻辑沉到自包含脚本,平台与定时器都只是调用者。** —— 需求 4(平台不可用时手动恢复)因此不是另写一套,而是复用同一套脚本,逻辑不重复、兜底路径永远可用。
2. **备份清单的权威来源是文件系统(备份目录 + 每目录内的 `manifest.json`),不是数据库任务表。** —— 任务表只记录「平台发起的操作过程」用于展示进度与审计;删掉任务表不影响任何一份备份的可用性与可恢复性。最需要恢复的时刻(DB 崩了)恰恰是任务表也读不到的时刻,而恢复根本不读任务表。
### 3.2 分层架构
```
┌─────────────────────────────────────────┐
│ 核心执行层(自包含,不依赖平台/DB 任务表) │
│ scripts/db_backup.sh ← 生成备份 │
│ scripts/db_restore.sh ← 从备份恢复到旁库 │
│ scripts/db_promote.sh ← 旁库扶正生产(高危)│
└─────────────────────────────────────────┘
▲ ▲ ▲
┌─────────────┘ │ └──────────────┐
① 每日定时 ② 平台手动备份/恢复 ④ 平台挂了,人 SSH
systemd timer admin API → setsid 起脚本 跑 db_restore.sh
(不经过平台) (任务表记状态,web 不阻塞) (交互式选备份)
```
### 3.3 平台执行方式:方案 B(轻量版)
平台触发的「几分钟重操作」不放进 web 请求,采用**任务表解耦 + 独立进程执行**:
| | 采用 | 说明 |
|---|---|---|
| 谁拉起脚本 | admin API 插任务后 `setsid` 起一个**脱离 web 进程组**的 subprocess 跑脚本 | 不新增常驻 worker 服务,最贴合项目「脚本 + systemd」风格;web 重启不影响已起的脚本 |
| 并发控制 | `db_backup.sh``flock` 文件锁天然串行 | 防定时与手动同时跑、上一轮未完下一轮又起 |
| 升级路径 | 任务量大、需排队/限流/重试时,升级为常驻 worker 轮询任务表 | 本期低频操作用不上 |
**被否决的方案 A**admin API 里同步 `subprocess` 跑 dump/restore,请求一直挂着。否决原因:几分钟操作阻塞 web 进程、易请求超时;admin 进程需持有 DB 高权限;平台一挂手动备份路径也没了(解耦度反而更差)。
## 4. 备份设计
### 4.1 备份产物:一个自描述目录
```
$BACKUP_DIR/20260803_030000_daily/
├── shaguabijia.dump # pg_dump -Fc(自定义格式,内建压缩)
├── globals.sql # pg_dumpall --globals-only --no-role-passwords
├── manifest.json # 元信息(见下)
└── SHA256SUMS # 上述两个文件的校验和
```
目录名(= `backup_id`)格式:`YYYYMMDD_HHMMSS_<reason>`(北京时 `Asia/Shanghai`,纯数字 + 下划线),如 `20260803_030000_daily``reason ∈ {daily, manual}`。刻意不用连字符/冒号/字母,使其**同时是合法文件名与合法 PG 库名**——旁库名由它直接拼接(见 §5.1)。
`manifest.json` 字段:
| 字段 | 说明 |
|---|---|
| `backup_id` | = 目录名,全局唯一标识 |
| `created_at` | ISO8601 北京时 |
| `reason` | `daily` / `manual` |
| `triggered_by` | `systemd-timer` / admin 用户名 / `ssh-manual` |
| `pg_version` | 备份时 PG 版本 |
| `database` | `shaguabijia` |
| `format` | `custom` |
| `alembic_version` | 备份时的 schema 版本(恢复时对齐迁移用;旁库 `alembic_version` 表亦可查) |
| `files` | 每个文件的 `name` / `size_bytes` / `sha256` |
| `dump_duration_sec` | 导出耗时 |
| `status` | `success` / `failed`**仅自检通过才写 success** |
### 4.2 `db_backup.sh` 关键行为
| 环节 | 做法 | 理由 |
|---|---|---|
| 导出业务库 | `pg_dump -Fc``shaguabijia.dump` | 单文件、压缩,`pg_restore` 支持并行(`-j`)与按表选择性恢复 |
| 导出全局角色 | `pg_dumpall --globals-only --no-role-passwords``globals.sql` | 恢复到新机器时角色/权限齐全;`--no-role-passwords` 免超级用户读 `pg_authid`、旁库核对也不需要密码 |
| 资源降级 | `nice -n 19 ionice -c3` 包裹 pg_dump | 让业务优先,压制备份对 CPU/IO 的抢占 |
| **备份后自检** | dump 完立刻 `pg_restore --list *.dump >/dev/null` 解析 TOC + 校验 sha256 | 不验证的备份是薛定谔的备份;能抓到文件截断/损坏,只有通过才写 `status=success` |
| 并发锁 | `flock` 独占锁(仿 daily-exchange 30min 锁) | 串行化,防重入 |
| 保留清理 | **仅自检通过后**执行(§4.4) | 保证先有新备份、再删旧的 |
| 远程上传 | 末尾 `upload_to_remote()``BACKUP_S3_ENABLED=false` 时 return 0 | S3 可插拔扩展点 |
| 通知 | 末尾 `notify()`(§8),成功/失败都打结构化日志 | 贴合现状、可插拔 |
| 退出码 | 失败非 0 退出 + `SyslogIdentifier=pg-backup` | systemd 与 openobserve 可感知 |
| 触发标签 | `--reason daily|manual`,写进 manifest | 定时与手动共用一个脚本 |
### 4.3 对线上读写的影响与缓解
`pg_dump` 基于 MVCC 一致性快照,**不阻塞正常增删改查**,导出的是某一时刻的一致性视图。要点:
- **唯一会互斥的例外是 DDL**dump 持 `ACCESS SHARE` 锁,与 `ALTER TABLE` / `DROP` / `TRUNCATE` / `VACUUM FULL` / 非并发建索引冲突。规避:alembic 迁移仅在手动部署时跑,与凌晨备份天然错开。
- **真正的影响是资源争抢**(1–10GB 下的关注点):磁盘 I/O(整库顺序读)、CPU(zlib 压缩)、缓存冲刷(热数据被挤出 shared_buffers,备份后短时命中率下降)、长事务期间 VACUUM 暂时回收不了死元组。
- **缓解**:定时放凌晨低谷(03:00,与 0 点 daily-exchange 错开)+ `nice`/`ionice` 降级 + 控制 dump 时长在几分钟内。
### 4.4 保留策略
清理只在**备份成功且自检通过后**执行,且**只删本地、不碰 S3**(S3 用自身生命周期策略):
- `BACKUP_KEEP_DAILY_DAYS=14` —— 保留最近 14 天每日备份。
- `BACKUP_KEEP_MONTHLY_COUNT=6` —— 更早的备份中,每月 1 号那份额外保留 6 个月(防「问题两周前就埋下」)。
- `BACKUP_MIN_KEEP=3` —— **硬底线:无论配置如何,永远至少保留最近 3 份**,防时间跳变/配置错误把备份删光。
- 清理前检查磁盘剩余空间,不足则打告警日志而非静默继续。
## 5. 恢复设计(恢复到旁库)
**铁律:平台与默认脚本永远只恢复到旁库,绝不自动碰生产。**
### 5.1 `db_restore.sh` —— 恢复到旁库(常规路径)
输入一个 `backup_id`,产出可供核对的旁库 `shaguabijia_restore_<backup_id>`(如 `shaguabijia_restore_20260803_030000_daily`):
| 步 | 动作 | 护栏 |
|---|---|---|
| 1 | **恢复前先验完整性**:校验 `SHA256SUMS` + `pg_restore --list` 确认可解析 | 不把损坏备份恢复到一半才发现 |
| 2 | 创建旁库 `shaguabijia_restore_<backup_id>`;已存在则提示换名或显式 `--force` 重建 | 不撞库、不误删 |
| 3 | 灌 `globals.sql`(角色已存在则跳过,幂等) | |
| 4 | `pg_restore -j <并行度,默认 2>` 到旁库(custom 格式自动识别,无需 `-Fc`),`nice`/`ionice` 降级 | 往旁库写,**不锁生产表**,对生产仅轻微资源争抢 |
| 5 | **恢复后自检报告**:表数量、关键业务表(user / wallet / withdraw 等)行数、`alembic_version`、最新记录时间戳 | 人凭报告判断「数据对不对、新不新」 |
**磁盘峰值提醒**:旁库与生产库共存于同一实例,恢复期间磁盘占用约为 2×(生产 + 旁库),恢复前脚本检查磁盘余量。
### 5.2 兜底:交互式选择(需求 4)
`db_restore.sh` **不带参数**运行时 → 扫描 `$BACKUP_DIR`,列表打印所有备份(`backup_id` / 时间 / 大小 / `status`),提示输入序号选一个恢复到旁库。这就是「平台挂了 SSH 登录手动选备份恢复」——与平台调用的是同一个脚本,零额外学习成本。
### 5.3 深度核对的边界
平台展示 §5.1 第 5 步的自检报告即可支撑「数据对不对」的判断。**逐行深度核对不在本方案范围**——需要时 DBA 直接连旁库跑只读 SQL,旁库为此存在。
## 6. 扶正生产:`db_promote.sh`(高危、人工、不进平台)
核对无误后把旁库扶正为生产。**故意不做成平台按钮**,必须人工 SSH 执行,因为不可逆且涉及短暂停机:
1. **二次确认**:要求手输生产库名 `shaguabijia` 才继续(仿 GitHub 删仓库确认)。
2. **先自动备份当前生产库**(调 `db_backup.sh --reason manual`)—— 最关键的兜底,扶错了能退回来。
3. **断开生产连接**:停 app 服务或 `pg_terminate_backend` 清连接(执行者自身不能连在待改名的库上)。**此步有短暂停机。**
4. **rename 切换**(不用 drop):`shaguabijia``shaguabijia_old_<ts>``shaguabijia_restore_xxx``shaguabijia`。出错可立即换回,比删库重建安全得多。
5. 重启 app 服务,人工验证。
**为何用脚本而非裸敲 SQL**:切换不是「只有 rename 一句」,而是一串不能漏的动作,其中「先备份当前生产库」裸敲最易漏、漏了就没退路。脚本把这串封装成带护栏的原子操作。极端情况下连脚本都不可用时,rename 本质是两条 `ALTER DATABASE` SQL,DBA 亦可手工执行,但会失去自动兜底备份保护,属下策。
## 7. 平台 API + 数据模型 + 权限
### 7.1 数据模型(两张任务表,均新建 + alembic 迁移)
`db_backup_job`
| 字段 | 类型 | 说明 |
|---|---|---|
| `id` | PK | |
| `triggered_by_admin_id` | FK admin_user, nullable | 手动才有 |
| `status` | str | `pending` / `running` / `success` / `failed` |
| `backup_id` | str, nullable | 成功后填产出目录名 |
| `error_msg` | text, nullable | |
| `created_at` / `started_at` / `finished_at` | datetime | 后两者 nullable |
| `duration_sec` | int, nullable | |
`db_restore_job`
| 字段 | 类型 | 说明 |
|---|---|---|
| `id` | PK | |
| `backup_id` | str | 恢复哪个备份 |
| `target_db` | str | 旁库名 |
| `triggered_by_admin_id` | FK admin_user | |
| `status` | str | `pending` / `running` / `success` / `failed` |
| `sanity_report` | json, nullable | §5.1 第 5 步的自检报告 |
| `error_msg` | text, nullable | |
| `created_at` / `started_at` / `finished_at` | datetime | 后两者 nullable |
模型需在 [app/models/__init__.py](../../../app/models/__init__.py) 导入以便 Alembic 发现。
**备份清单不查任务表**:定时备份不写任务表(不经过平台),但其产出的备份**出现在备份列表**(列表来自扫盘)。平台因此有两个视图:备份文件列表(扫盘,含所有来源)、操作记录(任务表 + 审计)。
### 7.2 API 端点(新建 `app/admin/routers/db_backup.py`,注册进 `app/admin/main.py`
| 方法 | 路径 | 作用 | 返回 |
|---|---|---|---|
| GET | `/admin/db-backup/backups` | 列所有备份(扫盘读 manifest) | `list[BackupItem]` |
| GET | `/admin/db-backup/health` | 新鲜度:最近成功备份时间与年龄 | `{last_success_at, age_hours, healthy}` |
| POST | `/admin/db-backup/backups` | 手动触发备份 | `202 {job_id}` |
| GET | `/admin/db-backup/jobs/{id}` | 轮询备份任务状态 | `BackupJobStatus` |
| POST | `/admin/db-backup/restores` | 恢复 `backup_id` 到旁库(body: `{backup_id}` | `202 {restore_job_id}` |
| GET | `/admin/db-backup/restores/{id}` | 轮询恢复状态 + 自检报告 | `RestoreJobStatus` |
`BackupItem``backup_id / created_at / reason / triggered_by / size_bytes / database / pg_version / alembic_version / status`
契约放 `app/schemas/db_backup.py`Pydantic)。POST 端点为**异步**语义:插任务表 → `setsid` 起脚本 → 立即 `202` 返回 `job_id`,前端轮询对应 GET 端点。
### 7.3 权限与审计
- **权限**:备份与恢复端点**均限 `super_admin`**(复用 `require_role("super_admin")`)。恢复能触及全库数据,不宜下放;将来给运维岗再引入 `require_page("db_backup")` 细分。
- **审计**:每次备份 / 恢复写现有审计日志(触发人、`backup_id`、结果),与项目其它高危操作一致。
- **不提供下载**(见 §2)。
## 8. 监控告警(贴合「日志 + openobserve + 平台徽标」现状)
三层,从被动到主动:
1. **结构化日志**:脚本成功/失败均打 `SyslogIdentifier=pg-backup` 日志,`journalctl` 可查、openobserve 可抓,失败为 ERROR 级 → 若已配 openobserve 告警规则即命中。
2. **平台新鲜度徽标**(最直观、不依赖推送):admin 备份页顶部显示「最近成功备份:X 小时前」,超 `BACKUP_FRESH_MAX_HOURS`(默认 26h)红色高亮。数据源为 `GET /admin/db-backup/health`
3. **可插拔 `notify()` hook**:脚本内通知点,现在只打日志(照搬 `notifier.py``LogNotifier` 模式),将来推送能力(飞书/短信)就绪再填,主流程不改。
**新鲜度自检兜底**`deploy/pg-backup-check.timer`+service)每天 09:00 检查最新成功备份是否在 26h 内,过期打 ERROR 日志 + 调 `notify()`。这覆盖 `OnFailure` 抓不到的盲区(timer 被禁 / 宕机没补跑)——因为那种情况脚本根本没运行,靠「有没有新备份」反向判断。
## 9. 脚本连库身份与权限
备份/恢复脚本**以本机 `postgres` 超级用户走 Unix socket(peer 认证,无需密码)**执行,而非业务的 `DATABASE_URL`TCP + 密码)。理由:
- `pg_dumpall --globals-only``CREATE DATABASE`(建旁库)、`ALTER DATABASE ... RENAME``pg_terminate_backend` 都需高权限;
- 本机 socket peer 认证是运维脚本标准做法,权限最省心,且不用把超级用户密码写进任何配置文件。
脚本以 `root``postgres` 系统用户运行(systemd service 内 `User=postgres` 或经 `sudo -u postgres`)。SSH 手动执行时同理。
## 10. 配置项
新增到 [app/core/config.py](../../../app/core/config.py) 的 `Settings`(供平台侧 §8 新鲜度等读取),shell 脚本经 systemd `EnvironmentFile=.env` 注入、SSH 手动跑时脚本内有默认值兜底:
| 配置 | 默认 | 说明 |
|---|---|---|
| `BACKUP_DIR` | `/opt/pg_backups` | 备份根目录(建议独立数据盘,与 PG 数据文件不同物理盘) |
| `BACKUP_KEEP_DAILY_DAYS` | `14` | 每日备份保留天数 |
| `BACKUP_KEEP_MONTHLY_COUNT` | `6` | 月度长留份数 |
| `BACKUP_MIN_KEEP` | `3` | 硬底线,永远至少保留份数 |
| `BACKUP_FRESH_MAX_HOURS` | `26` | 新鲜度阈值 |
| `BACKUP_PG_SUPERUSER` | `postgres` | 脚本连库超级用户 |
| `BACKUP_PGHOST` | `/var/run/postgresql` | socket 目录(peer 认证) |
| `BACKUP_S3_ENABLED` | `false` | 远程上传总开关(预留) |
| `BACKUP_S3_*` | 空 | bucket / endpoint / 凭证(预留,S3 确认后填) |
同步更新 `.env.example`
## 11. 完整文件清单
### 本仓(后端 + 脚本)
**新增:**
| 文件 | 作用 |
|---|---|
| `scripts/db_backup.sh` | 备份核心(dump + 自检 + 保留清理 + upload hook + notify |
| `scripts/db_restore.sh` | 恢复到旁库(含无参交互式选择 = 需求 4) |
| `scripts/db_promote.sh` | 旁库扶正生产(高危、人工、二次确认) |
| `app/models/db_backup_job.py` | `db_backup_job` + `db_restore_job` 模型 |
| `app/admin/repositories/db_backup_job.py` | 任务表数据访问 + 扫盘读 manifest |
| `app/admin/routers/db_backup.py` | 平台 API6 个端点) |
| `app/schemas/db_backup.py` | API 契约 |
| `alembic/versions/xxxx_add_db_backup_jobs.py` | 建两张任务表 |
| `deploy/pg-backup.service` / `pg-backup.timer` | 每日定时备份 |
| `deploy/pg-backup-check.service` / `pg-backup-check.timer` | 新鲜度自检 |
| `deploy/pg-backup.md` | 部署文档(仿 [deploy/daily-exchange.md](../../../deploy/daily-exchange.md) |
**改动:** `app/models/__init__.py`(导模型)、`app/core/config.py`(§10 配置项)、`app/admin/main.py`(注册 router)、`.env.example`
### 另一仓(`shaguabijia-admin-web`)配套
备份管理前端页:备份列表 + 一键备份 + 选备份恢复到旁库 + 恢复进度/自检报告展示 + 新鲜度徽标。本 spec 定义 API 契约(§7.2),前端据此实现,单独走该仓的开发流程。
## 12. 需求覆盖对照
| 需求 | 落地 |
|---|---|
| 每日定时备份 | `pg-backup.timer``db_backup.sh --reason daily`,凌晨 3 点,不依赖平台 |
| 平台手动备份 | `POST /admin/db-backup/backups` → 插任务 + `setsid``db_backup.sh --reason manual` |
| 平台指定备份快速恢复 | `POST /admin/db-backup/restores``db_restore.sh` 恢复到旁库 + 自检报告 |
| 平台不可用手动恢复 | SSH 跑 `db_restore.sh`(交互式选备份),同一套脚本 |
| (延伸)扶正生产 | 人工 `db_promote.sh`,二次确认 + 先备份当前 + rename 切换 |
## 13. 部署时需确认的运维参数
以下取决于服务器实际情况,部署或本 spec 复审时确认,不阻碍设计:
1. **备份目录位置**:默认 `/opt/pg_backups`;若有独立数据盘,建议放数据盘且与 PG 数据文件不同物理盘(避免一盘挂掉数据与备份同亡)。
2. **保留量匹配磁盘余量**:1–10GB 压缩后每份约几百 MB~2GB,14 天约 3–30GB;核对与磁盘余量是否匹配。
3. **备份时间**:默认凌晨 03:00(与 0 点 daily-exchange 错开);确认无其它凌晨任务撞车。
4. **脚本运行用户**:确认以 `postgres`(或可 `sudo -u postgres`)运行、socket peer 认证可用。
-3
View File
@@ -35,9 +35,6 @@ dependencies = [
# multipart form (FastAPI 表单上传依赖)
"python-multipart>=0.0.9",
# 用户反馈截图缩略图,避免 App 历史页为 48dp 小图下载数 MB 原图
"pillow>=11.0.0",
# admin 后台账号密码 hash(用户侧是手机号+验证码登录,不需要密码;admin 才用)
"bcrypt>=4.0.0",
+11 -48
View File
@@ -1,13 +1,8 @@
"""发奖公式下限:看了就保底 1 金币,与前端展示公式完全对齐
"""发奖公式下限:有真实正 eCPM 时最低发 1 金币,eCPM 缺失/为 0 仍发 0
产品口径(2026-08看了就保底 1):`calculate_ad_reward_coin` **任何输入**都至少返回 1,
镜像客户端 FeedRewardFormula.singleUnitCoin(那边:eCPM /非法/正数都返 1)此前 eCPM
缺失/ 0 0,与前端小球显示的 +1 不一致,且信息流侧把这类看满一份的广告记成 too_short 零发
公式是发奖与后台审计对账的唯一口径,改这一处两边同源
:激励视频(S2S 回调)路径在公式之前还有一道 `if not ecpm_raw` 早退**完全没上报 eCPM**
的回调仍记 ecpm_missing 零发( test_ad_reward.test_callback_without_ecpm_records_exception),
那是"回调缺字段"的数据完整性闸,区别于"广告如实上报 eCPM=0"(=看了真广告 保底 1)
针对eCPM 过低时公式四舍五入成 0 金币 被记 too_short 不发的问题:只要广告有真实
eCPM,单份金币至少 1(不再 0); eCPM 缺失/ 0/(没有真实广告价值)仍发 0,
不凭空铸币不破坏 ecpm_missing 语义公式是发奖与后台审计对账的唯一口径,改这一处两边同源
"""
from __future__ import annotations
@@ -15,23 +10,19 @@ from app.core.rewards import calculate_ad_reward_coin
def test_low_positive_ecpm_floors_to_one_coin() -> None:
"""真实但极低的 eCPM 原本四舍五入成 0 金币,兜底为 1(线上 record 4667 的 43 分场景)。"""
"""真实但极低的 eCPM 原本四舍五入成 0 金币,现在兜底为 1(线上 record 4667 的 43 分场景)。"""
# 43 分 = ¥0.43 CPM,因子1=0.1,重度用户 LT 第 69 条=1.0 → 0.43/1000×0.1×1.0×10000=0.43 → 旧口径 round=0
assert calculate_ad_reward_coin("43", 69) == 1
# 更低的 5 分同理:算出来 <0.5,旧口径也是 0
assert calculate_ad_reward_coin("5", 11) == 1
def test_zero_or_missing_ecpm_also_floors_to_one() -> None:
"""eCPM 为 0 / 缺失 / 非法都保底 1(2026-08「看了就保底 1」,与前端 FeedRewardFormula 对齐)。
尤其 "0"(广告如实上报零价值)此前返 0,导致小球显示 +1后端信息流记 too_short 零发的
前后端不一致 现在两端都 1
"""
assert calculate_ad_reward_coin("0", 1) == 1
assert calculate_ad_reward_coin(None, 1) == 1
assert calculate_ad_reward_coin("", 1) == 1
assert calculate_ad_reward_coin("abc", 1) == 1
def test_zero_or_missing_ecpm_stays_zero() -> None:
"""eCPM 缺失 / 为 0 / 非法(没有真实广告价值)不兜底,仍发 0。"""
assert calculate_ad_reward_coin("0", 1) == 0
assert calculate_ad_reward_coin(None, 1) == 0
assert calculate_ad_reward_coin("", 1) == 0
assert calculate_ad_reward_coin("abc", 1) == 0
def test_normal_ecpm_value_unchanged() -> None:
@@ -66,31 +57,3 @@ def test_feed_reward_low_ecpm_grants_one_coin_instead_of_too_short() -> None:
assert rec.coin == 1
finally:
db.close()
def test_feed_reward_zero_ecpm_grants_one_coin() -> None:
"""端到端:eCPM 如实上报 0(用户反馈的真实广告返回 eCPM=0 场景),看满一份也保底 1、
状态 granted 与前端小球显示的 +1 一致,不再前显示后零发"""
from app.db.session import SessionLocal
from app.models.user import User
from app.repositories.ad_feed_reward import grant_feed_reward
db = SessionLocal()
try:
user = User(phone="19900000044", username="feedzero44", register_channel="sms")
db.add(user)
db.commit()
db.refresh(user)
rec = grant_feed_reward(
db, user.id,
client_event_id="feed-zero-ecpm-0001",
ecpm="0", # 广告如实上报 eCPM=0
duration_seconds=15, # 看满一份
ad_type="draw",
feed_scene="comparison",
)
assert rec.status == "granted"
assert rec.coin == 1
finally:
db.close()
+1 -94
View File
@@ -1,7 +1,7 @@
"""Admin M2 读接口测试:大盘聚合 + 用户/流水/提现/反馈列表 + 鉴权拦截。"""
from __future__ import annotations
from datetime import UTC, date, datetime
from datetime import UTC, datetime
import pytest
from fastapi.testclient import TestClient
@@ -12,7 +12,6 @@ from app.admin.repositories import admin_user as admin_repo
from app.admin.repositories import queries
from app.db.session import SessionLocal, engine
from app.models.comparison import ComparisonRecord
from app.models.coupon_state import CouponClaimRecord, CouponSession
from app.models.feedback import Feedback
from app.models.invite import InviteRelation
from app.models.savings import SavingsRecord
@@ -121,86 +120,6 @@ def test_dashboard_period_comparison_is_aggregated_by_backend(
assert comparison["token_cost_total_yuan"] == pytest.approx(1.0)
def test_dashboard_coupon_success_rate_excludes_abandoned_sessions(
admin_client: TestClient, admin_token: str
) -> None:
started_date = date(2038, 1, 16)
started_at = datetime(2038, 1, 16, 8, tzinfo=UTC)
sessions = [
("coupon-rate-completed-1", "coupon-rate-device-1", "completed"),
("coupon-rate-completed-2", "coupon-rate-device-2", "completed"),
("coupon-rate-failed", "coupon-rate-device-3", "failed"),
("coupon-rate-abandoned", "coupon-rate-device-4", "abandoned"),
]
db = SessionLocal()
try:
for trace_id, device_id, status in sessions:
db.add(
CouponSession(
trace_id=trace_id,
device_id=device_id,
status=status,
app_env="prod",
platforms=["meituan-waimai"],
started_at=started_at,
started_date=started_date,
)
)
for index, device_id in enumerate(("coupon-rate-device-1", "coupon-rate-device-2")):
db.add(
CouponClaimRecord(
device_id=device_id,
coupon_id=f"mt_dashboard_rate_{index}",
claim_date=started_date,
status="success",
app_env="prod",
)
)
db.commit()
finally:
db.close()
response = admin_client.get(
"/admin/api/stats/overview",
params={"date_from": "2038-01-16", "date_to": "2038-01-16"},
headers=_auth(admin_token),
)
assert response.status_code == 200, response.text
coupon = response.json()["period"]["coupon"]
assert coupon["started"] == 4
assert coupon["abandoned"] == 1
assert coupon["success_denominator"] == 3
assert coupon["all_success"] == 2
assert coupon["success_rate"] == pytest.approx(2 / 3, abs=0.0001)
db = SessionLocal()
try:
db.add(
CouponSession(
trace_id="coupon-rate-only-abandoned",
device_id="coupon-rate-device-only-abandoned",
status="abandoned",
app_env="prod",
platforms=["meituan-waimai"],
started_at=datetime(2038, 1, 17, 8, tzinfo=UTC),
started_date=date(2038, 1, 17),
)
)
db.commit()
finally:
db.close()
empty_denominator_response = admin_client.get(
"/admin/api/stats/overview",
params={"date_from": "2038-01-17", "date_to": "2038-01-17"},
headers=_auth(admin_token),
)
assert empty_denominator_response.status_code == 200
only_abandoned = empty_denominator_response.json()["period"]["coupon"]
assert only_abandoned["success_denominator"] == 0
assert only_abandoned["success_rate"] is None
def test_user_list_and_detail(admin_client: TestClient, admin_token: str) -> None:
uid = _seed_user_with_data("13800000002")
r = admin_client.get("/admin/api/users", headers=_auth(admin_token))
@@ -801,16 +720,6 @@ def test_comparison_records_show_readable_device_and_rom_version(
rom_name="OriginOS",
rom_version=4,
android_version="14",
platforms=[
{
"platform_id": "meituan",
"platform_name": "美团",
"status": "ok",
"role": "target",
"price": 18.8,
"is_best": True,
}
],
)
db.add(record)
db.commit()
@@ -838,8 +747,6 @@ def test_comparison_records_show_readable_device_and_rom_version(
assert detail.status_code == 200, detail.text
assert detail.json()["device_model_name"] == "vivo Y77e"
assert detail.json()["rom_version"] == 4
assert detail.json()["platforms"][0]["status"] == "ok"
assert detail.json()["platforms"][0]["is_best"] is True
def test_comparison_records_show_real_order_status(
+2 -29
View File
@@ -48,10 +48,7 @@ def test_compare_start_is_idempotent_by_trace_id(client) -> None:
retry = client.post("/api/v1/compare/start", json=payload, headers=_headers(token))
assert first.status_code == 200, first.text
# trace_id 回显客户端带来的值(老协议幂等路径)
assert first.json() == {
"limit": 100, "used": 1, "remaining": 99, "trace_id": payload["trace_id"],
}
assert first.json() == {"limit": 100, "used": 1, "remaining": 99}
assert retry.status_code == 200, retry.text
assert retry.json() == first.json()
with SessionLocal() as db:
@@ -72,28 +69,6 @@ def test_compare_start_is_idempotent_by_trace_id(client) -> None:
assert record.device_id == "quota-device"
def test_compare_start_issues_trace_id_when_absent(client) -> None:
"""新客户端不带 trace_id → 服务端签发并随响应返回,running 行以签发 id 建。"""
token, user_id = _login(client)
response = client.post(
"/api/v1/compare/start",
json={"business_type": "food", "device_id": "quota-device-issue"},
headers=_headers(token),
)
assert response.status_code == 200, response.text
body = response.json()
issued = body["trace_id"]
assert issued # 非空签发
assert body["used"] == 1
with SessionLocal() as db:
record = db.execute(
select(ComparisonRecord).where(ComparisonRecord.trace_id == issued)
).scalar_one()
assert record.user_id == user_id
assert record.status == "running"
assert record.device_id == "quota-device-issue"
def test_compare_start_rejects_101st_beijing_day_attempt(client) -> None:
token, user_id = _login(client)
now = datetime.now(CN_TZ).replace(tzinfo=None)
@@ -126,9 +101,7 @@ def test_compare_start_rejects_101st_beijing_day_attempt(client) -> None:
headers=_headers(token),
)
assert allowed.status_code == 200, allowed.text
assert allowed.json() == {
"limit": 100, "used": 100, "remaining": 0, "trace_id": final_allowed_trace,
}
assert allowed.json() == {"limit": 100, "used": 100, "remaining": 0}
rejected_trace = f"quota-rejected-{user_id}"
response = client.post(
-131
View File
@@ -13,7 +13,6 @@ import uuid
from unittest.mock import MagicMock, patch
import httpx
import pytest
from sqlalchemy import select
from app.db.session import SessionLocal
@@ -26,26 +25,6 @@ def _tid() -> str:
return uuid.uuid4().hex
@pytest.mark.parametrize(
("raw_status", "record_status"),
[
("success", "success"),
("below_minimum", "success"),
("failed", "failed"),
("store_closed", "failed"),
("store_not_found", "failed"),
("items_not_found", "failed"),
("no_delivery", "failed"),
("unsupported", "failed"),
("cancelled", "cancelled"),
("running", "running"),
(None, None),
],
)
def test_record_status_normalization(raw_status, record_status) -> None:
assert crud._normalize_record_status(raw_status) == record_status
def _done_params() -> dict:
"""一份典型 done 帧 params:美团 25 元 vs 源淘宝闪购 30 元 → 省 5 元、success。"""
return {
@@ -132,116 +111,6 @@ def test_harvest_done_derives_and_newly_success_once(client) -> None:
assert newly2 is False
def test_harvest_done_below_minimum_counts_as_completed_success(client) -> None:
"""未达起送是可信业务结论:主状态/完成奖励归 success,细分结局仍留在 raw_payload。"""
tid = _tid()
done_below_minimum = {
"record_status": "below_minimum",
"comparison_results": [
{
"platform_id": "meituan",
"platform_name": "美团",
"package": "com.sankuai.meituan",
"price": 59.0,
"is_source": True,
"rank": 1,
"store_name": "测试店",
"items": [{"name": "红乌苏", "qty": 1}],
},
],
"platform_results": {
"meituan": {"is_source": True, "status": "source", "price": 59.0},
"taobao_flash": {
"is_source": False,
"status": "below_minimum",
"reason": "购物车未达起送门槛(差 ¥25.2",
},
},
"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_below_minimum
)
assert rec.status == "success"
assert newly is True
assert rec.fail_reason is None
assert rec.raw_payload["record_status"] == "below_minimum"
assert (
rec.raw_payload["platform_results"]["taobao_flash"]["status"]
== "below_minimum"
)
def test_harvest_done_platforms_below_minimum_counts_as_success(client) -> None:
"""新 platforms 单源派生路径也必须执行同一 below_minimum → success 归一化。"""
tid = _tid()
with SessionLocal() as db:
rec, newly = crud.harvest_done(
db,
trace_id=tid,
user_id=None,
done_params={
"record_status": "below_minimum",
"platforms": [
{
"role": "source",
"platform_id": "meituan",
"platform_name": "美团",
"price": 59.0,
"store_name": "测试店",
"items": [{"name": "红乌苏", "qty": 1}],
},
{
"role": "target",
"platform_id": "taobao_flash",
"platform_name": "淘宝",
"status": "below_minimum",
"price": None,
},
],
"platform_results": {
"taobao_flash": {
"is_source": False,
"status": "below_minimum",
"reason": "购物车未达起送门槛",
}
},
},
)
assert rec.status == "success"
assert newly is True
assert rec.raw_payload["record_status"] == "below_minimum"
def test_legacy_upsert_below_minimum_counts_as_success(client) -> None:
"""灰度期客户端直报路径无论走 status 还是 record_status 都不能落第四种主状态。"""
tid = _tid()
payload = ComparisonRecordIn(
trace_id=tid,
business_type="food",
status="below_minimum",
comparison_results=[],
platform_results={
"taobao_flash": {
"is_source": False,
"status": "below_minimum",
"reason": "购物车未达起送门槛",
}
},
information="淘宝未达起送门槛,可加菜凑单后下单",
)
with SessionLocal() as db:
rec = crud.upsert_record(db, user_id=987654, payload=payload)
assert rec.status == "success"
assert rec.fail_reason is None
assert rec.raw_payload["status"] == "below_minimum"
def test_harvest_done_failed_derives_fail_reason(client) -> None:
"""failed 记录:记录级 information 笼统,但 fail_reason 从 platform_results 救出具体原因
(id 3030 :美团系统失败 + 京东 items_not_found 展示京东那条)"""
+13 -34
View File
@@ -23,7 +23,7 @@ def test_postgresql_duration_summary_uses_ordered_set_aggregates() -> None:
)
assert sql.count("percentile_cont") == 4
assert "comparison_record.status IN ('success', 'below_minimum')" in sql
assert "comparison_record.status = 'success'" in sql
def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
@@ -34,9 +34,6 @@ def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
("summary-success-b", "success", 3000, 2.0, 0),
("summary-failed", "failed", 100_000, 3.0, 0),
("summary-cancelled", "cancelled", 5000, 4.0, 0),
("summary-below-minimum", "below_minimum", 5000, None, 0),
("summary-store-closed", "store_closed", 200_000, None, 0),
("summary-running", "running", 4000, None, 0),
]
for trace_id, status, total_ms, cost, saved in rows:
db.add(ComparisonRecord(
@@ -59,44 +56,26 @@ def test_summary_uses_only_success_durations_and_filters_beijing_date() -> None:
db, date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert summary["started"] == 7
assert summary["completed"] == 5
assert summary["success"] == 3
assert summary["success_rate"] == pytest.approx(3 / 6)
assert summary["started"] == 4
assert summary["completed"] == 3
assert summary["success"] == 2
assert summary["success_rate"] == pytest.approx(2 / 3)
assert summary["avg_token_cost"] == pytest.approx(2.5)
assert summary["lower_price_rate"] == pytest.approx(1 / 3)
assert summary["avg_duration_ms"] == 3000
assert summary["p5_duration_ms"] == 1200
assert summary["p50_duration_ms"] == 3000
assert summary["p95_duration_ms"] == 4800
assert summary["p99_duration_ms"] == 4960
assert summary["lower_price_rate"] == 0.5
assert summary["avg_duration_ms"] == 2000
assert summary["p5_duration_ms"] == 1100
assert summary["p50_duration_ms"] == 2000
assert summary["p95_duration_ms"] == 2900
assert summary["p99_duration_ms"] == 2980
assert summary["cancelled"] == 1
assert summary["cancelled_rate"] == pytest.approx(1 / 7)
assert summary["cancelled_rate"] == 0.25
assert summary["cancelled_p50_ms"] == 5000
items, _next_cursor, total = queries.list_comparison_records(
db, date_from=date(2038, 1, 15), date_to=date(2038, 1, 15), limit=20
)
assert total == 7
assert total == 4
assert {item.trace_id for item in items} == {row[0] for row in rows}
success_items, _next_cursor, success_total = queries.list_comparison_records(
db, status="success", date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert success_total == 3
assert {item.status for item in success_items} == {"success", "below_minimum"}
failed_items, _next_cursor, failed_total = queries.list_comparison_records(
db, status="failed", date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert failed_total == 2
assert {item.status for item in failed_items} == {"failed", "store_closed"}
running_items, _next_cursor, running_total = queries.list_comparison_records(
db, status="running", date_from=date(2038, 1, 15), date_to=date(2038, 1, 15)
)
assert running_total == 1
assert running_items[0].status == "running"
finally:
db.rollback()
db.close()
@@ -1,95 +0,0 @@
"""Regression coverage for comparison terminal-status normalization."""
from __future__ import annotations
import importlib.util
from pathlib import Path
import sqlalchemy as sa
from alembic.migration import MigrationContext
from alembic.operations import Operations
def _migration_module():
path = (
Path(__file__).parents[1]
/ "alembic"
/ "versions"
/ "comparison_below_minimum_as_success.py"
)
spec = importlib.util.spec_from_file_location(path.stem, path)
assert spec is not None and spec.loader is not None
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
def test_business_status_migration_upgrade_and_downgrade() -> None:
engine = sa.create_engine("sqlite:///:memory:")
metadata = sa.MetaData()
records = sa.Table(
"comparison_record",
metadata,
sa.Column("id", sa.Integer, primary_key=True),
sa.Column("status", sa.String(16), nullable=False),
sa.Column("fail_reason", sa.String(256)),
sa.Column("raw_payload", sa.JSON, nullable=False),
)
metadata.create_all(engine)
module = _migration_module()
with engine.begin() as connection:
connection.execute(
records.insert(),
[
{
"status": "below_minimum",
"fail_reason": "旧失败原因",
"raw_payload": {"record_status": "below_minimum"},
},
{
"status": "success",
"fail_reason": None,
"raw_payload": {"record_status": "success"},
},
{
"status": "failed",
"fail_reason": "技术异常",
"raw_payload": {"status": "failed"},
},
{
"status": "store_closed",
"fail_reason": "店铺打烊",
"raw_payload": {"record_status": "store_closed"},
},
{
"status": "items_not_found",
"fail_reason": "商品未找到",
"raw_payload": {"status": "items_not_found"},
},
],
)
module.op = Operations(MigrationContext.configure(connection))
module.upgrade()
upgraded = connection.execute(
sa.select(records.c.status, records.c.fail_reason).order_by(records.c.id)
).all()
assert upgraded == [
("success", None),
("success", None),
("failed", "技术异常"),
("failed", "店铺打烊"),
("failed", "商品未找到"),
]
module.downgrade()
downgraded = connection.execute(
sa.select(records.c.status, records.c.fail_reason).order_by(records.c.id)
).all()
assert downgraded == [
("below_minimum", None),
("success", None),
("failed", "技术异常"),
("store_closed", "店铺打烊"),
("items_not_found", "商品未找到"),
]
+3 -111
View File
@@ -10,7 +10,6 @@ from app.admin.repositories.coupon_data import (
_point_scores_by_trace,
coupon_data_report,
coupon_point_details,
coupon_user_records,
)
from app.admin.security import create_admin_token
from app.db.session import SessionLocal
@@ -39,7 +38,6 @@ def test_point_scores_by_trace() -> None:
stats = _point_scores_by_trace(db, [trace])[trace]
assert stats["succeeded"] == 2
assert stats["tried"] == 3
assert stats["events"] == 4
details = coupon_point_details(db, trace_id=trace)
assert [item["status"] for item in details] == [
"success", "already_claimed", "failed", "skipped"
@@ -50,8 +48,8 @@ def test_point_scores_by_trace() -> None:
db.close()
def test_skipped_detail_is_distinguished_from_no_events() -> None:
"""仅有 skipped 时分数仍为0/0,但保留事件数供前端开放明细"""
def test_skipped_detail_does_not_create_a_score() -> None:
"""仅有 skipped 时按需明细仍可查到,但列表没有虚假的 0/0 分数"""
db = SessionLocal()
trace = "point-score-skipped"
try:
@@ -65,7 +63,7 @@ def test_skipped_detail_is_distinguished_from_no_events() -> None:
db.flush()
scores = _point_scores_by_trace(db, [trace, "missing-trace"])
assert scores[trace] == {"succeeded": 0, "tried": 0, "events": 1}
assert trace not in scores
assert "missing-trace" not in scores
assert coupon_point_details(db, trace_id=trace)[0]["status"] == "skipped"
finally:
@@ -116,112 +114,6 @@ def test_coupon_data_report_returns_scores_without_embedding_details() -> None:
db.close()
def test_coupon_data_report_marks_abandoned_without_point_results() -> None:
"""中途退出且没有逐券终态时返回0/0,其他状态缺埋点仍保持为空。"""
db = SessionLocal()
report_date = date(2020, 1, 6)
user_id = 910006
try:
db.add_all([
CouponSession(
trace_id="point-score-abandoned-without-result",
device_id="score-abandoned-device",
user_id=user_id,
status="abandoned",
app_env="prod",
platforms=["meituan-waimai"],
started_at=datetime(2020, 1, 6, tzinfo=UTC),
started_date=report_date,
),
CouponSession(
trace_id="point-score-completed-without-result",
device_id="score-completed-device",
user_id=user_id,
status="completed",
app_env="prod",
platforms=["meituan-waimai"],
started_at=datetime(2020, 1, 6, 1, tzinfo=UTC),
started_date=report_date,
),
CouponSession(
trace_id="point-score-abandoned-with-result",
device_id="score-abandoned-result-device",
user_id=user_id,
status="abandoned",
app_env="prod",
platforms=["meituan-waimai"],
started_at=datetime(2020, 1, 6, 2, tzinfo=UTC),
started_date=report_date,
),
CouponSession(
trace_id="point-score-abandoned-skipped-only",
device_id="score-abandoned-skipped-device",
user_id=user_id,
status="abandoned",
app_env="prod",
platforms=["meituan-waimai"],
started_at=datetime(2020, 1, 6, 3, tzinfo=UTC),
started_date=report_date,
),
])
db.add_all([
CouponClaimEvent(
trace_id="point-score-abandoned-with-result",
device_id="score-abandoned-result-device",
coupon_id=f"mt-abandoned-{status}",
claim_date=report_date,
status=status,
)
for status in ("success", "failed")
])
db.add(CouponClaimEvent(
trace_id="point-score-abandoned-skipped-only",
device_id="score-abandoned-skipped-device",
coupon_id="mt-abandoned-skipped",
claim_date=report_date,
status="skipped",
))
db.flush()
report = coupon_data_report(
db,
date_from=report_date.isoformat(),
date_to=report_date.isoformat(),
app_env="prod",
)
rows = {item["trace_id"]: item for item in report["items"]}
abandoned = rows["point-score-abandoned-without-result"]
assert abandoned["point_success_count"] == 0
assert abandoned["point_total_count"] == 0
assert abandoned["point_event_count"] == 0
abandoned_with_result = rows["point-score-abandoned-with-result"]
assert abandoned_with_result["point_success_count"] == 1
assert abandoned_with_result["point_total_count"] == 2
assert abandoned_with_result["point_event_count"] == 2
abandoned_skipped = rows["point-score-abandoned-skipped-only"]
assert abandoned_skipped["point_success_count"] == 0
assert abandoned_skipped["point_total_count"] == 0
assert abandoned_skipped["point_event_count"] == 1
completed = rows["point-score-completed-without-result"]
assert completed["point_success_count"] is None
assert completed["point_total_count"] is None
user_rows = {
item["trace_id"]: item
for item in coupon_user_records(db, user_id=user_id)["items"]
}
assert user_rows["point-score-abandoned-without-result"]["point_total_count"] == 0
assert user_rows["point-score-abandoned-with-result"]["point_total_count"] == 2
assert user_rows["point-score-abandoned-skipped-only"]["point_event_count"] == 1
assert user_rows["point-score-completed-without-result"]["point_total_count"] is None
finally:
db.rollback()
db.close()
def test_coupon_point_details_endpoint() -> None:
"""前端点击使用的接口按约定返回 trace_id 和逐券 items。"""
db = SessionLocal()
+6 -12
View File
@@ -8,7 +8,6 @@ mock 掉对 pricebot 的 httpx 调用,验证:
"""
from __future__ import annotations
import json
import time
from unittest.mock import MagicMock, patch
@@ -63,8 +62,7 @@ def test_coupon_step_no_auth_required(client) -> None:
def test_coupon_step_passes_body_through(client, access_token) -> None:
"""带 token + pricebot 200 → 请求 body 原样转发到 /api/coupon/step;
响应在透传基础上顶层回显本次任务 trace_id(setdefault 注入,其余字段原样)"""
"""带 token + pricebot 200 → 响应原样透传,请求 body 原样转发到 /api/coupon/step"""
fake_pricebot_resp = {
"success": True,
"action": {
@@ -85,12 +83,9 @@ def test_coupon_step_passes_body_through(client, access_token) -> None:
captured: dict = {}
# ⚠️ coupon_step 转发用 content=raw(原始字节透传,不重新 dumps),不是 json= ——
# fake 必须捕 content。旧 fake 只捕 json= 导致 captured["json"] 恒 None,本测试
# 自 content=raw 优化后一直红着(pre-existing),本次顺手修正。
async def fake_post(self, url, content=None, **kw):
async def fake_post(self, url, json=None, **kw):
captured["url"] = url
captured["content"] = content
captured["json"] = json
mock_resp = MagicMock()
mock_resp.status_code = 200
mock_resp.json = lambda: fake_pricebot_resp
@@ -104,10 +99,9 @@ def test_coupon_step_passes_body_through(client, access_token) -> None:
)
assert r.status_code == 200, r.text
# 顶层多出 trace_id 回显(值=请求带的;不 mint,签发点唯一在 /coupon/session started)
assert r.json() == {**fake_pricebot_resp, "trace_id": "test-trace-1"}
# 验证请求被原样转发(body 字节不动 + URL 指向 pricebot)
assert json.loads(captured["content"]) == _stub_request_body()
assert r.json() == fake_pricebot_resp
# 验证请求被原样转发(body 不动 + URL 指向 pricebot)
assert captured["json"] == _stub_request_body()
assert captured["url"].endswith("/api/coupon/step")
-65
View File
@@ -1,65 +0,0 @@
"""POST /api/v1/coupon/session 的 trace_id 签发行为(统一 trace_id 由后端下发)。
- started 不带 trace_id 服务端签发并返回,行以签发 id ;
- started trace_id 回显沿用(老客户端兼容);
- started trace_id 不签发不写库,trace_id=null(防孤儿行)
"""
from __future__ import annotations
from sqlalchemy import func, select
from app.db.session import SessionLocal
from app.models.coupon_state import CouponSession
def _base_payload(**overrides) -> dict:
payload = {
"device_id": "cs-issue-device",
"status": "started",
"started_at_ms": 1_722_000_000_000,
"platforms": ["meituan-waimai"],
"app_env": "dev",
}
payload.update(overrides)
return payload
def test_session_started_issues_trace_id_when_absent(client) -> None:
response = client.post("/api/v1/coupon/session", json=_base_payload())
assert response.status_code == 200, response.text
body = response.json()
assert body["ok"] is True
issued = body["trace_id"]
assert issued
with SessionLocal() as db:
row = db.execute(
select(CouponSession).where(CouponSession.trace_id == issued)
).scalar_one()
assert row.device_id == "cs-issue-device"
assert row.status == "started"
def test_session_started_echoes_client_trace_id(client) -> None:
response = client.post(
"/api/v1/coupon/session", json=_base_payload(trace_id="cs-legacy-1")
)
assert response.status_code == 200, response.text
assert response.json()["trace_id"] == "cs-legacy-1"
with SessionLocal() as db:
count = db.scalar(
select(func.count(CouponSession.id)).where(
CouponSession.trace_id == "cs-legacy-1"
)
)
assert count == 1
def test_session_terminal_without_trace_id_skips_write(client) -> None:
response = client.post(
"/api/v1/coupon/session",
json=_base_payload(status="completed", elapsed_ms=1234, claimed_count=2),
)
assert response.status_code == 200, response.text
body = response.json()
assert body["ok"] is True
assert body["trace_id"] is None
-184
View File
@@ -1,184 +0,0 @@
from __future__ import annotations
import importlib.util
from datetime import datetime
from pathlib import Path
import sqlalchemy as sa
def _load_migration():
path = (
Path(__file__).parents[1]
/ "alembic"
/ "versions"
/ "deepseek_v4_flash_price.py"
)
spec = importlib.util.spec_from_file_location("deepseek_v4_flash_price", path)
assert spec is not None and spec.loader is not None
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
def test_migration_follows_comparison_status_normalization():
migration = _load_migration()
assert migration.down_revision == "comparison_below_min_success"
def test_migration_adds_price_and_corrects_only_mispriced_snapshot(monkeypatch):
migration = _load_migration()
engine = sa.create_engine("sqlite://")
metadata = sa.MetaData()
app_config = sa.Table(
"app_config",
metadata,
sa.Column("key", sa.String(64), primary_key=True),
sa.Column("value", sa.JSON, nullable=False),
sa.Column("updated_at", sa.DateTime),
)
comparison = sa.Table(
"comparison_record",
metadata,
sa.Column("id", sa.Integer, primary_key=True),
sa.Column("llm_calls", sa.JSON),
sa.Column("llm_cost_yuan", sa.Float),
sa.Column("llm_price_snapshot", sa.JSON),
)
metadata.create_all(engine)
calls = [
{
"model": "deepseek-v4-flash",
"error": None,
"usage": {"prompt_tokens": 12031, "completion_tokens": 125},
},
{
"model": "qwen3.5-flash",
"error": None,
"usage": {"prompt_tokens": 2416, "completion_tokens": 119},
},
]
snapshot = {
"mode": "per_model",
"prices": {
"deepseek-v4-flash": {
"input_per_1m": 3.0,
"output_per_1m": 15.0,
"_source": "default",
},
"qwen3.5-flash": {
"input_per_1m": 0.2,
"output_per_1m": 2.0,
"_source": "per_model",
},
},
}
with engine.begin() as conn:
original_updated_at = datetime(2026, 7, 13, 18, 20, 10)
conn.execute(
app_config.insert().values(
key="llm_token_price",
value={
"per_model": {
"qwen3.5-flash": {
"input_per_1m": 0.2,
"output_per_1m": 2.0,
}
},
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
},
updated_at=original_updated_at,
)
)
conn.execute(
comparison.insert().values(
id=1,
llm_calls=calls,
llm_cost_yuan=0.038689,
llm_price_snapshot=snapshot,
)
)
monkeypatch.setattr(migration.op, "get_bind", lambda: conn)
migration.upgrade()
config = conn.execute(
sa.select(app_config.c.value).where(
app_config.c.key == "llm_token_price"
)
).scalar_one()
assert config["per_model"]["deepseek-v4-flash"] == {
"input_per_1m": 1.0,
"output_per_1m": 2.0,
}
assert conn.execute(
sa.select(app_config.c.updated_at).where(
app_config.c.key == "llm_token_price"
)
).scalar_one() == original_updated_at
corrected = conn.execute(sa.select(comparison)).mappings().one()
assert corrected["llm_cost_yuan"] == 0.013002
assert corrected["llm_price_snapshot"]["prices"]["deepseek-v4-flash"] == {
"input_per_1m": 1.0,
"output_per_1m": 2.0,
"_source": "per_model",
}
assert corrected["llm_price_snapshot"]["pricing_correction"] == (
"deepseek_v4_flash_price"
)
migration.downgrade()
reverted_config = conn.execute(
sa.select(app_config.c.value).where(
app_config.c.key == "llm_token_price"
)
).scalar_one()
assert "deepseek-v4-flash" not in reverted_config["per_model"]
reverted = conn.execute(sa.select(comparison)).mappings().one()
assert reverted["llm_cost_yuan"] == 0.038689
assert "pricing_correction" not in reverted["llm_price_snapshot"]
def test_downgrade_preserves_price_that_existed_before_upgrade(monkeypatch):
migration = _load_migration()
engine = sa.create_engine("sqlite://")
metadata = sa.MetaData()
app_config = sa.Table(
"app_config",
metadata,
sa.Column("key", sa.String(64), primary_key=True),
sa.Column("value", sa.JSON, nullable=False),
sa.Column("updated_at", sa.DateTime),
)
sa.Table(
"comparison_record",
metadata,
sa.Column("id", sa.Integer, primary_key=True),
sa.Column("llm_calls", sa.JSON),
sa.Column("llm_cost_yuan", sa.Float),
sa.Column("llm_price_snapshot", sa.JSON),
)
metadata.create_all(engine)
explicit_price = {"input_per_1m": 1.0, "output_per_1m": 2.0}
with engine.begin() as conn:
conn.execute(
app_config.insert().values(
key="llm_token_price",
value={
"per_model": {"deepseek-v4-flash": explicit_price},
"default": {"input_per_1m": 3.0, "output_per_1m": 15.0},
},
)
)
monkeypatch.setattr(migration.op, "get_bind", lambda: conn)
migration.upgrade()
migration.downgrade()
config = conn.execute(sa.select(app_config.c.value)).scalar_one()
assert config["per_model"]["deepseek-v4-flash"] == explicit_price
-4
View File
@@ -108,10 +108,6 @@ def test_get_llm_prices_falls_back_to_default_then_uses_override():
# 无 override → CONFIG_DEFS 默认(含 per_model / default)
prices = get_llm_prices(db)
assert "per_model" in prices and "default" in prices
assert prices["per_model"]["deepseek-v4-flash"] == {
"input_per_1m": 1.0,
"output_per_1m": 2.0,
}
# 有 override → 用 DB 值
app_config.set_value(
db, "llm_token_price",
-77
View File
@@ -1,77 +0,0 @@
from __future__ import annotations
from app.core.config import settings
from app.db.session import SessionLocal
from app.models.comparison import ComparisonRecord
from app.models.user import User
from app.repositories import ops_marquee
def _clear_real_rows_cache() -> None:
ops_marquee._real_rows_cache.update(at=None, rows=None, test_phones=None)
def test_real_marquee_records_exclude_configured_test_account(monkeypatch) -> None:
legacy_test_phone = "19900009991"
listed_test_phone = "19900009992"
real_phone = "19900009993"
monkeypatch.setattr(settings, "TEST_ACCOUNT_PHONE", legacy_test_phone)
monkeypatch.setattr(settings, "TEST_ACCOUNT_PHONES", listed_test_phone)
_clear_real_rows_cache()
try:
with SessionLocal() as db:
legacy_test_user = User(
phone=legacy_test_phone,
username="29900009991",
register_channel="sms",
)
listed_test_user = User(
phone=listed_test_phone,
username="29900009992",
register_channel="sms",
)
real_user = User(
phone=real_phone,
username="29900009993",
register_channel="sms",
)
db.add_all([legacy_test_user, listed_test_user, real_user])
db.flush()
db.add_all([
ComparisonRecord(
user_id=legacy_test_user.id,
trace_id="marquee-legacy-test-account-trace",
status="success",
saved_amount_cents=29_991,
items=[],
comparison_results=[],
),
ComparisonRecord(
user_id=listed_test_user.id,
trace_id="marquee-listed-test-account-trace",
status="success",
saved_amount_cents=29_992,
items=[],
comparison_results=[],
),
ComparisonRecord(
user_id=real_user.id,
trace_id="marquee-real-account-trace",
status="success",
saved_amount_cents=29_993,
items=[],
comparison_results=[],
),
])
db.flush()
app_rows = ops_marquee._recent_real_rows(db)
browse_rows, _ = ops_marquee.list_real_records(db, mode="real", limit=1_000)
finally:
_clear_real_rows_cache()
assert all(saved not in {29_991, 29_992} for _uid, saved, _nickname in app_rows)
assert any(saved == 29_993 for _uid, saved, _nickname in app_rows)
assert all(row["saved_amount_cents"] not in {29_991, 29_992} for row in browse_rows)
assert any(row["saved_amount_cents"] == 29_993 for row in browse_rows)